59 / BROKER · MESSAGE BUS / PRODUCTION FABRIC
59 / PRODUCTION / PUB-SUB · ROUTING · DELIVERY · DECOUPLING

BROKER
/ MESSAGE BUS.

Broker / Message Bus — транспортный слой, через который producers публикуют сообщения, а один или несколько consumers получают их по routing rules, topics, subscriptions или consumer groups без прямой жёсткой связи producer → конкретный process.

Главный принцип: broker отвечает за доставку сообщений между компонентами. Он не должен подменять semantics события, workflow orchestration или business state. Сообщение — транспортный envelope; смысл принадлежит domain contract.
00. ARCHITECTURAL STATUS

НЕ НУЖЕН, ПОКА ПРОСТЫЕ DIRECT CALLS И DB QUEUE ДОСТАТОЧНЫ

Message bus появляется, когда есть много независимых producers/consumers, fan-out, asynchronous domain events, cross-service decoupling, replay/retention requirements или routing становится сложнее, чем один work queue.
TYPEPRODUCTIONMessaging / transport infrastructure.
DEFAULTCONDITIONALНе обязательный компонент малого монолита.
ENABLE WHENMESSAGE TOPOLOGY GROWSFan-out, multiple consumers, service decoupling.
SEPARATE COMPONENTYESLogical transport boundary.
LIVES INPRODUCTION FABRICSupporting infrastructure.
COMPLEXITYMEDIUM → HIGHStart only when justified.
IMPLEMENT: WHEN DECOUPLING PAYS
Минимум 80% ценности: versioned message envelope, explicit topics/channels, producer/consumer ownership, at-least-once assumptions, idempotent consumers, durable subscriptions where needed, retention/replay policy, partition/routing key, dead-letter handling, schema compatibility and observability. Не нужен «универсальный enterprise event bus» для трёх процессов.
01A. ARCHITECTURE BOUNDARIES & OPERATIONS

EXPLICIT SYSTEM CONTRACT

A. BOUNDARY WITH NEIGHBORS

№42 Events & Triggers владеет domain/event semantics: что произошло и какая reaction возможна; №59 переносит сообщение. №57 Queues & Workers владеет executable jobs, leases, attempts и worker lifecycle; broker может доставлять jobs, но это отдельная responsibility. №58 Scheduler создаёт time-originated occurrence; broker лишь транспортирует его. №50 Contracts задаёт message schemas/version compatibility. №61 Provenance/Lineage отслеживает origin/transformations data; broker хранит message metadata/correlation, но не весь lineage. №68 Durable Workflow владеет process history/replay; broker replay не равен workflow replay.

B. PREREQUISITES / CROSS-REFERENCES

Prerequisites: №42 Events, №46 Observability, №50 Contracts, №57 Queues & Workers, №58 Scheduler. Forward references: №61 Provenance, №63 Retry/Circuit Breakers, №64 Rate Limits/Budgets, №68 Durable Workflow, №69 Distributed Reliability.

C. PLANE PLACEMENT

REQUEST-TIME: producer may publish synchronously, but consumers usually async. CONTROL PLANE: topics, subscriptions, routing, retention, schemas, ACLs, partitions. DATA PLANE: message envelopes and delivery/ack state. OFFLINE: replay, re-drive, schema migration tests, load/failure tests.

D. FAILURE & OPERATIONS CONTRACT

Success: accepted message reaches all required durable subscriptions according to delivery contract. Retryable: transient broker/consumer/network outage. Permanent: invalid schema, forbidden topic, unsupported version. Delivery: duplicates expected unless effect-level guarantees prove otherwise. Idempotency: message_id/event_id + consumer-side dedupe/effect key. Persist: message metadata, topic, partition/routing key, schema version, offsets/ack state where relevant. Trace: publish→broker→subscription→consume→ack.

E. WHAT THIS TOPIC DOES NOT OWN

№59 не владеет domain event meaning, task execution lifecycle, timers, workflow state, data warehouse, source ingestion semantics или business compensation. Она владеет MESSAGE ROUTING, DELIVERY, FAN-OUT AND SUBSCRIPTION TRANSPORT.

01. WHY A MESSAGE BUS

PRODUCER НЕ ДОЛЖЕН ЗНАТЬ ВСЕХ CONSUMERS

DIRECT COUPLING

Service A вызывает B, C и D напрямую.

Добавление E требует менять A; failure одного consumer влияет на producer path.

BROKER

A публикует stable event/message в topic.

Broker хранит routing/delivery semantics.

INDEPENDENT CONSUMERS

B/C/D/E подписываются самостоятельно и развиваются независимо.

Message bus особенно ценен при fan-out: одно domain событие может независимо запустить indexing, analytics, notification, audit и downstream sync.
02. BROKER VS QUEUE

ОДНО СООБЩЕНИЕ — ОДИН WORKER ИЛИ МНОГО INDEPENDENT CONSUMERS?

WORK QUEUE

Competing consumers

Одна logical job должна быть выполнена одним из workers. Несколько workers конкурируют за work item.

Пример: document.parse.

PUB-SUB BUS

Independent subscriptions

Одно событие может получить каждый заинтересованный consumer/subscription.

Пример: document.updated → search index + audit + analytics.

Некоторые brokers поддерживают обе модели. Но выбирайте semantics по задаче, а не по названию продукта.
03. MESSAGE ENVELOPE

СООБЩЕНИЕ ДОЛЖНО ИМЕТЬ СТАБИЛЬНУЮ ИДЕНТИЧНОСТЬ И CONTRACT

{
  "message_id": "MSG-...",
  "message_type": "document.updated",
  "schema_version": "2.0",
  "occurred_at": "...",
  "published_at": "...",
  "tenant_id": "tenant_A",
  "producer": "ingestion-service",
  "subject_ref": "doc://tenant_A/123",
  "correlation_id": "TRACE-...",
  "causation_id": "MSG-parent-...",
  "routing_key": "tenant_A.document",
  "payload": {
    "source_version": "v42"
  }
}
ENVELOPE FIELDS

Separate transport from domain payload

  • message/event id;
  • message type;
  • schema version;
  • tenant/security context reference;
  • occurred vs published time;
  • producer identity;
  • subject/resource reference;
  • correlation + causation;
  • routing/partition key;
  • payload or payload_ref.
04. EVENT VS COMMAND VS MESSAGE

НЕ ВСЕ ENVELOPES ОЗНАЧАЮТ ОДНО И ТО ЖЕ

TypeMeaningTypical ownership
EVENTFact: something already happened.Producer owns fact; consumers choose reaction.
COMMANDRequest that a specific capability perform an action.Usually one logical handler / work queue.
NOTIFICATIONInformational signal, may be lossy depending contract.Consumer may ignore.
STATE SNAPSHOTCurrent representation of object/state.Useful for projection/cache/update, but may be large/stale.
Event = past tense fact. Не называть send_email событием, если это command. Точная семантика сильно влияет на retry, idempotency и ownership.
05. TOPICS & SUBSCRIPTIONS

ROUTING ТОПОЛОГИЯ ДОЛЖНА БЫТЬ ПОНЯТНОЙ, А НЕ МАГИЧЕСКОЙ

PRODUCERS ingestion / workflow connectors / scheduler MESSAGE BUS topic: document.events topic: workflow.events topic: audit.events SEARCH INDEXER durable subscription ANALYTICS independent offset AUDIT separate consumer
Subscription — самостоятельный delivery state. Если search indexer отстал, analytics не должен ждать его и наоборот.
06. ROUTING KEY

КАК BROKER ПОНИМАЕТ, КУДА И В КАКОМ PARTITION ИДЁТ MESSAGE

TOPIC

Broad domain

document.events, workflow.events, audit.events.

MESSAGE TYPE

Specific contract

document.updated, document.deleted.

ROUTING KEY

Selective delivery

tenant/resource/category/region where topology requires.

PARTITION KEY

Ordering locality

Messages with same key land in same ordered partition if broker supports it.

Не кодировать всю бизнес-логику в topic names. Topic taxonomy должна быть достаточно стабильной и coarse-grained, а точная semantics — в message type/schema.
07. DELIVERY SEMANTICS

ДОСТАВКА СООБЩЕНИЯ И BUSINESS EFFECT — НЕ ОДНО И ТО ЖЕ

MODEL
LOSS
DUPLICATES
CONSUMER
USE
NOTE
AT-MOST-ONCE
possible
low
simple
telemetry/noncritical notification
Loss acceptable.
AT-LEAST-ONCE
low
expected
idempotent
domain events / commands
Best practical default.
EXACTLY-ONCE LOGICAL EFFECT
desired
hidden by effect design
transaction/idempotency required
critical effects
Achieved at boundaries, not by marketing flag.
Даже если broker имеет exactly-once features, external APIs, DB writes и downstream systems требуют собственного idempotency/transaction design.
08. ACK & CONSUMER OFFSET

CONSUMER ДОЛЖЕН ПОДТВЕРДИТЬ, ЧТО ОН ДОШЁЛ ДО DURABLE EFFECT

DELIVERBroker exposes message to subscription.
VALIDATESchema/security/tenant/version.
PROCESSConsumer performs local effect.
PERSIST EFFECTUpsert/project/result durable.
ACK / COMMIT OFFSETOnly after required effect.
NEXTAdvance subscription progress.
Ack-before-effect создаёт loss window. Effect-before-ack создаёт duplicate window. Поэтому idempotent consumer нужен независимо от broker.
09. DURABLE VS EPHEMERAL SUBSCRIPTIONS

НЕ КАЖДЫЙ CONSUMER ДОЛЖЕН ПОЛУЧИТЬ MESSAGE ПОСЛЕ ДНЯ OFFLINE

DURABLE

Catch up later

Subscription progress retained. Consumer after restart continues from last acknowledged offset/message.

Подходит для index, audit, business projections.

EPHEMERAL

Only while connected

Old messages may be irrelevant. Подходит для live UI hints, transient monitoring, noncritical notifications.

Durability должна быть explicit per subscription. Хранить всё навсегда «на всякий случай» — дорого и усложняет privacy/governance.
10. RETENTION & REPLAY

REPLAY МОЖЕТ БЫТЬ СИЛЬНЫМ ИНСТРУМЕНТОМ — И ОПАСНЫМ

RETENTION

How long messages live

Hours/days/weeks based on recovery, audit and governance needs.

REPLAY

Rebuild consumer state

Новый/починенный consumer может пройти historical events заново.

SIDE EFFECT RISK

Don't resend the world

Replay public/purchase/send commands without idempotency may repeat real-world actions.

Event replay безопаснее для deterministic projections. Для external side effects нужен separate command/effect ledger and idempotency.
11. ORDERING

ГЛОБАЛЬНЫЙ ORDER ПОЧТИ ВСЕГДА СЛИШКОМ ДОРОГ И НЕ НУЖЕН

Ordering scopeExampleRecommendation
NoneIndependent analytics events.Max throughput; consumer handles concurrency.
Per entityUpdates for one document/customer/workflow.Partition key = entity id.
Per tenantTenant-local ordered projection.Can create hotspot for large tenants.
GlobalTotal system sequence.Avoid unless domain truly requires it.
Чаще всего нужен order per aggregate/entity, а не единый порядок всей системы.
12. OUT-OF-ORDER EVENTS

CONSUMER ДОЛЖЕН УМЕТЬ ПОНЯТЬ, ЧТО ПРИШЛА СТАРАЯ ВЕРСИЯ

VERSION

Entity version

Event carries source/entity version so projection can ignore stale update.

OCCURRED_AT

Temporal hint

Useful but wall-clock alone is weaker than authoritative sequence/version.

REFETCH

Recover authoritative state

For critical mutable state, event can trigger live read of current object.

Не полагаться на arrival order как на business truth, если producer/network/broker can reorder or retry messages.
13. SCHEMA EVOLUTION

PRODUCER И CONSUMER ДЕПЛОЯТСЯ НЕ ОДНОВРЕМЕННО

ADD FIELD

Usually compatible

Consumers ignore unknown optional fields.

DEPRECATE

Transition window

Keep old field/schema until all consumers migrate.

BREAKING

New version/type

Use explicit schema version and compatibility strategy.

REGISTRY

Optional later

Schema registry becomes useful at scale; simple repo contracts are enough initially.

№50 Structured Outputs & Contracts owns the contract discipline. Broker should preserve schema version and reject/route unsupported messages predictably.
14. TRANSACTIONAL OUTBOX

КАК НЕ ПОТЕРЯТЬ EVENT МЕЖДУ DB COMMIT И BROKER PUBLISH

DUAL WRITE PROBLEM

DB + broker separately

Application updates DB, then publishes event. Crash between operations leaves state changed without event — or event published before DB commit.

OUTBOX

One local transaction

Business state and outbox row commit together. Separate relay publishes outbox messages and marks them delivered/retries idempotently.

BEGIN;

UPDATE documents
SET version = 42
WHERE id = 123;

INSERT INTO outbox(
  message_id,
  topic,
  payload,
  status
) VALUES (..., 'document.events', ..., 'PENDING');

COMMIT;

OUTBOX RELAY:
  claim pending row
  publish to broker
  mark SENT / retry

CONSUMER:
  dedupe by message_id / effect key
Outbox — один из ключевых distributed reliability patterns. №69 позже свяжет его с inbox, idempotency и delivery semantics.
15. INBOX / CONSUMER DEDUPE

MESSAGE ID ДОЛЖЕН ДОЙТИ ДО EFFECT BOUNDARY

consumer(message):

  begin transaction

  if inbox.exists(
      consumer="search-indexer",
      message_id=message.id
  ):
      return ACK

  apply_projection(message)

  inbox.insert(
      consumer="search-indexer",
      message_id=message.id
  )

  commit

  ACK
INBOX VALUE

One logical application

Consumer-side inbox/dedupe table помогает:

  • переживать duplicate delivery;
  • atomic apply + mark processed;
  • строить deterministic replay;
  • отлаживать конкретный message;
  • измерять duplicates.
16. DEAD LETTERS

НЕСОВМЕСТИМОЕ СООБЩЕНИЕ НЕ ДОЛЖНО БЛОКИРОВАТЬ SUBSCRIPTION НАВСЕГДА

POISON MESSAGE

Always fails

Unsupported schema, invalid payload, deterministic consumer bug.

DLQ

Isolate

After bounded attempts, message/subscription delivery moves to dead-letter destination/state.

REDRIVE

Fix then replay

After consumer/schema repair, controlled re-drive preserves original id/correlation.

DLQ ownership must be explicit. «Сообщения лежат в DLQ» без alert/owner/runbook означает скрытую потерю бизнеса.
17. BACKPRESSURE & SLOW CONSUMERS

ОДИН МЕДЛЕННЫЙ CONSUMER НЕ ДОЛЖЕН ТОРМОЗИТЬ ВСЕХ

SEPARATE OFFSET

Independent progress

Каждая durable subscription идёт своим темпом.

LAG

Measure backlog

Consumer lag/oldest unprocessed age — primary health signal.

SCALE

Partition-aware

Добавлять consumers while preserving required ordering and downstream capacity.

SHED

Noncritical only

Ephemeral/low-value consumers may drop or sample under pressure if contract allows.

18. SECURITY & TENANCY

MESSAGE BUS — НЕ ДОВЕРЕННАЯ ВНУТРЕННЯЯ ПОМОЙКА

PUBLISH ACL

Who may write

Producer identity scoped to allowed topics/message types.

SUBSCRIBE ACL

Who may read

Consumer access restricted by topic/tenant/data classification.

SECRETS

Do not embed

Use references/capabilities; avoid raw credentials in messages.

TENANT

Boundary

Tenant metadata validated and applied to routing/storage/consumer authorization.

Internal message content may still contain untrusted external text. Broker transport does not elevate its trust level.
19. CORRELATION & CAUSATION

ПОСЛЕДОВАТЕЛЬНОСТЬ СООБЩЕНИЙ ДОЛЖНА БЫТЬ ОБЪЯСНИМОЙ

MSG A

document.updated

correlation = TRACE-7

CONSUMER

Indexer processes A and emits document.indexed.

MSG B

correlation = TRACE-7
causation_id = MSG-A

Correlation связывает одну business/run цепочку. Causation показывает, какое конкретное сообщение вызвало следующее. Это критично для observability и incident debugging.
20. EVENT LOOPS & STORMS

ASYNC SYSTEM МОЖЕТ СЛУЧАЙНО НАЧАТЬ ГЕНЕРИРОВАТЬ САМУ СЕБЯ

LOOP

A → B → A

Consumer события A пишет state, который снова публикует A без change detection.

FAN-OUT STORM

One message → thousands

Unbounded downstream fan-out exceeds quotas or creates cascading backlog.

GUARDS

Control amplification

Change detection, causation tracing, hop limits where relevant, dedupe, rate/admission controls.

Message bus делает coupling слабее, но hidden feedback loops — сложнее. Causation metadata и metrics по amplification factor очень полезны.
21. BROKER REPLAY ≠ EVENT SOURCING

RETENTION LOG НЕ АВТОМАТИЧЕСКИ СТАНОВИТСЯ AUTHORITATIVE DATABASE

MESSAGE RETENTION

Transport history

Messages retained to recover consumers, replay processing or audit transport.

EVENT-SOURCED STATE

Domain state model

Events are canonical source of truth and current state is derived by replay. Requires stronger domain/version/invariant discipline.

Не выбирать event sourcing только потому, что broker умеет хранить сообщения неделями.
22. BROKER VS WORKFLOW

CHAIN OF EVENTS НЕ РАВЕН DURABLE ORCHESTRATION

BROKER CHOREOGRAPHY

Loose reactions

Services independently react to facts. Хорошо, когда no single component owns full end-to-end process.

№68 ORCHESTRATION

Owned process

Нужны explicit state, steps, timers, retries, compensation, completion criteria, history/replay.

Если бизнес спрашивает «на каком шаге находится заявка и кто отвечает за завершение?», одного message bus обычно недостаточно.
23. OBSERVABILITY

СМОТРЕТЬ НА PUBLISH, DELIVERY И CONSUMER LAG ОТДЕЛЬНО

PUB

Publish rate

Messages/sec by topic/type/producer.

P95

Publish latency

Producer → broker accepted.

LAG

Consumer lag

Oldest unprocessed age / offset distance.

DEL

Delivery latency

published_at → consumer process/ack.

DUP

Duplicate rate

Repeated message ids delivered/observed.

DLQ

Dead letters

Count, age, types, ownership.

AMP

Amplification

Messages produced per input event; detect storms.

DROP

Loss / reject

Rejected schema/ACL/retention expiration where measurable.

24. TESTING & FAILURE INJECTION

ПРОВЕРЯТЬ DUPLICATES, REORDERING И DOWNTIME

DUP

Redelivery

Same message twice → one logical consumer effect.

ORDER

Reordering

Old entity version after new one does not corrupt projection.

OFFLINE

Consumer downtime

Durable subscription catches up after restart.

POISON

Bad schema

Moves to dead-letter after bounded attempts, does not block forever.

OUTBOX

Producer crash

Crash between DB write and publish is eventually reconciled.

BROKER OUTAGE

Backpressure

Producer/outbox buffers safely, no silent loss.

SCHEMA

Compatibility

Old/new producer-consumer versions coexist.

TENANT

Isolation

Unauthorized subscriber cannot read another tenant's protected stream.

25. FAILURE MODES

КАК MESSAGE BUS ЛОМАЕТ АРХИТЕКТУРУ

BUS FOR EVERYTHING
Даже простой request/response превращается в async choreography.
USE DIRECT CALLS WHEN SIMPLE
EVENT = COMMAND
Consumer semantics and ownership становятся неоднозначными.
NAME SEMANTICS
NO MESSAGE ID
Consumer не может dedupe/replay/audit.
STABLE ID
ACK BEFORE EFFECT
Consumer crash теряет logical processing.
ACK AFTER DURABILITY
ASSUME GLOBAL ORDER
Scale/retry breaks hidden ordering assumptions.
ORDER PER KEY
NO SCHEMA VERSION
Independent deployments silently break consumers.
VERSIONED CONTRACTS
DB + PUBLISH DUAL WRITE
State and event diverge on crash.
OUTBOX
REPLAY SIDE EFFECTS
Historical messages resend emails/orders/publishes.
IDEMPOTENT EFFECT LEDGER
MESSAGE CHAINS = WORKFLOW
End-to-end process ownership/state disappears.
№68 WHEN PROCESS MATTERS
26. METRICS

ЧТО ИЗМЕРЯТЬ

PUB

Publish Rate

Message throughput by topic/type/producer.

LAG

Consumer Lag

Oldest unprocessed age/offset distance per subscription.

P95

End-to-End Latency

published → consumer effect/ack.

DUP

Duplicate Delivery

Repeated message IDs and duplicate logical effects.

DLQ

Dead-Letter Rate

Poison/incompatible messages by contract and consumer.

AMP

Amplification Factor

Downstream messages per source message/event.

OUT

Outbox Age

Oldest committed event not yet published.

SCH

Schema Reject

Messages rejected for unsupported/invalid contract version.

27. MVP IMPLEMENTATION

НЕ ОБЯЗАТЕЛЬНО НАЧИНАТЬ С KAFKA

messages/
├── contracts/
├── publisher.py
├── subscriptions.py
├── outbox.py
├── inbox.py
├── handlers/
└── tests/

outbox(
  message_id,
  topic,
  message_type,
  schema_version,
  tenant_id,
  payload_json,
  status,
  created_at,
  published_at
)

consumer_inbox(
  consumer_id,
  message_id,
  processed_at,
  primary key(consumer_id, message_id)
)

MVP options:
  A) PostgreSQL outbox + polling consumers
  B) DB queue + fan-out tables
  C) lightweight broker when multiple
     independent consumers justify it
80% VALUE MVP

Contracts before infrastructure

  • Define event/command semantics.
  • Stable message_id + schema version.
  • Topic/message taxonomy.
  • Transactional outbox at producer seam.
  • Idempotent consumer/inbox.
  • One durable subscription per independent consumer.
  • Bounded retries + dead-letter state.
  • Correlation/causation IDs.
  • Lag and delivery metrics.
  • Only then choose dedicated broker if needed.

Для небольшого монолита PostgreSQL outbox + background dispatcher может дать основную decoupling/reliability ценность без отдельного cluster.

28. WHEN TO USE A DEDICATED BROKER

ИНФРАСТРУКТУРА ДОЛЖНА СЛЕДОВАТЬ ЗА ТРЕБОВАНИЯМИ

SignalWhy broker helps
Many independent consumers / fan-outSubscriptions, consumer groups, routing become first-class.
High throughput / sustained message volumeDedicated log/broker handles streaming better than OLTP DB tables.
Need retention/replayHistorical stream can rebuild projections/consumers.
Service boundaries across hosts/teamsDecoupled transport and ownership become valuable.
Partitioned ordered streamsBroker partitions/offsets support scalable ordered consumption.
Independent backpressure per consumerEach subscription advances at its own rate.
29. PRACTICAL DECISION

СТОИТ ЛИ ДЕЛАТЬ ОТДЕЛЬНЫЙ КОМПОНЕНТ?

ВопросОтвет
Стоит ли реализовывать?Только когда есть реальная multi-consumer/event-driven topology. Для малого монолита direct calls + PostgreSQL queue/outbox часто лучше.
Separate Component?YES как transport responsibility. Физически может появиться позже, когда dedicated broker оправдан.
Минимум 80% ценности?Stable message contracts, outbox, idempotent consumers, topics/subscriptions, schema versioning, correlation, retries/DLQ, lag metrics.
Когда overkill?Kafka cluster ради двух фоновых handlers и десятка событий в час.
Trigger?One producer needs multiple independent consumers, replay/retention, decoupled services, high-volume asynchronous message routing.
Как измерить uplift?Producer/consumer coupling, fan-out reliability, consumer lag, recovery/replay time, duplicate effects, incident isolation, time to add new consumer.
Можно ли rule/tool/code вместо LLM-agent?Полностью. Broker/message bus — deterministic infrastructure. LLM may produce domain data, but never owns delivery semantics.
30. DESIGN RULES

ПРАВИЛА ДЛЯ РЕАЛЬНОЙ СИСТЕМЫ

RULE 01

Semantics before transport

Сначала решить event/command meaning, затем выбирать broker/topic.

RULE 02

Stable message identity

message_id/schema_version обязательны для replay/dedupe/audit.

RULE 03

Assume duplicates

Consumer effects должны быть idempotent.

RULE 04

Order only where needed

Prefer entity/partition ordering over global ordering.

RULE 05

Outbox at dual-write seams

DB state and publish intent commit together.

RULE 06

Independent subscriptions

Slow consumer should not block unrelated consumers.

RULE 07

Replay consciously

Replaying transport history must not repeat unsafe side effects.

RULE 08

Observe causation

Correlation + causation reveal event chains and storms.

RULE 09

Don't build bus too early

Monolith + DB outbox is a valid intermediate architecture.

31. FINAL MAP

MOVE MESSAGES WITHOUT LOSING DOMAIN MEANING

DOMAIN / SYSTEM COMPONENT
        ↓
SOMETHING HAPPENS
        ↓
CREATE VERSIONED MESSAGE
  message_id
  message_type
  schema_version
  tenant
  subject
  occurred_at
  correlation
  causation
        ↓
[ TRANSACTIONAL OUTBOX IF DB STATE CHANGED ]
        ↓
PUBLISH
        ↓
BROKER / MESSAGE BUS
  topics
  routing
  partitions
  retention
  subscriptions
        ↓
        ├─ CONSUMER A
        │    ↓
        │  validate
        │  idempotent effect
        │  ACK / offset
        │
        ├─ CONSUMER B
        │    ↓
        │  own independent lag/state
        │
        └─ CONSUMER C
             ↓
           own independent lag/state

FAILURES:
  duplicate → inbox/idempotency
  stale/out-of-order → version check
  poison → DLQ
  consumer offline → durable subscription
  producer crash → outbox replay
  message storm → backpressure / quotas

BOUNDARIES:

№42 EVENT/TRIGGER
  = WHAT HAPPENED / WHAT IT MEANS

№57 QUEUE/WORKER
  = WHAT WORK IS WAITING / WHO EXECUTES

№58 SCHEDULER
  = WHEN A TIME OCCURRENCE EXISTS

№59 BROKER
  = HOW MESSAGES ARE ROUTED
    AND DELIVERED TO CONSUMERS

CORE PRINCIPLE:

A MESSAGE BUS SHOULD
REDUCE COUPLING,
NOT HIDE OWNERSHIP.

IT MOVES VERSIONED,
TRACEABLE MESSAGES.

IT DOES NOT DECIDE
WHAT THE BUSINESS MEANS,
WHETHER A WORKFLOW IS COMPLETE,
OR WHETHER AN EXTERNAL EFFECT
IS SAFE TO REPEAT.

ECC RETROFIT / PRACTICAL HARNESS INTEGRATION

A. Related ECC ideas. Context-as-cache, scoped memory, lifecycle hooks, selective capabilities, feature flags, deterministic enforcement, provider-neutral adapters and eval-gated learning are applied only where relevant to №59 Broker / Message Bus.

B–E. Existing boundary and placement. The existing conceptual boundary, class PRODUCTION, default CONDITIONAL and owner Production Fabric remain authoritative. Runtime/control/data/offline placement is unchanged; durable state stays outside model context.

F–H. Hooks and contracts. Use bounded PRE_MODEL/POST_MODEL, PRE_TOOL/POST_TOOL, CHECKPOINT and TASK_COMPLETED events as applicable. Illustrative fields and canonical contracts are defined in NEW_CONTRACTS_SPEC.md; no universal schema is implied.

I–J. Security and evaluation. Host-side schema, permission, secret, budget, idempotency and audit checks take precedence over LLM output. Optional mechanisms require a feature flag and WITH/WITHOUT ablation; measure quality, acceptance, correction, latency, cost, escalations and severe errors.

K–L. Task profiles and cross-references. A TaskProfile selects the relevant skill, tool/context slice, memory scope and enforcement profile independently from FAST/STANDARD/DEEP. See cross-reference map, hook spec and ablation plan. Provider adapters remain outside the core.