Skip to content

ADR-014 — Обработка сообщений (SendTask / ReceiveTask, message events, брокер и producer/consumer-шов)

Поле Значение
Статус Принято
Версия v.2
Дата 2026-06-14 (v.2 2026-08-08)
Владелец Руслан Габитов
Уточняет ADR-001 v.5 Execution Model

Принято — реализовано сопровождающими его приземляющими SRD (task- половина: SendTask/ReceiveTask + MessageWaiter; event-половина: промежуточные throw/catch message events + producer/consumer-шов). Фазовые отсрочки (вывод correlation-key, message-triggered инстанцирование — §2.8) остаются открытыми для своих именованных последователей. Решает, как gobpm отправляет и принимает BPMN-сообщения: исполнители SendTask и ReceiveTask, message throw/catch-события, их разделяемый producer/consumer-шов (отложенный сюда из ADR-011 v.5 §2.6) и как оба ездят на брокере сообщений и машинерии ожидания событий. Охват фазовый: этот ADR решает базовую send/receive-модель; вывод correlation-key и message-triggered инстанцирование отложены к именованным последователям (§2.8). Приземляющий SRD делает файловую работу и code-grounded.

1. Контекст

1.1 Что требует стандарт

BPMN 2.0 (§8.3.2 Correlation, §8.4 Messages & Operations, §10.3.x Tasks, §10.4.2 Data, §11 Collaboration) определяет обмен сообщениями между процессом и внешним участником:

  • Message — это контент, обмениваемый между двумя участниками. Message ссылается на ItemDefinition (структуру его payload'а). MessageFlow соединяет отправителя с получателем через границы пулов; внутри процесса send/receive моделируется задачами или message-событиями.
  • SendTask отправляет сообщение и завершается, как только оно отправлено. Он ссылается на Message (контент) и, опционально, на Operation (когда отправка — service-вызов). Он не ждёт ответа.
  • ReceiveTask ждёт сообщения и завершается, когда оно приходит. Он ссылается на Message; его флаг instantiate, выставленный на первой активности процесса, позволяет пришедшему сообщению запустить новый инстанс процесса.
  • Throw/catch message events зеркалят задачи. Throw message event эмитит сообщение (как SendTask); catch message event ждёт его (как ReceiveTask). Семантика сообщения идентична — различается лишь моделирующий элемент (активность vs событие).
  • Сообщения коррелируются к инстансам (§8.3.2). CorrelationKey вычисляется из свойств сообщения через retrieval-выражения; пришедшее сообщение маршрутизируется к инстансу, чей ключ совпадает. Корреляция — это то, что делает "правильное сообщение достигает правильного запущенного процесса".
  • Message flow асинхронен отправителю, но синхронен жизненному циклу получателя. Отправитель эмитит и продолжает; активность получателя не завершается (и не эмитит токенов), пока подходящее сообщение не придёт и его данные не будут связаны.

1.2 Что у движка сейчас

Части существуют, но message-путь не подключён:

  • Брокер сообщений существует. pkg/messaging определяет MessageBroker (Publish(Envelope) / Subscribe(name, key) → <-chan Envelope) с in-memory- дефолтом; Envelope несёт payload, имя сообщения и плоский correlation-key. Рантайм движка уже выставляет его исполнениям.
  • Машинерия ожидания событий существует. EventHub регистрирует waiter на каждое event-определение и возобновляет ожидающий трек при срабатывании события; трек переходит в состояние "ожидание события" и обратно. Реализован только timer-waiter — message-waiter'а нет, так что catch message event или ReceiveTask нечего ждать.
  • Throw работает только из end-события, внутренне. Throwing-событие эмитит через propagate-путь EventHub'а, не через брокер — так что сегодня брошенное сообщение никогда не покидает движок во внешний канал.
  • SendTask и ReceiveTask — заглушки. Это структуры только из полей (Message, рудиментарный Operation, флаг ReceiveTask.instantiate) без исполнителя — они не могут работать.
  • Корреляция смоделирована, но не подключена. Типы CorrelationKey / CorrelationSubscription существуют; ничто не вычисляет и не матчит ключ в run time. Брокер матчит по имени сообщения плюс плоскому ключу.

1.3 Почему сейчас

ADR-011 v.5 §2.6 решил Operation у ServiceTask и явно отложил node-уровневый message-handling-шов (MessageProducer / MessageConsumer) к "исполнительскому SRD, где есть несколько реализаторов, чтобы зафиксировать его форму." Этот момент здесь: SendTask, ReceiveTask и throw/catch message events — ровно те реализаторы. Решение message-модели сейчас даёт движку его первую реальную cross-participant-способность и чистый шов, который четыре элемента разделяют, без переоткрытия data-flow-ADR'ов.

2. Решение

2.1 Сообщения ходят через брокер; EventHub остаётся внутренней машиной ожидания

gobpm держит два различных канала и мост между ними:

  • Брокер сообщений — внешний канал сообщений. Каждое сообщение, которое процесс отправляетSendTask или throw message event — публикуется в MessageBroker как Envelope. Каждое сообщение, которое процесс ждётReceiveTask или catch message event — принимается подпиской на брокер. Брокер — граница, через которую сообщения пересекают к и от внешних участников (и между инстансами).
  • EventHub остаётся внутренней машиной возобновления токена. Он владеет per-definition-waiter'ом и жизненным циклом wait/resume трека; это то, чем ожидающий узел паркуется и будится. Он не транспорт сообщений.
  • MessageWaiter мостит эти двое. Ожидающий message-узел регистрирует MessageWaiter (недостающий замковый камень), который подписывается на брокер на своё сообщение и, когда приходит Envelope, срабатывает событием в EventHub — так что трек возобновляется тем же путём, что использует timer. Это держит одну единообразную wait/resume-модель для каждого вида триггера и замыкает знание о брокере на waiter.

Обоснование: messaging по своей природе — boundary-забота (он пересекает участников), так что ему место на брокере, а не на внутренней шине событий; но ожидающий ReceiveTask должен интегрироваться с существующим жизненным циклом track-resume, так что waiter, а не узел, владеет broker-подпиской и адаптирует её к EventHub'у.

2.2 Producer/consumer-шов, разделяемый задачами и message-событиями

Обработка сообщений разделяется по направлению, как два узких контракта (шов, отложенный ADR-011 v.5):

  • MessageProducer — биндит свой Message из process scope и публикует его в брокер. Реализуется SendTask и throw message event.
  • MessageConsumerподписывается на свой Message, и по приходу биндит payload в process scope. Реализуется ReceiveTask и catch message event.

Шов — это два интерфейса, не один: BPMN message handling действительно направленный (send/throw vs receive/catch), и узел — то или другое. Задача и событие одного направления разделяют одну и ту же producer/consumer-реализацию, так что message-хореография (bind → publish, или subscribe → bind) живёт ровно в одном месте на направление, а activity/event-обёртка лишь адаптирует её к своей форме исполнения (активность завершается; событие выстреливает токен).

2.3 SendTask — это message-producer, синхронный своему жизненному циклу

SendTask исполняется как любой node-executor: он биндит свой Message из данных исполнения (payload сообщения заполняется из process scope input-data- ассоциациями активности), публикует получившийся Envelope в брокер (имя сообщения + payload; correlation-key по §2.6) и завершается, эмитя свои исходящие потоки. Он не ждёт ответа (request/reply-обмен — это два узла, send затем receive — диаграмма показывает ожидание). Отправка синхронна жизненному циклу активности: узел не Complete, пока publish не вернётся.

2.4 ReceiveTask — это message-consumer, который ждёт, затем биндит

ReceiveTask регистрирует MessageWaiter на свой Message и переводит свой трек в состояние ожидания (§2.1). Когда подходящий Envelope приходит, waiter срабатывает событием; трек возобновляется, исполнитель биндит payload в process scope (через output-data-ассоциации активности — обратное input-биндингу SendTask), и узел завершается, эмитя свои исходящие потоки. ReceiveTask никогда не таймаутится сам по себе (дедлайн моделируется boundary timer event'ом — диаграмма показывает ожидание, по принципу no-hidden-wait ADR-011).

2.5 MessageWaiter параллелит TimerWaiter

MessageWaiter построен так же, как timer-waiter: реестр waiter'ов получает TriggerMessage-билдер; waiter крутит свой service-loop, подписывается на брокер на (message name, correlation key) и на первом подходящем Envelope зовёт event-processing-путь, который возобновляет зарегистрированный(е) трек(и). Waiter несёт пришедший payload вместе со сработавшим определением к ДОСТАВКЕ: принимающее выполнение захватывает его и связывает из собственного контекста выполнения — никогда из состояния узла, которое N конкурентных выполнений одного разделяемого узла портили бы по построению (это решает ADR-006 v.5 §2.9.1; v.2 синхронизирует формулировку). Это делает ожидание сообщения равноправным ожиданию таймера — одна waiter-абстракция, много видов триггеров — а не специальным путём.

2.6 Корреляция фазовая: name-match сейчас, вывод ключа позже

Фаза 1 (этот ADR) маршрутизирует сообщение по его имени (Envelope.Name брокера), с correlation-key, проброшенным дословно при наличии, но не выведенным из модели. Этого достаточно для распространённого single-message-single-waiter-обмена и для примеров.

Полный вывод correlation-key — вычисление CorrelationKey из retrieval- выражений CorrelationSubscription над payload'ом сообщения и его матч к целевому инстансу — отложен к последователю (§2.8). Шов и поле Envelope.CorrelationKey спроектированы нести ключ, так что добавление вывода позже не переформирует producer/consumer-контракты.

2.7 Message-triggered инстанцирование отложено

ReceiveTask (или start message event) с выставленным instantiate, запускающий новый инстанс на пришедшем сообщении, отложен (§2.8). Ему нужно, чтобы движок маршрутизировал broker-сообщение к определению процесса (не к запущенному инстансу) и порождал инстанс — thresher-уровневая забота за пределами per-instance-исполнительской модели, которую решает этот ADR. ReceiveTask фазы 1 работает внутри уже запущенного инстанса.

2.8 Не-цели и вне охвата (у каждой именованный дом)

  • Вывод correlation-key (выражения CorrelationSubscription → ключ, составные ключи, key-based-маршрутизация) — последующий Correlation SRD/ADR; шов уже несёт ключ.
  • Message-triggered инстанцирование (instantiate ReceiveTask / start message event, порождающий инстанс) — последователь, с thresher message- routing-работой.
  • Send на базе service-operation (SendTask/ReceiveTask, чей транспорт — вызов service.Operation, а не broker-сообщение) — отложенная альтернатива (§4); message handling фазы 1 — это Message + брокер, так что рудиментарное поле Operation на этих задачах убирается приземляющим SRD и вновь вводится только когда понадобится service-backed messaging.
  • Reply/timeout/transactional доставка, гарантии порядка сообщений, dead-letter-обработка — broker-quality-заботы, принадлежащие реализации брокера и будущему ADR Distribution & Scale, не model-слою.
  • Durable subscriptions / персистентность ожидающего ReceiveTask через рестарт — Persistence ADR.

3. Последствия

  • Движок получает реальный cross-participant messaging. Send и receive, задачи и события, все ездят на одном брокере и одной wait-модели — первая способность, что позволяет процессу gobpm говорить с внешним миром (и с другими инстансами).
  • Одна хореография на направление. Producer/consumer-шов означает, что bind→publish и subscribe→bind каждый живут один раз; SendTask/throw-event и ReceiveTask/catch-event — тонкие обёртки, так что четыре элемента не могут разойтись.
  • Wait-модель остаётся единообразной. Message-ожидание равноправно timer-ожиданию через MessageWaiter; жизненный цикл трека, история и resume-путь не изменены — нет второго wait-механизма.
  • Никаких скрытых ожиданий. ReceiveTask ждёт видимо (это узел на диаграмме); дедлайн — это boundary timer, согласно ADR-011 §2.3. Движок никогда не блокируется на невидимом условии.
  • Корреляция честна о своей фазе. Name-match сейчас задокументирован как ограничение; шов несёт ключ, так что вывод приземляется позже без изменения контракта. Один инстанс на имя сообщения — допущение фазы 1.
  • Цена: новый waiter, два исполнителя, шов и broker-round-trip в hot path. Приземляющий SRD стейджит это (waiter сначала, затем исполнители) и держит make ci зелёным на каждом шаге.

4. Рассмотренные альтернативы

  • Маршрутизировать сообщения через EventHub вместо брокера. Переиспользовать внутренний propagate-путь и для сообщений. Отвергнуто: messaging — это boundary-забота (он пересекает участников и инстансы); сворачивание его во внутреннюю шину токенов смешивает "возобновить мой собственный трек" с "доставить через участников" и не даёт места для внешних транспортов. Брокер — правильная граница; EventHub остаётся внутренним.
  • Один унифицированный интерфейс MessageHandler вместо producer/consumer. Единый контракт с send и receive. Отвергнуто: BPMN message handling направленный, и узел — ровно одно направление; унифицированный handler заставляет каждого реализатора заглушать половину, которую он не использует (та же причина, по которой ADR-011 разделил виды, а не приболтил оба к одному типу).
  • Сделать SendTask/ReceiveTask тонкими шимами над throw/catch message events. Отвергнуто как первичная структура: жизненные циклы активности и события различаются (активность завершается и гоняет data-ассоциации; событие выстреливает токен). Разделение происходит на producer/consumer-слое (§2.2), а не делая задачу притворяющейся событием.
  • Реализовать вывод correlation-key сейчас. Конформно и полно, но это самодостаточная под-задача (вычисление выражений над payload'ом, составные ключи, маршрутизация), которая не меняет send/receive-форму. Отложено (§2.6), чтобы держать охват этого ADR на базовом пути; шов несёт ключ.
  • Реализовать message-triggered инстанцирование сейчас. Отвергнуто для этой фазы: это thresher/routing-забота (сообщение → определение → новый инстанс), ортогональная per-instance-исполнительской модели, решённой здесь (§2.7).
  • Send через service.Operation (web-service-стиль). Стандарт позволяет SendTask вызвать Operation. Отвергнуто для фазы 1: канал сообщений gobpm — брокер; Operation-backed-транспорт — отложенная альтернатива (§2.8), в которую producer-шов может врасти без переформирования вызывающих.

5. Рекомендации enterprise-готовности

Совещательно, не гейтинг — для приземляющего(их) SRD и последующей работы:

  • Поднимайте провалившийся publish / никогда-не-приходящее сообщение как инцидент. Ошибка broker-publish или ReceiveTask, ждущий за (смоделированным) дедлайном, — операционное событие, которое владелец процесса должен видеть — структурированный, классифицированный failure, несущий активность и имя сообщения, а не тихий стопор.
  • Сделайте ограничение name-match фазы 1 явным моделлерам. Пока не приземлится вывод correlation-key, два живых инстанса, ждущих то же имя сообщения, неоднозначны; user-facing-доки должны заявлять single-waiter-допущение, чтобы моделлер не был удивлён.
  • Логируйте обработку сообщений по name/key/ids, никогда payload'ы. Payload'ы сообщений бизнес-чувствительны; логируйте имя сообщения, correlation-key, id item'ов и состояния — согласно рекомендации маскирования ADR-010/011.
  • Держите broker-контракт сменяемым. In-memory-брокер — дефолт; интерфейс MessageBroker должен оставаться достаточно узким, чтобы реальный транспорт (Kafka, NATS, …) вставлялся под Distribution & Scale ADR без касания model-слоя.

6. Открытые вопросы

  • Нет. Модель send=publish / receive=MessageWaiter-subscribe, разделение broker-vs-EventHub с мостящим waiter'ом, направленный producer/consumer-шов, name-match-корреляция фазы 1 и отсрочка вывода ключа и инстанцирования решены выше. Точные сигнатуры интерфейсов, форма service-loop'а waiter'а и data-binding-проводка — заботы реализации для приземляющего SRD, не открытые концептуальные вопросы.

7. Ссылки

  • SAD-001 v.1 Vision & Architecture — §14 Conformance & Compliance Scope; цель BPMN Process Execution Conformance, которой служит этот messaging.
  • ADR-001 v.5 Execution Model — двухслойный рантайм, жизненный цикл трека и node-executor-контракт, в который этот messaging включается.
  • ADR-006 v.1 Events & Subscriptions — концепция event-доставки и wait-узла; message events — это message-типизованный случай его catch/throw-модели. Sibling — этот ADR решает message-специфику, которую тот оставляет открытой.
  • ADR-011 v.5 Process Data Flow — §2.6 отложил node-уровневый MessageProducer/MessageConsumer-шов сюда; data-ассоциации, которые биндят сообщение к/из scope, — это его §2.4.
  • BPMN 2.0 §8.3.2 (Correlation), §8.4 (Messages & Operations), §10.3 (Tasks — Send/Receive), §10.4.2 (Data associations), §11 (Collaboration & message flow) — message-модель, которую этот ADR кодирует (и, для корреляции/инстанцирования, фазирует).

История документа

Версия Дата Автор Изменение
v.2 2026-08-08 Руслан Габитов Синхронизация формулировки с ADR-006 v.5 §2.9.1 — «waiter несёт payload к узлу» в §2.5 описывала отставленный захват в состоянии узла: со времени SRD-085 payload — свойство ДОСТАВКИ, захватываемое принимающим выполнением и связываемое из его собственного контекста (слот на узле портился N конкурентными выполнениями по построению — случай параллельного MI). Ни один механизм, которым владеет ЭТОТ ADR, не меняется; контракт доставки живёт в ADR-006 v.5 §2.9.
v.1 2026-06-14 Руслан Габитов Draft. Решает обработку сообщений: сообщения ходят через MessageBroker (внешний канал), пока EventHub остаётся внутренней машиной ожидания, смостённые новым MessageWaiter (равноправным timer-waiter'у); направленный producer/consumer-шов (MessageProducer = SendTask + throw message event; MessageConsumer = ReceiveTask + catch message event) несёт одну хореографию bind→publish / subscribe→bind на направление (шов, отложенный из ADR-011 v.5 §2.6). SendTask публикует и завершается; ReceiveTask ждёт через MessageWaiter, затем биндит payload в scope. Фазовое ядро: корреляция маршрутизирует по имени сообщения сейчас (вывод ключа отложен), а message-triggered инстанцирование отложено (§2.8); рудиментарное поле Operation на задачах убрано (send на базе service-operation — отложенная альтернатива). Уточняет ADR-001 v.5; sibling к ADR-006 v.1 и ADR-011 v.5.
v.1 2026-06-16 Руслан Габитов Статус Draft → Accepted: решённый охват полностью реализован task-half-SRD (SendTask/ReceiveTask + MessageWaiter) и event-half-SRD (промежуточные throw/catch message events + шов MessageProducer/MessageConsumer в pkg/model/msgflow). Без изменения содержания; отсрочки §2.8 (вывод correlation-key, message-triggered инстанцирование, boundary message events, messaging на базе service-operation) остаются открытыми для своих именованных последователей. Добавлен RU-twin.