docs(P02): complete messaging domain phase — v0.4

---ci---
project: atelier
phase: 2
milestone: v0.4
status: complete
phase_role: execution
phase_tag: v0.3.2
requirements:
  covered: [ATELIER-97, ATELIER-98, ATELIER-99, ATELIER-100, ATELIER-101]
  partial: []
---/ci---
This commit is contained in:
Jon Chery
2026-08-05 16:00:31 +00:00
parent f7f007dce8
commit 013bd7258e
5 changed files with 1602 additions and 0 deletions
+379
View File
@@ -0,0 +1,379 @@
# Delivery Semantics — Derived Rules
> Derives from `domains/messaging/first-principles.md`. Applies P4
> (Delivery Semantics are Explicit), P3 (Consumers are Idempotent),
> P2 (Ordering is a Property, Not an Assumption), P5 (Dead-Letter
> Handling is Defined), and P6 (Backpressure is Bounded) primarily,
> with P10 (DLQ depth as an alert). For the dead-letter strategy
> decision, see the comparison table below. Includes a fenced
> idempotency-key dedup-store example (IDEATE-39). Cross-links
> `domains/concurrency/patterns` for the in-process retry/backoff
> analog, `domains/errors/patterns` for errors as data for message
> failures, and `domains/observability/metrics` for DLQ depth as an
> alert.
## The Three Delivery Semantics (P4 Delivery Semantics are Explicit)
- **At-most-once**: a message is delivered 0 or 1 times; loss is
possible, duplication is not. The producer fires-and-forgets; the
broker does not ack; the consumer does not dedup. Lowest latency,
lowest implementation cost, lossy. Fits telemetry where a dropped
sample is acceptable (MQTT QoS 0, fire-and-forget metrics).
- **At-least-once**: a message is delivered 1 or more times;
duplication is possible, loss is not. The producer sends and
waits for a broker ack; the consumer processes and acks; a crash
before the consumer's ack triggers redelivery. The consumer must
be idempotent (P3). The default for side-effecting operations
(orders, payments, commands). The vast majority of broker-backed
queues (SQS standard, RabbitMQ ack, MQTT QoS 1).
- **Exactly-once**: a message is delivered exactly 1 time; no
loss, no duplication. In practice, this is at-least-once plus
idempotency (P3), or a transactional consume-process-produce loop
(P4 — see `domains/messaging/streams.md`). Jepsen analyses verify
broker claims: "exactly-once" requires independent verification;
the durable engineering practice is at-least-once with idempotent
consumers (P3), which collapses to exactly-once under correct
dedup.
- The semantic is declared per channel (P4), not emergent. An
unstated semantic is a defect: the consumer guesses, and the
guess is wrong under the first failure. See
`domains/messaging/queues.md` for the three-semantics comparison
table (semantics, latency cost, implementation cost, when each
fits).
## Idempotency (P3 Consumers are Idempotent)
- Idempotency is the correctness property that makes at-least-once
safe. A consumer that processes the same message twice has the
same effect as processing it once. The mechanism is the
idempotency key: a per-message unique identifier the consumer
uses to dedup redeliveries.
- The idempotency key is per-message, not per-producer or
per-session. A consumer that dedups by producer alone drops
distinct messages issued in the same window. Use a UUID per
message, or a deterministic key derived from the message content
(e.g., `(entity, operation, version)`).
- The dedup store is bounded (P6 — Backpressure is Bounded): a
dedup store that grows without bound is a memory leak. Use a TTL
window longer than the broker's max-redelivery window, or a
bounded LRU. The TTL is the P6 bound: a key seen within the TTL
is a redelivery; a key older than the TTL is expired (the broker
has given up redelivering it).
## Idempotency-Key Dedup Store (IDEATE-39, P3 — fenced example)
The dedup store is the concrete mechanism that makes a consumer
idempotent under at-least-once delivery. This is the fenced
example required by IDEATE-39 (parallel to the v0.3 IDEATE-29
signed-attestation fenced example): it is NOT prose-only — the
consumer-with-dedup-store demonstrates P3 concretely.
```python
# Idempotent consumer with a dedup store (P3 Consumers are
# Idempotent, P6 Backpressure is Bounded — the dedup store is
# TTL-bounded). The consumer dedups by idempotency key before
# processing; a redelivered message is a no-op, not a double-apply.
import time
# P6: the dedup store is bounded by a TTL window. The TTL must
# exceed the broker's max-redelivery window; beyond the TTL, the
# key is expired (the broker has given up). A dedup store with no
# TTL is a memory leak (P6 violation).
DEDUP_TTL_SECONDS = 24 * 3600 # longer than max-redelivery window
class DedupStore:
"""A TTL-bounded idempotency-key dedup store (P3, P6).
seen(key): True if the key was processed within the TTL.
mark(key): Record the key as processed (with a timestamp).
"""
def __init__(self, backend):
# backend is a Redis, a DB, or an in-process LRU. The
# backend must be shared across consumer instances if the
# subscription is shared (P2/P3 — see domains/messaging/
# pubsub.md on shared vs independent subscriptions).
self.backend = backend
def seen(self, key: str) -> bool:
ts = self.backend.get(key)
if ts is None:
return False
if time.time() - ts > DEDUP_TTL_SECONDS:
# P6: expired. The broker has given up redelivering;
# this key is no longer a redelivery signal.
self.backend.delete(key)
return False
return True
def mark(self, key: str):
self.backend.set(key, time.time(), ttl=DEDUP_TTL_SECONDS)
# The idempotent consumer: dedup before process, mark after
# process, ack after mark. A crash before mark re-processes (the
# dedup store does not have the key); a crash before ack
# redelivers (the broker did not see the ack) and the dedup store
# makes the redelivery a no-op (P3).
def consume_idempotent(broker, dedup: DedupStore, process):
for message in broker.receive():
# P3: dedup BEFORE process. A redelivered message is a
# no-op, not a double-apply.
if dedup.seen(message["idempotencyKey"]):
broker.ack(message) # already processed; skip
continue
try:
process(message["payload"])
# P3: mark AFTER process succeeds. A crash between
# process and mark re-processes (acceptable: the
# process must be idempotent OR the mark must be
# transactional with the process — see below).
dedup.mark(message["idempotencyKey"])
broker.ack(message)
except TransientError as exc:
# P6: bounded retry with backoff. Nack for redelivery;
# the broker redelivers after exponential backoff.
broker.nack(message, delay=backoff(message["attempt"]))
except PoisonError as exc:
# P5: poison message. Route to DLQ, do NOT retry
# forever (see the DLQ routing rule below).
route_to_dlq(broker, message, exc)
broker.ack(message) # remove from the origin queue
```
- The order `process → mark → ack` gives at-least-once with
idempotent dedup: a crash before `mark` re-processes (the dedup
store does not have the key), and a crash before `ack`
redelivers (the dedup store makes the redelivery a no-op). If
`process` is not itself idempotent, the `process → mark` window
must be transactional (e.g., process and mark in one DB
transaction) — otherwise a crash in the window double-applies.
- For a non-idempotent `process` (e.g., a payment that must not
double-charge), use a transactional dedup: process and mark in
one DB transaction, so the mark commits iff the process
commits. This is the "exactly-once via idempotency" pattern
(P4): at-least-once delivery plus a transactional
process-and-mark collapses to exactly-once under correct
transactional semantics.
## Ordering (P2 Ordering is a Property, Not an Assumption)
- The delivery semantic interacts with ordering (P2). At-least-once
with per-partition ordering: a redelivery within a partition
preserves order (the redelivered message re-appears in its
original position relative to other messages the consumer has
not yet seen). At-least-once with no ordering: a redelivery may
appear out of order relative to messages delivered after it.
- The consumer must not assume an ordering property the broker
does not provide (P2). A standard queue delivers per-receive-node
arrival order but no global order and no order across redeliveries;
a FIFO queue delivers strict per-group order including across
redeliveries; a partitioned stream delivers strict per-partition
order across redeliveries. Document the property; do not assume
it. See `domains/messaging/queues.md` (FIFO vs standard) and
`domains/messaging/streams.md` (per-partition order).
## Dead-Letter Strategies (P5 Dead-Letter Handling is Defined)
- A poison message (unparseable, repeatedly failing, or exhausting
the retry budget) must be routed to a dead-letter queue, not
retried forever or silently dropped (P5). The DLQ is observable
(P10 — depth is an alert) and drainable (an operator can
inspect, replay, or discard with audit).
- The dead-letter strategy determines when a message is
dead-lettered and what the operator sees. The choice is the
decision matrix below (D-069).
## Dead-Letter Strategy Comparison (D-069)
| Strategy | When It Applies | Failure Visibility | Operational Cost |
|----------|-----------------|-------------------|------------------|
| **Retry-count-limit** | A fixed max-redeliveries count (e.g., 5). After N redeliveries, route to DLQ. Simple, predictable. | The redelivery count is visible in the DLQ entry; the operator sees how many times it was retried. | Low — a counter per message; no backoff tuning. Risk: retries fire as fast as the broker redelivers, hammering a downstream that is already failing (P6 — no backoff = no backpressure escape). |
| **TTL-with-backoff** | A max time-to-live for redelivery (e.g., 30 minutes) with exponential backoff between retries. After the TTL, route to DLQ. | The TTL and the backoff schedule are visible; the operator sees the retry timeline. | Medium — backoff tuning per message type. Benefit: backoff gives the downstream time to recover (P6 — the retry rate is bounded); fits transient failures (a downstream that is briefly unavailable). |
| **Poison-queue** | A separate queue for messages that fail a specific check (unparseable, schema-invalid, unknown type) before any processing retry. Routed immediately, not retried. | The poison queue is a separate signal from the DLQ; the operator sees parse-vs-process failures distinctly. | Low — a routing rule per check. Benefit: distinguishes "never going to succeed" (poison) from "might succeed on retry" (DLQ). Use for unparseable messages that no retry will fix. |
| **DLQ + alert** | Any of the above strategies, plus an alert on DLQ depth. The DLQ is observable (P10) — depth, age, and rate are alerted. | Highest — the operator is paged on DLQ growth; the DLQ is a first-class signal, not a graveyard. | Medium — alerting setup per DLQ. This is the P5/P10 floor: a DLQ without an alert is a silent correctness defect (poison messages accumulate invisibly). |
- The default for transient failures is **TTL-with-backoff + DLQ
+ alert**: backoff gives the downstream time to recover (P6),
the TTL bounds the retry budget (P5), the DLQ captures the
unprocessable, and the alert makes it visible (P10). The default
for unparseable messages is **poison-queue + alert**: route
immediately, do not retry a message no retry will fix.
- **Retry-count-limit alone** (no backoff, no alert) is the
`messaging-unbounded-retry` chaos anti-pattern's cousin: it
caps the count but hammers the downstream at full retry rate,
and the DLQ grows silently if no alert is wired. Always pair
a retry budget with backoff (P6) and an alert (P10).
- The failure-visibility column is the P10 check: every strategy
must surface the failure to the operator. A strategy with no
visibility is a P10 violation regardless of its retry semantics.
- The operational-cost column is the C8 tradeoff: more visibility
and more backoff cost more to set up but pay back in operational
stability. The P5/P10 floor is "DLQ + alert"; below that, the
strategy is a silent defect waiting to grow.
## DLQ Routing Rule (P5, P10)
```python
# DLQ routing rule (P5 dead-letter handling, P6 bounded retry with
# backoff, P10 DLQ depth alert). Combines TTL-with-backoff for
# transient failures and poison-queue for unparseable messages,
# with an alert on DLQ depth.
import json, time
MAX_RETRY_TTL_SECONDS = 30 * 60 # 30 min total retry window
DLQ = "orders-dlq"
POISON = "orders-poison" # unparseable; never retried
def route_to_dlq(broker, message, reason):
"""Route a message to the DLQ with audit metadata (P5)."""
broker.send(DLQ, body=json.dumps({
"original": message,
"reason": str(reason),
"deadLetteredAt": now_iso(),
"redeliveryCount": message.get("attempt", 0),
}))
# P10: emit a metric so DLQ depth alerts fire. A DLQ that
# grows with no alert is a silent correctness defect (P5/P10).
metrics.increment("dlq.depth", tags={"queue": "orders"})
# The ack removes the message from the origin queue; the DLQ
# is the durable record (P5 — observable and drainable).
def consume_with_dlq(broker, dedup, process):
for message in broker.receive():
# P1: parse first. An unparseable message is poison —
# route immediately, do NOT retry (no retry will fix it).
try:
payload = json.loads(message["body"])
except (ValueError, SchemaError) as exc:
broker.send(POISON, body=json.dumps({
"original": message["body"],
"reason": f"parse-failed: {exc}",
"poisonedAt": now_iso(),
}))
metrics.increment("poison.depth", tags={"queue": "orders"})
broker.ack(message) # remove from origin; poison queue holds it
continue
# P3: dedup before process.
if dedup.seen(payload["idempotencyKey"]):
broker.ack(message); continue
attempt = payload.get("attempt", 0)
first_attempt_ts = payload.get("firstAttemptTs", time.time())
try:
process(payload)
dedup.mark(payload["idempotencyKey"])
broker.ack(message)
except TransientError as exc:
# P6: TTL-with-backoff. If the retry window is
# exhausted, route to DLQ; otherwise redeliver with
# exponential backoff.
if time.time() - first_attempt_ts > MAX_RETRY_TTL_SECONDS:
route_to_dlq(broker, message, exc) # P5
broker.ack(message)
else:
broker.nack(message, delay=backoff(attempt))
except PermanentError as exc:
# A permanent error (e.g., a not-found dependency)
# does not benefit from retry — route to DLQ now.
route_to_dlq(broker, message, exc)
broker.ack(message)
```
- The routing rule distinguishes three failure modes: **poison**
(unparseable — route immediately, no retry), **transient**
(retry with backoff until the TTL, then DLQ), and **permanent**
(a retry will not fix it — DLQ now). This distinction is the P5
discipline: not every failure is a retry; some are immediate
DLQs.
- The DLQ entry carries `reason`, `deadLetteredAt`, and
`redeliveryCount` — it is auditable (the operator knows why each
message was dead-lettered and how many times it was retried).
This is the errors-as-data discipline — see
`domains/errors/patterns` for the general principle a DLQ entry
instantiates.
## Retry Budgets and Backoff (P6 Backpressure is Bounded)
- The retry budget is the cap on redelivery: a count, a TTL, or
both. After the budget, the message routes to the DLQ (P5). An
unbounded retry budget is the `messaging-unbounded-retry` chaos
anti-pattern (P5 breach): the consumer never makes progress past
the poison message.
- **Exponential backoff** spaces retries: 1s, 2s, 4s, 8s, ... with
a jitter to avoid thundering-herd synchrony. Backoff gives the
downstream time to recover (P6 — the retry rate is bounded,
giving the downstream a chance to catch up). A retry with no
backoff hammers the downstream at full rate, making the failure
worse.
- The retry budget × the backoff schedule is the P6 bound: the
consumer's retry load is bounded by design, not by luck. See
`domains/concurrency/patterns` Pattern 6 (Timeout on Every
Block) for the in-process retry/backoff analog; messaging owns
the broker-backed instance where the redelivery comes from the
broker across a network, not an in-process loop (D-062).
## Poison Messages (P5, P1)
- A poison message is one no retry will fix: unparseable (the
schema is wrong, P1), unknown type (the consumer does not handle
this version, P9), or a permanent failure (a not-found
dependency). Retrying a poison message wastes resources and
blocks the queue (P6 — the consumer never makes progress).
- Poison messages route to the poison queue immediately (no
retry), distinct from the DLQ (which holds messages that
exhausted their retry budget on transient failures). The
distinction is the P5 discipline: a poison queue is for
"never going to succeed"; a DLQ is for "might have succeeded
but didn't within the budget."
- A poison queue without an alert is the same defect as a DLQ
without an alert (P10): the operator cannot see the poison
accumulating. Wire both to `domains/observability/metrics`.
## Cross-Link to Concurrency (P6, cross-link concurrency/patterns)
- The retry/backoff discipline here is the cross-process analog
of `domains/concurrency/patterns` Pattern 6 (Timeout on Every
Block) for in-process retry. Concurrency owns the in-process
analog (a retry loop within one program, with a timeout per
attempt); messaging owns the broker-backed instance (the broker
redelivers across a network, the consumer applies backoff via
nack-with-delay). The failure model differs: in-process retry
fails by a thread crash or a timeout; broker-backed retry fails
by network partition, broker restart, or consumer crash-and-retry
(D-062).
- The cross-link is one-directional outward (messaging →
concurrency) per D-026 extended: messaging references concurrency
as the in-process foundation; concurrency does not back-link to
messaging.
## Cross-Link to Errors (P5, cross-link errors/patterns)
- A poison message is an errors-as-data instance: the failure is
captured as a DLQ entry (with `reason`, `redeliveryCount`,
`deadLetteredAt`), not swallowed. See `domains/errors/patterns`
for the general errors-as-data discipline a DLQ entry
instantiates. The DLQ is the async-messaging instance of an
error log — observable, auditable, drainable.
- The cross-link is one-directional outward (messaging → errors):
messaging references errors for the errors-as-data pattern;
errors does not back-link to messaging.
## What Violates Delivery-Semantics Discipline
| Violation | Principle |
|-----------|-----------|
| Unstated delivery semantic (at-least-once vs exactly-once guessed) | P4 Delivery Semantics are Explicit |
| Non-idempotent consumer under at-least-once delivery | P3 Consumers are Idempotent |
| Dedup store with no TTL (memory leak; grows without bound) | P6 Backpressure is Bounded |
| Dedup key per-producer (distinct messages in the same window deduped) | P3 Consumers are Idempotent |
| Retry with no backoff (hammers the downstream at full rate) | P6 Backpressure is Bounded |
| Unbounded retry budget (consumer never progresses past the poison) | P5 Dead-Letter Handling is Defined |
| DLQ with no depth alert (poison messages accumulate invisibly) | P10, `domains/observability/metrics` |
| Poison message retried forever (no poison queue, no immediate DLQ) | P5 Dead-Letter Handling is Defined |
| Process-not-idempotent with non-transactional mark (crash in window double-applies) | P3, P4 |
| DLQ entry with no reason/audit metadata (uninspectable failure) | P5, `domains/errors/patterns` |
| Ordering assumption the broker does not provide (FIFO assumed on standard queue) | P2 Ordering is a Property, Not an Assumption |
+307
View File
@@ -0,0 +1,307 @@
# Messaging — First Principles
## 1. The Principles
### P1. Messages are Contracts
A message has an explicit, versioned schema. Producer and consumer
agree on shape before exchange; the schema is the boundary, not a
guess. A schemaless message — a free-form JSON blob the consumer
parses by hope — is a defect: the consumer breaks silently on the
next shape change, and the producer has no contract to evolve
against. This derives from `C1 Correctness` (the exchange must
carry what the parties agreed to) and `C2 Clarity` (the schema
makes the boundary obvious to both sides). This is the
cross-process expression of the contract discipline that
`domains/api/rest` owns for synchronous request/response: the
message schema is to async exchange what the API contract is to
sync exchange. It is distinct from
`domains/concurrency/patterns` Pattern 1 (Message Passing), which
owns the *in-process* channel primitive — here the contract spans
separate systems and survives network failure (D-062). See
`domains/messaging/queues.md` for the queue-flavored application
and `domains/messaging/streams.md` for the durable-log-flavored
application.
### P2. Ordering is a Property, Not an Assumption
Ordering guarantees — per-partition strict, global, or none — are
explicit and documented. "It's FIFO" is a claim that must be backed
by the broker's partitioning contract, not an assumption the
consumer makes and the broker may not honor. A standard queue
delivers in arrival order per receive-node but offers no global
ordering across shards; a FIFO queue delivers strict per-message-
group order but at a latency cost; a partitioned stream delivers
strict per-partition order but only within a partition. Each is a
distinct, declared property. This derives from `C1 Correctness`
(order is a correctness property — a consumer that assumes order
the broker does not provide is wrong) and `C2 Clarity` (the
ordering guarantee is documented, not discovered in production).
This is distinct from in-process ordering, which
`domains/concurrency/patterns` Pattern 1 owns for channels within
one program: messaging ordering survives network failure, broker
restart, and consumer crash-and-retry — a stronger failure model
than thread-local channels (D-062). The `messaging-shared-
subscription` chaos anti-pattern breaches this rule: two consumers
sharing one subscription break per-consumer ordering because the
broker dispatches each message to an arbitrary consumer. See
`domains/messaging/streams.md` for the partition-order contract
and `domains/messaging/delivery-semantics.md` for the interaction
of ordering with the three delivery semantics.
### P3. Consumers are Idempotent
Delivery is at-least-once by default across the network; a
consumer deduplicates via idempotency keys or deterministic
processing. "Exactly-once" is idempotency plus at-least-once, not a
broker guarantee — Jepsen analyses of Kafka, RabbitMQ, and NATS
establish that exactly-once claims require independent
verification, and the durable engineering practice is to make
consumers idempotent under redelivery. A non-idempotent consumer
under at-least-once delivery doubles the effect on every retry; a
non-idempotent consumer under a claimed exactly-once broker is a
bug waiting for the broker's exactly-once invariant to break. This
derives from `C1 Correctness`: correctness under redelivery is the
contract, not a nice-to-have. This parallels
`domains/edge/P5 Edge Operations are Idempotent` (the
cross-partition device-and-cache-flavored analog) and is the
cross-process instance of the retry-safety discipline that
`domains/concurrency/patterns` Pattern 6 (Timeout on Every Block)
implies for in-process retry. It is distinct from in-process
retry because the redelivery comes from the broker across a
network, not from an in-process loop (D-062). See
`domains/messaging/delivery-semantics.md` for the idempotency-key
dedup-store pattern.
### P4. Delivery Semantics are Explicit
At-least-once / at-most-once / exactly-once is a declared choice
per channel, not an emergent behavior. The tradeoff — latency cost,
implementation complexity, operational cost — is made consciously
and documented. At-most-once is fire-and-forget (low latency, lossy);
at-least-once is acked with possible duplication (the default,
requires idempotent consumers per P3); exactly-once is at-least-once
plus idempotency or a transactional two-phase commit (highest cost,
narrowest fit). An unstated semantic is a defect: the consumer
guesses, and the guess is wrong under the first failure. This
derives from `C1 Correctness` (the chosen semantic must hold) and
`C2 Clarity` (the tradeoff is visible to the reader and the
operator). This is the cross-process analog of the explicit-failure-
mode discipline that `domains/errors/patterns` owns for synchronous
code — messaging makes the delivery-mode choice as explicit as an
error-handling choice. See `domains/messaging/queues.md` for the
three-semantics comparison table and `domains/messaging/delivery-
semantics.md` for the correctness properties of each.
### P5. Dead-Letter Handling is Defined
Poison messages — unparseable, repeatedly failing, or exhausting
the retry budget — are routed to a dead-letter queue, not retried
forever or silently dropped. The DLQ is observable and drainable: an
operator can inspect it, replay from it, or discard with audit. An
unbounded retry loop is a livelock: the consumer never makes
progress past the poison message. A silent drop is a correctness
defect: the message vanished with no record. This derives from `C1
Correctness` (poison messages must not livelock the consumer or
silently disappear) and `C5 Reversibility` (the DLQ is the
reversibility mechanism — a dead-lettered message can be reprocessed
after the bug is fixed). This is the cross-process analog of the
bounded-error discipline that `domains/errors/patterns` owns for
synchronous code: a poison message is an error-as-data instance
that must be observable and recoverable, not swallowed. It is
distinct from in-process error handling because the failure spans a
network and a consumer restart (D-062). See
`domains/messaging/delivery-semantics.md` for the dead-letter
strategy comparison table and the DLQ routing rule pattern.
### P6. Backpressure is Bounded
A slow consumer cannot unbounded-buffer the broker or the
producer. Backpressure is explicit: consumer lag is visible,
max-unacked is bounded, the retry budget is capped. A consumer
that falls behind without a visible signal is a silent backlog —
the operator cannot fix what they cannot see, and the broker's
memory grows without bound until it fails. This derives from `C1
Correctness` (a backlog that grows until OOM is a correctness
failure) and `C8 Economy` (the broker's memory is bounded by
design, not by luck). This is distinct from
`domains/concurrency/P9 Bounded Queues` and
`domains/concurrency/patterns` Pattern 5 (Bounded Queue with
Backpressure), which own the *in-process* analog: concurrency's
bounded queue fails by OOM or thread crash; messaging's bounded
backpressure fails by network partition, broker restart, or
consumer crash-and-retry (D-062). The Reactive Streams
specification (`request(n)`, `onNext` bounded) is the in-process
instance; messaging's broker-backed backpressure is the
cross-process instance above it. See `domains/messaging/queues.md`
for prefetch and max-unacked and `domains/observability/metrics`
for consumer-lag as an alert.
### P7. Partitioning is Intentional
The partition key determines ordering, parallelism, and hotspots.
Key choice is a design decision with documented rationale, not a
default. A key that hashes unevenly creates a hot partition that
limits throughput; a key that does not match the ordering need
breaks per-key semantics; a key that is too coarse (one partition
for the whole topic) serializes all the traffic. The partition
count is a capacity bound: too few partitions cap parallelism, too
many partition overhead the broker. This derives from `C4 Locality`
(ordering and parallelism are co-located with the partition) and
`C6 Composability` (the partition is the unit of parallelism and
scaling — consumer groups compose from per-partition workers).
This is the cross-process analog of the locality discipline that
`domains/performance/` owns for generic data-near-compute
optimization: performance's locality is algorithmic (data near
compute); messaging's locality is partitional (order and
parallelism near the partition). See `domains/messaging/streams.md`
for the partitioned-log model and consumer-group rebalance
strategies.
### P8. Replay and Retention are Configured
Retention windows and replay-from-offset are explicit. A message
is not ephemeral by default; the broker is a durable log, not a
pipe. A topic with no retention is a fire-and-forget stream — a
consumer that falls behind loses data permanently; a topic with
infinite retention is an unbounded log — the broker grows until
disk exhaustion. Both are defects: the retention window is a
declared bound, and replay-from-offset is the mechanism that makes
the log durable (re-consumable) rather than ephemeral. This derives
from `C5 Reversibility` (a retained message is reversible — it can
be re-consumed; an ephemeral message is not) and `C7 Observability`
(the durable log is itself an observable record of what happened —
the offset is the position from which to replay). This is the
foundation for `domains/messaging/streams.md` and the rule that
distinguishes a stream from a queue (a queue deletes on ack; a
stream retains for replay). See `domains/messaging/pubsub.md` for
the pub/sub-vs-stream durability boundary.
### P9. Schemas Evolve Compatibly
Schema changes are backward- and forward-compatible by
construction. Breaking changes are versioned migrations, not
silent shape edits. A producer that ships a new field the old
consumer ignores is backward-compatible; a consumer that handles a
missing field the new producer omits is forward-compatible. A
silent schema change — the producer renames a field and the
consumer parses `undefined` — is a P1 violation (the contract was
broken) compounded here as an evolution defect. This derives from
`C5 Reversibility` (a schema change is reversible by versioning —
the old shape is still readable) and `C6 Composability` (producers
and consumers of different versions compose because the schema
evolves compatibly). This parallels `domains/data/migrations`
(schema migration for databases) and `domains/api/versioning`
(API contract evolution): messaging's schema evolution is the
async instance of the same compatibility discipline. See
`domains/messaging/streams.md` for the stream-schema-evolution
angle.
### P10. Messaging is Observable
Consumer lag, DLQ depth, throughput, and consumer-group health are
first-class signals. Silent backlog is a bug, not a feature: a
consumer that falls behind with no lag metric is invisible until
the downstream effect surfaces — by which time the backlog may be
hours or days. A DLQ that grows without an alert is a silent
correctness defect: poison messages are accumulating and no one
knows. This derives from `C7 Observability` (the broker's behavior
is visible to the operator) and `C1 Correctness` (backlog
detection is a correctness bound — unbounded lag is a failure).
This is distinct from `domains/observability/metrics`, which owns
*generic* structured metrics; messaging owns the *broker-specific*
signals — lag, DLQ depth, partition imbalance, consumer-group
rebalance events. See `domains/observability/metrics` for the
generic SLI/SLO discipline and `domains/observability/tracing` for
cross-partition traces.
## 2. Core Principle Trace
Each messaging P-rule derives from one or more core C-rules
(C1C8). The matrix extension lands in P4 of the v0.4 plan; the
traces below are authoritative. Messaging is a broad-derivation
domain touching 7 of 8 core principles (C1, C2, C4, C5, C6, C7,
C8); C3 (Simplicity) is not a primary derivation — messaging is
inherently a tradeoff domain where simplicity yields to the
correctness of delivery guarantees (a simpler-than-necessary
delivery model does not handle the failure cases, per C3's
"simpler than necessary is also a violation").
| P-rule | Core | Why |
|--------|------|-----|
| P1 Messages are Contracts | C1, C2 | Correctness of the exchange; clarity of the schema boundary |
| P2 Ordering is a Property, Not an Assumption | C1, C2 | Correctness of order; clarity of the guarantee |
| P3 Consumers are Idempotent | C1 | Correctness under redelivery |
| P4 Delivery Semantics are Explicit | C1, C2 | Correctness of the chosen semantic; clarity of the tradeoff |
| P5 Dead-Letter Handling is Defined | C1, C5 | Correctness of poison-message routing; reversibility of reprocessing |
| P6 Backpressure is Bounded | C1, C8 | Correctness of bounded backlog; economy of broker memory |
| P7 Partitioning is Intentional | C4, C6 | Locality of order; composability of parallelism |
| P8 Replay and Retention are Configured | C5, C7 | Reversibility of replay; observability of the durable log |
| P9 Schemas Evolve Compatibly | C5, C6 | Reversibility of schema changes; composability of versions |
| P10 Messaging is Observable | C7, C1 | Observability of lag/DLQ; correctness of backlog detection |
## 3. What Violates These Principles
| Violation | Principle Breached |
|-----------|-------------------|
| Schemaless message (no versioned contract; consumer parses by guess) | P1 Messages are Contracts |
| "It's FIFO" with no documented partition contract | P2 Ordering is a Property, Not an Assumption |
| Non-idempotent consumer under at-least-once delivery | P3 Consumers are Idempotent |
| Unstated delivery semantic (at-least-once vs exactly-once guessed) | P4 Delivery Semantics are Explicit |
| No dead-letter queue (poison message retried forever or silently dropped) | P5 Dead-Letter Handling is Defined |
| Unbounded retry budget (no cap; slow consumer stalls the partition) | P6 Backpressure is Bounded |
| Default partition key (no rationale; hotspot or wrong-order) | P7 Partitioning is Intentional |
| Ephemeral broker (no retention; no replay) | P8 Replay and Retention are Configured |
| Silent schema change (producer breaks consumers with no version bump) | P9 Schemas Evolve Compatibly |
| Silent backlog (no lag metric; consumer falls behind invisibly) | P10 Messaging is Observable |
| Shared subscription (two consumers share one subscription; per-consumer ordering breaks) | P2 Ordering is a Property, Not an Assumption (P3 compounding) |
| Blocking consumer (slow downstream call with no timeout; broker redelivers to the stuck consumer) | P6 Backpressure is Bounded |
## 4. Relationship to Other Domains
Messaging systems are the engineering discipline of
**cross-process, cross-system asynchronous communication via
brokers**. Producer and consumer are separate systems; the broker
is the intermediary that brokers delivery, ordering, retention,
and failure semantics. The distinguishing constraints are a
cross-process failure model (network, not crash), explicit
delivery semantics, decoupled producer/consumer lifecycle, and
replay-and-retention as a durable-log property. Messaging overlaps
`domains/concurrency/` by *subject* (messages, queues,
backpressure) but not by *failure model*: per D-062, messaging
owns the cross-process/network-failure-model angle; concurrency
owns the in-process/crash-failure-model angle. The discriminator
is the failure model: concurrency's queue fails by OOM or thread
crash; messaging's queue fails by network partition, broker
restart, or consumer crash-and-retry. Messaging extends
concurrency's bounded-queue/backpressure model to the network-
partition regime. Cross-links are one-directional outward (per
D-026 extended); no back-link edits to v0.1/v0.2/v0.3 content.
- `domains/concurrency/patterns` ← P6 (the broker-backed bounded
queue is the cross-process analog of the in-process bounded
buffer — concurrency Pattern 5 owns in-process; messaging owns
the network-failure-model instance above it, per D-062)
- `domains/concurrency/patterns` ← P3 (idempotent retry is the
cross-process analog of in-process retry-safety — the failure
model differs: broker redelivery across a network vs in-process
loop)
- `domains/observability/metrics` ← P10 (consumer lag and DLQ
depth as alerts; observability owns the generic SLI/SLO
discipline, messaging owns the broker-specific signals)
- `domains/observability/tracing` ← P10 (cross-partition traces
for stream processing; observability owns the generic tracing
discipline, messaging owns the cross-partition propagation)
- `domains/data/schema-design` ← P1, P9 (message schema design
and evolution; data owns the generic schema discipline,
messaging owns the cross-process message-shape instance)
- `domains/errors/patterns` ← P5 (errors as data for message
failures; a poison message is an error-as-data instance that must
be observable and recoverable, not swallowed)
- `domains/edge/iot` ← P4 (the edge↔messaging cross-link
resolves bidirectionally here: edge/iot.md links outward to
messaging/queues for MQTT QoS parallels to delivery semantics;
this first-principles doc acknowledges the back-link — the
edge/iot.md → messaging/queues link from P1 now resolves because
messaging/queues.md exists, completing the bidirectionality per
IDEATE-40)
> Note: the edge/iot.md → messaging/queues cross-link (MQTT QoS
> parallels for delivery semantics) was authored in P1 with a
> dangling reference; this P2 authorship of messaging/queues.md
> resolves it. The bidirectionality is verified in P5
> (ATELIER-114 per IDEATE-40). The cross-link is one-directional
> outward from edge/iot.md; this first-principles doc
> acknowledges the resolution without editing edge/iot.md (per
> D-026 extended — no back-link edits to v0.1/v0.2/v0.3 or to
> P1-authored edge content).
+269
View File
@@ -0,0 +1,269 @@
# Pub/Sub — Derived Rules
> Derives from `domains/messaging/first-principles.md`. Applies P1
> (Messages are Contracts), P2 (Ordering is a Property, Not an
> Assumption), P3 (Consumers are Idempotent), and P4 (Delivery
> Semantics are Explicit) primarily, with P7 (partitioning), P10
> (per-subscription lag). The `messaging-shared-subscription` chaos
> anti-pattern lives here (pre-specified in P4 ATELIER-110).
> Cross-links `domains/messaging/streams` for the pub/sub-vs-stream
> durability boundary and `domains/observability/metrics` for
> per-subscription lag.
## What Pub/Sub Is (P1 Messages are Contracts)
- Pub/sub is the fan-out primitive: a producer publishes a message
to a topic; N independent subscriptions each receive a copy. The
message has an explicit, versioned schema (P1): the topic's
schema is the contract every subscription agrees to before
subscribing. A schemaless topic is a defect — every subscriber
breaks silently on the next shape change.
- The boundary with queues is the fan-out ratio. A queue is
point-to-point (one producer, one consumer); pub/sub is
one-to-many (one producer, N consumers, each with its own
subscription). The boundary with streams is the durability model
— see the cross-link below. Pub/sub is an async concern because
producer and consumers are separate systems and the failure model
is network, not crash (D-062).
- See `domains/messaging/queues.md` for the point-to-point variant
and `domains/messaging/streams.md` for the durable-log variant.
## Topic / Subscription Model (P1, P3, P4)
- A **topic** is the named stream of messages. A **subscription**
is a durable cursor over the topic: each subscription receives
every message published after it was created (subject to
retention and filtering). The subscription is independent — its
ack, redelivery, and DLQ are per-subscription, not shared.
- Each subscription is a consumer under at-least-once by default
(P4): the broker redelivers until the subscription acks, and the
subscriber must be idempotent (P3). A subscription with no
idempotency dedup duplicates every redelivered message.
- The topic's schema evolves compatibly (P9 — Schemas Evolve
Compatibly): a new field the old subscriber ignores is
backward-compatible; a renamed field the old subscriber parses
as `undefined` is a P1 violation.
```python
# Publish + two independent subscriptions (P1 contract, P3
# idempotency, P4 at-least-once per subscription). Each
# subscription is an independent durable cursor; acking one does
# not affect the other.
import json, uuid
# --- Publisher ---
def publish(topic, event, broker):
# P1: versioned schema on the topic. All subscribers must
# understand this schema (or a compatible superset — P9).
message = {
"schema": "user.signed-up.v1",
"id": str(uuid.uuid4()),
"idempotencyKey": f"user:{event['userId']}:signup",
"payload": event,
}
broker.publish(topic=topic, body=json.dumps(message))
# --- Subscription A: welcome-email service ---
def subscribe_welcome(broker, dedup_store, send_email):
sub = broker.subscribe(topic="users", subscription="welcome-email")
for message in sub.receive():
# P3: idempotent per subscription. A redelivered message is
# a no-op for THIS subscription, not for the others.
if dedup_store.seen(("welcome", message["idempotencyKey"])):
sub.ack(message)
continue
send_email(message["payload"]["email"], "Welcome!")
dedup_store.mark(("welcome", message["idempotencyKey"]))
sub.ack(message)
# --- Subscription B: analytics-ingest service ---
def subscribe_analytics(broker, dedup_store, ingest):
# Independent subscription: its own cursor, its own dedup,
# its own ack. Welcome-email acking does NOT advance this.
sub = broker.subscribe(topic="users", subscription="analytics")
for message in sub.receive():
if dedup_store.seen(("analytics", message["idempotencyKey"])):
sub.ack(message)
continue
ingest(message["payload"])
dedup_store.mark(("analytics", message["idempotencyKey"]))
sub.ack(message)
```
- The dedup key is scoped per subscription: `(subscription,
idempotencyKey)`. A redelivery to subscription A that was already
processed by A is a no-op for A; the same message delivered to
subscription B is processed by B independently. Scoping the dedup
key by subscription prevents one subscription's dedup from
masking another's redelivery.
## Fan-Out Semantics (P4, P7)
- Fan-out means every subscription receives every published message
(subject to filtering — see below). The broker duplicates the
message per subscription; each subscription's delivery is
independent. The fan-out ratio is the number of subscriptions; the
broker's cost scales with fan-out × message size.
- Partitioning (P7) applies to topics that are partitioned for
throughput: a partitioned topic delivers per-partition order, and
each subscription receives from every partition. A subscription
that consumes partitions in parallel must handle per-partition
ordering and cross-partition non-ordering (P2 — document the
property, do not assume global order).
- The delivery semantic is per-subscription (P4): subscription A
may be at-least-once, subscription B may be at-most-once (for a
loss-tolerant analytics feed). The choice is per subscription,
declared, not emergent.
## Shared vs Independent Subscriptions (P2, P3 — the chaos anti-pattern)
- An **independent subscription** is one durable cursor per
consumer group: each subscription receives every message in
topic order (per partition, P2) and acks independently. This is
the correct default: per-consumer ordering and per-consumer
idempotency hold.
- A **shared subscription** is one subscription shared by multiple
consumers: the broker dispatches each message to an arbitrary
consumer in the shared group. This breaks per-consumer ordering
(P2 — consumer A sees message 3 before consumer B sees message
1) and complicates idempotency (P3 — the dedup state must be
shared across consumers, not per-consumer). This is the
`messaging-shared-subscription` chaos anti-pattern
(pre-specified in P4 ATELIER-110): the primary breach is P2
(ordering); P3 (idempotency) is the compounding consequence.
- A shared subscription is correct ONLY when the consumers are
stateless, the per-message processing is order-independent, and
the dedup store is shared (a shared Redis, a shared DB). A shared
subscription for order-dependent or per-consumer-stateful
processing is the chaos anti-pattern: the broker's arbitrary
dispatch breaks the order the consumer assumes.
```python
# The messaging-shared-subscription chaos anti-pattern (P2
# ordering breach, P3 idempotency compounding). Two consumers
# share one subscription; the broker dispatches each message to
# an arbitrary consumer. Per-consumer ordering breaks; dedup
# must be shared (and often is not).
# BAD — shared subscription, per-consumer dedup (chaos):
def shared_subscription_bad(broker, send_email):
# Both consumers call subscribe with the SAME subscription
# name. The broker round-robins; consumer A gets msg 1, msg 3;
# consumer B gets msg 2, msg 4. Per-consumer order is broken.
# If each consumer has its OWN dedup store, a redelivery to
# the OTHER consumer re-processes (P3 breach).
sub = broker.subscribe(topic="users", subscription="shared")
for message in sub.receive():
# Per-consumer dedup — WRONG. A redelivered message may
# land on the other consumer, which has not seen it.
if local_dedup.seen(message["idempotencyKey"]): # per-consumer
sub.ack(message); continue
send_email(message["payload"]["email"], "Welcome!")
local_dedup.mark(message["idempotencyKey"])
sub.ack(message)
# CORRECT — independent subscriptions (per-consumer ordering,
# per-subscription dedup):
def independent_subscriptions_good(broker, send_email):
sub = broker.subscribe(topic="users", subscription="welcome-email")
for message in sub.receive():
if dedup_store.seen(("welcome", message["idempotencyKey"])):
sub.ack(message); continue
send_email(message["payload"]["email"], "Welcome!")
dedup_store.mark(("welcome", message["idempotencyKey"]))
sub.ack(message)
```
- If a shared subscription is genuinely required (stateless,
order-independent, shared dedup), document the choice and the
shared-dedup requirement (P2 — the ordering property is "none
across consumers"; P3 — the dedup is shared). The default is
independent subscriptions; shared is an opt-in for the narrow case.
## Filtering (P4, C8 Economy)
- **Subscription filtering** lets a subscription receive only
messages matching a filter (e.g., `event.type == "order"`).
Filtering at the broker saves bandwidth (C8 — the subscriber
does not receive and discard) and reduces subscriber load.
- **Server-side filtering** (broker evaluates the filter before
delivery) is more efficient than **client-side filtering**
(subscriber receives and discards). Server-side filtering is the
default where the broker supports it (GCP Pub/Sub, SNS filtering,
NATS subject filtering); client-side is the fallback.
- A filter that is too broad wastes bandwidth; a filter that is
too narrow drops messages the subscriber needed. The filter is
a P1 (contract) and P4 (semantic) decision: the subscription's
filter is part of its declared contract.
## Ordering Across Subscriptions (P2)
- A topic with per-partition ordering delivers per-partition order
to each subscription. Across subscriptions, there is no ordering
guarantee: subscription A may ack message 3 while subscription B
is still on message 1. This is correct and expected — each
subscription is independent.
- Within a subscription, ordering holds per partition (P2 — the
documented property). A subscription that processes partitions
in parallel must not assume cross-partition order. A subscription
that needs global order must use a single partition (sacrificing
parallelism, P7) or an external sequencing mechanism.
- The `messaging-shared-subscription` anti-pattern breaks even
per-partition order within a subscription: the broker's arbitrary
dispatch to consumers in the shared group breaks the per-
partition sequence each consumer sees.
## Pub/Sub vs Stream — The Durability Boundary (cross-link messaging/streams)
- Pub/sub and streams are both fan-out or one-to-many primitives,
but their durability model differs. Pub/sub is a
**push-to-subscription** model: each subscription is a cursor,
retention is short (the subscription's unacked window), and
replay is limited to the unacked messages. A subscription that
falls behind beyond the retention window loses messages
permanently.
- A stream is a **durable-log** model: messages are retained by
the log for a configured window (P8 — Replay and Retention are
Configured), and any consumer group can replay from any offset
within the window. A stream consumer that falls behind can
catch up by replaying; a pub/sub subscription that falls behind
beyond retention cannot.
- The choice is the durability requirement: if the consumer must
be able to replay (reprocessing, backfill, new consumer starting
from the beginning), use a stream. If the consumer only needs
the live feed (and can tolerate loss on a long fall-behind),
pub/sub is lighter. See `domains/messaging/streams.md` for the
durable-log model, offsets, and consumer groups.
## Observability — Per-Subscription Lag (P10)
- Per-subscription lag (messages published minus messages acked
for each subscription, or the age of the oldest unacked message
per subscription) is the primary pub/sub health signal. Each
subscription has its own lag — a fast subscription and a slow
subscription on the same topic are independent signals.
- A subscription whose lag grows beyond the retention window is a
silent data-loss risk: the broker will drop the oldest messages,
and the subscription will never see them. Alert on lag relative
to retention — lag approaching retention is the loss threshold.
- Wire per-subscription lag to `domains/observability/metrics` as
an SLI per subscription. A topic with N subscriptions has N lag
metrics; a single aggregate hides the slow one. See
`domains/observability/metrics` for the generic SLI/SLO
discipline.
## What Violates Pub/Sub Discipline
| Violation | Principle |
|-----------|-----------|
| Shared subscription for order-dependent processing (broker dispatch breaks per-consumer order) | P2 Ordering is a Property, Not an Assumption |
| Shared subscription with per-consumer dedup (redelivery to the other consumer re-processes) | P3 Consumers are Idempotent |
| Schemaless topic (no versioned contract; subscribers parse by guess) | P1 Messages are Contracts |
| Subscription with no idempotency dedup (redelivered message duplicates the effect) | P3 Consumers are Idempotent |
| Subscription whose lag approaches retention (silent data loss) | P10, `domains/observability/metrics` |
| Unstated delivery semantic per subscription (at-least-once vs at-most-once guessed) | P4 Delivery Semantics are Explicit |
| Filter that is too narrow (drops messages the subscriber needed) | P1, P4 |
| Partitioned topic with no documented per-partition ordering contract | P2 Ordering is a Property, Not an Assumption |
| No per-subscription lag metric (slow subscription invisible) | P10 Messaging is Observable |
| Cross-partition order assumption within a subscription (no global order guarantee) | P2, P7 |
+318
View File
@@ -0,0 +1,318 @@
# Queues — Derived Rules
> Derives from `domains/messaging/first-principles.md`. Applies P1
> (Messages are Contracts), P3 (Consumers are Idempotent), P4
> (Delivery Semantics are Explicit), P5 (Dead-Letter Handling is
> Defined), and P6 (Backpressure is Bounded) primarily, with P2
> (ordering), P7 (partitioning), and P10 (observable lag). For the
> at-least-once / at-most-once / exactly-once decision, see the
> comparison table below. Cross-links `domains/concurrency/patterns`
> for the in-process bounded-queue analog and
> `domains/observability/metrics` for consumer lag.
## What a Queue Is (P1 Messages are Contracts)
- A queue is a point-to-point async delivery primitive. A producer
enqueues a message; exactly one consumer dequeues and processes
it. The message has an explicit, versioned schema (P1): the
producer and consumer agree on shape before exchange, and the
schema is the boundary — a schemaless message is a defect (the
consumer breaks silently on the next shape change).
- The boundary is per D-062: messaging owns the cross-process /
network-failure-model angle; concurrency owns the in-process
analog. A queue is a messaging concern because producer and
consumer are separate systems, the broker is the intermediary,
and the failure model is network (the message can be lost,
duplicated, reordered, or delayed by the broker or the network,
not by a thread crash). The in-process bounded buffer
(`domains/concurrency/patterns` Pattern 5) is the analog below
this boundary — it fails by OOM; a broker-backed queue fails by
partition, broker restart, or consumer crash-and-retry.
- See `domains/messaging/pubsub.md` for the fan-out (one-to-many)
variant and `domains/messaging/streams.md` for the durable-log
(replay-from-offset) variant. A queue deletes on ack; a stream
retains for replay — the durability boundary is the
distinguishing trait.
## Producer / Consumer Model (P1, P3)
- The producer enqueues a message with an idempotency key (P3).
The consumer dequeues, processes, and acks. If the consumer
crashes before acking, the broker redelivers; the idempotency
key makes the redelivery safe (the consumer dedups, not the
broker).
- The idempotency key is per-message, not per-producer or
per-session. A consumer that dedups by producer alone will drop
distinct messages issued in the same window. Use a UUID per
message, or a deterministic key derived from the message content
(e.g., `(entity, operation, version)`).
```python
# Producer/consumer pair with idempotency key (P1 contract, P3
# idempotency). The producer tags each message with a versioned
# schema and a unique idempotency key; the consumer dedups by the
# key so a redelivered message is processed once (P3).
# --- Producer ---
import json, uuid
def enqueue(order, broker):
# P1: versioned schema. The message carries its schema version
# so the consumer can route by shape (P9 evolution discipline).
message = {
"schema": "order.created.v1",
"id": str(uuid.uuid4()),
"idempotencyKey": f"order:{order['id']}:{order['version']}",
"payload": order,
}
broker.send(queue="orders", body=json.dumps(message))
# At-least-once by default (P4): the broker acks the send; the
# consumer may see this message more than once under retry.
# --- Consumer ---
def consume(broker, dedup_store, process_order):
for message in broker.receive(queue="orders"):
# P3: idempotent consumer. Dedup by idempotency key before
# processing; a redelivered message is a no-op, not a
# double-apply.
if dedup_store.seen(message["idempotencyKey"]):
broker.ack(message) # already processed; skip
continue
try:
process_order(message["payload"])
dedup_store.mark(message["idempotencyKey"])
broker.ack(message) # success; broker drops it
except Exception:
broker.nack(message) # redeliver (at-least-once, P4)
```
- The dedup store is bounded (P6 — Backpressure is Bounded): a
dedup store that grows without bound is a memory leak. Use a TTL
window longer than the broker's max-redelivery window, or a
bounded LRU. See `domains/messaging/delivery-semantics.md` for
the full idempotency-key dedup-store pattern.
## Ack / Nack (P4 Delivery Semantics are Explicit)
- **Ack** tells the broker the message was processed; the broker
drops it. **Nack** (negative ack) tells the broker the
processing failed; the broker redelivers (at-least-once) or
routes to a DLQ (after the retry budget — P5).
- A consumer that neither acks nor nacks within the visibility
timeout causes the broker to redeliver (the broker assumes the
consumer died). This is the at-least-once default: the broker
prefers duplication to loss.
- The semantic is explicit (P4): at-least-once is the default; the
consumer must be idempotent (P3). At-most-once is fire-and-forget
(no ack; the broker drops on send) — lossy but lowest latency.
Exactly-once is at-least-once plus idempotency, or a
transactional two-phase commit — see the comparison table below.
## Visibility Timeouts and Redelivery (P4, P5)
- The visibility timeout is the window the broker hides a message
after delivery, waiting for the ack. If the consumer does not
ack within the window, the broker makes the message visible
again and redelivers it (to the same consumer or another). This
is the at-least-once mechanism: the broker assumes a
no-ack-in-time consumer is dead.
- The timeout must be longer than the processing time, or the
broker redelivers a message the consumer is still processing —
causing duplicate processing (which P3 idempotency makes safe,
but which wastes resources). A timeout shorter than processing
time is a P6 (backpressure) smell: the consumer is too slow for
the configured timeout.
- Redelivery has a budget (P5): after N redeliveries or a TTL, the
message routes to the DLQ. An unbounded retry budget is the
`messaging-unbounded-retry` chaos anti-pattern: the consumer
never makes progress past the poison message.
## FIFO vs Standard Queues (P2 Ordering is a Property, Not an Assumption)
- A **standard queue** delivers in arrival order per receive-node
but offers no global ordering across shards, no per-message-
group ordering, and may redeliver out of order under retry. It
is the high-throughput default; ordering is *not* guaranteed
(P2: the ordering property is "none" — explicitly documented).
- A **FIFO queue** delivers strict per-message-group order: all
messages with the same group ID are delivered to one consumer
in send order. The cost is throughput (FIFO queues cap at lower
TPS) and latency (the broker must sequence per group). The
ordering property is "per-group strict" — explicitly documented
(P2).
- The choice is a P2 decision (which ordering guarantee) and a C8
decision (throughput cost). A consumer that assumes FIFO on a
standard queue is a P2 violation: the broker does not provide
the guarantee the consumer assumes. Document the property; do
not assume it.
## Prefetch and Concurrency (P6 Backpressure is Bounded)
- **Prefetch** (or max-unacked) bounds how many messages the
broker delivers to one consumer without an ack. A prefetch of 1
is strict stop-and-wait (lowest throughput, tightest backpressure);
a prefetch of N allows the consumer to process N in flight
(higher throughput, more memory). An unbounded prefetch is a P6
violation: the broker floods the consumer's memory.
- **Consumer concurrency** is the number of parallel workers
processing from the queue. More workers increase throughput up to
the downstream's limit; beyond that, the workers stall the
downstream (P6 — the backpressure propagates to the
downstream, not the broker).
- The prefetch × concurrency product is the in-flight cap. Declare
it (P6): an undeclared cap is a defect — the consumer either
underutilizes the broker (prefetch too low) or OOMs under load
(prefetch too high). This is the cross-process analog of
`domains/concurrency/patterns` Pattern 5 (bounded queue with
backpressure): concurrency owns the in-process analog;
messaging owns the broker-backed instance.
## Long Polling (P6, C8 Economy)
- Long polling (or `ReceiveMessage` with a wait-time-seconds)
holds the receive request open until a message arrives or the
wait expires. This reduces empty-receive round trips (C8
economy of the constrained link) and reduces latency-to-first-
message (the message is delivered when it arrives, not on the
next poll cycle).
- Long polling is the default for low-throughput queues: short
polling burns CPU on empty receives; long polling waits for
work. For high-throughput queues, the broker is usually full
enough that long polling adds no latency; for low-throughput
queues, long polling is the difference between 20ms and 20s
latency-to-first-message.
## Redelivery + DLQ Flow (P5 Dead-Letter Handling is Defined)
- A poison message (unparseable, repeatedly failing, or exhausting
the retry budget) routes to the dead-letter queue. The DLQ is
observable (P10 — DLQ depth is an alert) and drainable (an
operator can inspect, replay, or discard with audit).
- The retry budget is bounded (P6): N redeliveries, or a TTL with
exponential backoff. After the budget is exhausted, the message
is moved to the DLQ, not retried forever. An unbounded retry is
the `messaging-unbounded-retry` chaos anti-pattern (P5 breach).
- The DLQ routing rule is a redelivery-count or TTL threshold
plus a target queue. See `domains/messaging/delivery-semantics.md`
for the dead-letter strategy comparison table.
```python
# Redelivery + DLQ flow (P5 dead-letter handling, P6 bounded
# retry budget). The consumer tracks redelivery count; after the
# budget, the message routes to the DLQ. The DLQ is observable
# (P10 — depth is an alert) and drainable.
MAX_REDELIVERIES = 5
DLQ = "orders-dlq"
def consume_with_dlq(broker, dedup_store, process_order):
for message in broker.receive(queue="orders"):
# P3 idempotency: a redelivered, already-processed message
# is acked and skipped (not re-processed, not DLQ'd).
if dedup_store.seen(message["idempotencyKey"]):
broker.ack(message)
continue
try:
process_order(message["payload"])
dedup_store.mark(message["idempotencyKey"])
broker.ack(message)
except Exception as exc:
# P5: bounded retry budget. After MAX_REDELIVERIES,
# route to DLQ — do NOT retry forever.
count = message.get("redeliveryCount", 0) + 1
if count >= MAX_REDELIVERIES:
broker.send(DLQ, body=json.dumps({
"original": message,
"reason": str(exc),
"deadLetteredAt": now_iso(),
"redeliveryCount": count,
}))
broker.ack(message) # remove from the origin queue
# P10: the DLQ depth must alert. A DLQ that grows
# with no alert is a silent correctness defect.
else:
# Nack with backoff: the broker redelivers after a
# delay. The backoff caps the retry rate (P6).
broker.nack(message, delay=exponential_backoff(count))
```
- The `deadLetteredAt` and `reason` fields make the DLQ entry
observable and auditable: an operator inspecting the DLQ sees
why each message was dead-lettered and when. See
`domains/errors/patterns` for the errors-as-data discipline the
DLQ entry follows.
## Delivery Semantics Comparison (D-069)
| Semantic | Guarantee | Latency Cost | Implementation Cost | When It Fits |
|----------|-----------|-------------|---------------------|--------------|
| **At-most-once** | A message is delivered 0 or 1 times; loss is possible, duplication is not | Lowest (no ack; fire-and-forget) | Lowest (no ack, no dedup) | Telemetry where a dropped sample is acceptable; high-throughput metrics; MQTT QoS 0; logs where a lost line is tolerable. Never for billing, orders, or any side-effecting operation. |
| **At-least-once** | A message is delivered 1 or more times; duplication is possible, loss is not | Low (one ack round-trip) | Medium (consumer must be idempotent — P3; dedup store required) | The default for side-effecting operations: orders, payments, commands. The consumer dedups via idempotency keys (P3); the broker guarantees delivery. Fits the vast majority of broker-backed queues (SQS standard, RabbitMQ ack, MQTT QoS 1). |
| **Exactly-once** | A message is delivered exactly 1 time; no loss, no duplication | Highest (two-phase commit or transactional producer+consumer) | Highest (requires transactions, a transactional producer, and a transactional consumer — or at-least-once plus idempotency, which collapses to at-least-once with dedup) | Rare. Kafka transactions (consume-process-produce in one transaction); MQTT QoS 2 (four-step handshake). In practice, "exactly-once" is usually at-least-once plus idempotency (P3) — the broker does not guarantee it; the consumer enforces it. Jepsen analyses verify broker claims. |
- The default for side-effecting operations is **at-least-once with
idempotent consumers** (P3). At-most-once is for loss-tolerant
telemetry. Exactly-once is reserved for the narrow case where
the consume-process-produce loop must be transactional (Kafka
transactions) — and even then, the consumer should be idempotent
as defense-in-depth.
- The latency cost column is the C8 tradeoff: at-most-once is
cheapest, exactly-once is most expensive. The implementation
cost column is the C1/C3 tradeoff: at-most-once is simplest,
exactly-once is most complex (and most fragile — a transactional
consumer that partially fails is a bug source). The "when it
fits" column is the P4 decision: declare the semantic per
channel, do not let it emerge.
- This is the decision matrix required by D-069 for queues; see
`domains/messaging/delivery-semantics.md` for the correctness
properties of each semantic and the dead-letter strategy
comparison.
## Observability — Consumer Lag (P10 Messaging is Observable)
- Consumer lag (messages enqueued minus messages acked, or the
age of the oldest unacked message) is the primary queue health
signal. A lag that grows without bound is a P6 violation (the
consumer is slower than the producer) and a P10 violation if it
is not alerted.
- DLQ depth is the secondary signal: a DLQ that grows is a P5
signal (poison messages are accumulating) and a P10 signal if
not alerted. Wire both to `domains/observability/metrics` as
SLIs with SLOs (e.g., lag < 1000 messages, DLQ depth < 10).
- A queue with no lag metric is operating blind (P10 violation):
the operator cannot see the consumer falling behind until the
downstream effect surfaces.
## Cross-Link to Concurrency (P6, cross-link concurrency/patterns)
- The broker-backed queue is the cross-process analog of the
in-process bounded buffer. `domains/concurrency/patterns`
Pattern 5 (Bounded Queue with Backpressure) owns the in-process
instance ("producer is blocked or signaled" within one program);
messaging owns the broker-backed instance above it (the producer
is the broker's enqueue, the consumer is the broker's dequeue, the
backpressure is the prefetch cap and the lag signal). The
failure model differs: in-process fails by OOM; broker-backed
fails by network partition, broker restart, or consumer
crash-and-retry (D-062).
- The cross-link is one-directional outward (messaging →
concurrency) per D-026 extended: messaging references concurrency
as the in-process foundation; concurrency does not back-link to
messaging.
## What Violates Queue Discipline
| Violation | Principle |
|-----------|-----------|
| Schemaless message (no versioned contract; consumer parses by guess) | P1 Messages are Contracts |
| Non-idempotent consumer under at-least-once delivery (redelivery doubles the effect) | P3 Consumers are Idempotent |
| Unstated delivery semantic (at-least-once vs exactly-once guessed) | P4 Delivery Semantics are Explicit |
| No DLQ (poison message retried forever or silently dropped) | P5 Dead-Letter Handling is Defined |
| Unbounded prefetch (broker floods consumer memory) | P6 Backpressure is Bounded |
| Unbounded retry budget (no cap; consumer never progresses past the poison) | P5, P6 |
| Consumer that assumes FIFO on a standard queue (ordering not guaranteed) | P2 Ordering is a Property, Not an Assumption |
| Visibility timeout shorter than processing time (redeliver while still processing) | P4, P6 |
| DLQ with no depth alert (poison messages accumulate invisibly) | P10, `domains/observability/metrics` |
| No consumer-lag metric (consumer falls behind invisibly) | P10, `domains/observability/metrics` |
| Dedup store that grows without bound (memory leak) | P6 Backpressure is Bounded |
| Default prefetch with no rationale (underutilizes or OOMs) | P6, `domains/concurrency/patterns` |
+329
View File
@@ -0,0 +1,329 @@
# Streams — Derived Rules
> Derives from `domains/messaging/first-principles.md`. Applies P8
> (Replay and Retention are Configured) primarily, with P2
> (per-partition ordering), P3 (idempotent consumers), P4 (exactly-
> once via transactions), P7 (partitioning), and P10 (observable
> consumer-group health). For the Kafka / Kinesis / Pulsar /
> NATS JetStream decision, see the stream-platform comparison
> table below. For the consumer-group rebalance strategy choice,
> see the rebalance enumeration below (IDEATE-41). Cross-links
> `domains/messaging/delivery-semantics` for exactly-once via
> transactions, `domains/data/schema-design` for stream schema,
> and `domains/observability/tracing` for cross-partition traces.
## What a Stream Is (P8 Replay and Retention are Configured)
- A stream is a durable-log messaging primitive. Messages are
appended to a partitioned, replicated log; consumers read from an
offset and advance at their own pace. The log is retained for a
configured window (P8) — a stream is a durable log, not a pipe.
A consumer that falls behind can catch up by replaying from an
earlier offset; a consumer that starts fresh can replay from the
beginning (within retention).
- The boundary with queues and pub/sub is the durability model. A
queue deletes on ack; a pub/sub subscription retains only its
unacked window; a stream retains the whole log for the configured
retention. This makes a stream replayable (P8 — the
reversibility mechanism) and observable as a record (P10 — the
log is itself an audit of what happened). See
`domains/messaging/pubsub.md` for the pub/sub-vs-stream
durability boundary discussion.
- The boundary with concurrency is per D-062: messaging owns the
cross-process/network-failure-model angle. A stream is a
messaging concern because the log spans brokers and consumers
across a network, and the failure model is partition, broker
restart, or consumer crash-and-retry — not in-process OOM or
thread crash.
## Partitioned Log Model (P2, P7)
- A stream is partitioned for throughput and parallelism. Each
partition is an ordered, append-only log; messages within a
partition are strictly ordered (P2 — per-partition strict order
is the documented property). Across partitions, there is no
ordering guarantee: partition 0 and partition 1 are independent
logs.
- The partition key (P7 — Partitioning is Intentional) determines
which partition a message lands on. A key that hashes evenly
spreads load; a key that matches the per-entity ordering need
(e.g., `userId` for user events) keeps a user's events on one
partition in order; a key that is too coarse (one partition for
the whole topic) serializes all traffic.
- The partition count is a capacity bound: it caps the parallelism
(one consumer per partition per consumer group) and the
throughput (each partition has a write-throughput limit). Too
few partitions cap parallelism; too many partition overhead the
broker (file handles, replication, rebalance cost). The choice
is documented (P7), not defaulted.
## Offsets (P2, P8)
- An offset is a consumer's position in a partition. The consumer
reads from its last committed offset; acking (committing the
offset) advances it. A consumer that crashes before committing
re-reads from the last committed offset (at-least-once by
default, P4) — the consumer must be idempotent (P3).
- The offset is per-partition (P2): each partition has its own
position, and the consumer commits them independently (or
atomically across partitions in a transaction — see below).
- Replay (P8) is resetting the offset backward: a consumer can
replay from the beginning of retention, from a timestamp, or
from a specific offset. This is the durable-log property that
distinguishes a stream from a queue.
## Consumer Groups (P3, P7, P10)
- A consumer group is a set of consumers sharing the stream's
partitions: each partition is assigned to exactly one consumer
in the group. The group is the unit of parallelism and the unit
of offset tracking. Within a group, each consumer handles its
assigned partitions; across groups, each group independently
reads the whole stream (the pub/sub fan-out property, per
subscription/group).
- A consumer in the group is idempotent (P3): under at-least-once
(the default), a redelivery after a crash-and-retry re-processes
messages. The consumer dedups by idempotency key, or processes
deterministically (e.g., a stateful aggregation that overwrites
with the latest value).
- The consumer group's health is observable (P10): per-partition
lag (offset of the consumer vs the log's head), the group's
consumption rate, and rebalance events are first-class signals.
A group whose lag grows without bound is a P6 (backpressure)
smell and a P10 (observability) violation if not alerted.
```python
# Consumer-group reading from offsets (P2 per-partition order,
# P3 idempotent under at-least-once, P7 partition assignment,
# P8 replay from offset). Each consumer in the group handles its
# assigned partitions; the group commits offsets atomically or
# per-partition.
def consume_stream(stream, group, dedup_store, process_event):
# Assign partitions to this consumer by the group's
# rebalance strategy (see the enumeration below).
for partition in stream.assigned_partitions(group, consumer=ME):
# Read from the last committed offset (P8 — replay by
# resetting this offset).
offset = stream.committed_offset(group, partition)
for message in stream.read(partition, from_offset=offset):
# P3: idempotent under at-least-once. A redelivery
# after a crash-and-retry re-processes; dedup by key.
if dedup_store.seen(message["idempotencyKey"]):
stream.commit(group, partition, message["offset"])
continue
process_event(message["payload"])
dedup_store.mark(message["idempotencyKey"])
# Commit the offset to advance (P8 — the position is
# the replay pointer).
stream.commit(group, partition, message["offset"])
```
- The commit-after-process order gives at-least-once (a crash
before commit re-reads); the commit-before-process order gives
at-most-once (a crash after commit loses the unprocessed
message). The default is at-least-once with idempotent consumers
(P3, P4).
## Consumer-Group Rebalance Strategies (IDEATE-41, ATELIER-100 refinement)
When a consumer joins or leaves the group, the broker must
reassign partitions. The rebalance strategy determines the cost
and the use-case fit. This enumeration parallels the v0.3
IDEATE-30 drift-type enumeration (each strategy with its
stop-the-world cost and use-case fit).
| Strategy | Mechanism | Partition Stop-the-World Cost | Use-Case Fit |
|----------|-----------|-------------------------------|--------------|
| **Eager rebalance** (stop-the-world) | Every consumer in the group revokes ALL its partitions, the broker reassigns the full partition set, then consumers resume. Every rebalance pauses the whole group. | High — every partition pauses for every rebalance; the whole group stops processing during the revocation+reassignment window. Throughput drops to zero during rebalance. | Simple brokers, small groups, or rarely-rebalancing groups where the simplicity of full revocation outweighs the pause cost. Kafka's legacy protocol (pre-2.4). Avoid for large groups or frequent scale events. |
| **Sticky (incremental cooperative) rebalance** | The broker reassigns only the partitions that must move (the joining/leaving consumer's share); existing partitions stay assigned. The rebalance is incremental and cooperative — no full revocation. | Low — only the moving partitions pause; the rest of the group continues processing. The pause is proportional to the changed partition count, not the total. | The default for large groups, frequent scale events, and rolling deploys. Kafka's CooperativeStickyAssignor (2.4+), Pulsar, NATS JetStream. Prefer for any group where a full stop-the-world on every deploy is unacceptable. |
| **Cooperative (no-revoke) rebalance** | A subset of sticky where no partition is revoked unless the consumer leaves; only additions are incremental. The strictest minimization of stop-the-world. | Lowest — only added partitions pause; existing assignments are untouched. | Groups where partition assignment is append-only (consumers join but rarely leave). Useful for long-lived consumers with incremental scaling. |
- The default for any non-trivial group is **sticky/cooperative**:
a rolling deploy that triggers an eager rebalance pauses the
whole group on every pod restart, which is unacceptable at
scale. The eager strategy is a legacy default that survives
because it is simple; prefer sticky where the broker supports
it.
- The "partition stop-the-world cost" column is the P6
(backpressure) angle: a full stop-the-world during rebalance
causes lag to spike (the consumer is paused, the producer is
not). Sticky rebalance bounds the spike to the moving partitions.
- A rebalance that pauses without a lag alert is a P10 violation:
the operator cannot see the rebalance-induced lag. Wire rebalance
events to `domains/observability/metrics` as an event signal.
## Replay and Retention Windows (P8, C5 Reversibility)
- Retention is the configured window the log keeps messages: time-
based (e.g., 7 days), size-based (e.g., 10 GB per partition), or
compacted (keep the latest value per key — a changelog). A
stream with no retention is a pipe, not a log (P8 violation); a
stream with infinite retention grows until disk exhaustion (P6
violation — backpressure on the broker).
- Replay (P8) is resetting a consumer's offset to re-read from
within the retention window. Use cases: reprocessing after a
consumer bug (replay from the timestamp of the buggy deploy),
backfilling a new consumer (replay from the beginning), or
reindexing (replay to rebuild a derived store).
- Compacted topics (Kafka log-compaction, Pulsar compaction) keep
the latest value per key and discard older values for the same
key. This turns the log into a changelog — a durable
materialized view that replays to the current state. Compaction
is a P8 (retention) and C5 (reversibility) mechanism: the log
retains the current state per key and is replayable to it.
## Stream Processing (P2, P3, P4)
- Stream processing is computing over the stream as it arrives:
windowing (tumbling, sliding, session windows), joins (stream-
stream, stream-table), aggregations (count, sum, per-key
windows), and stateful transformations. The processing is
per-partition ordered (P2 — a window over a partition is
deterministic; a window across partitions is not unless the
window is global).
- Stream processing consumers are idempotent (P3): a redelivery
after a crash re-processes a window; the aggregation must
tolerate re-application (e.g., a sum is idempotent under replay
if the window is keyed by offset range, not by wall time).
- Exactly-once stream processing (P4) requires transactions: the
consume-process-produce loop is one transaction — the input
offset commit and the output produce are atomic. See the
transactional exactly-once producer below.
## Exactly-Once via Transactions (P4, P3)
- Exactly-once stream processing is at-least-once plus a
transaction: the consumer commits the input offset and produces
the output in one transactional operation. If the consumer
crashes mid-transaction, neither the offset commit nor the
output produce happens — the consumer re-reads from the last
committed offset and re-processes (at-least-once), but the
transaction ensures the output is produced exactly once.
- This is NOT a broker guarantee of exactly-once delivery; it is
at-least-once delivery plus idempotent/transactional processing
(P3, P4). Jepsen analyses of Kafka transactions confirm the
boundaries: the transaction is atomic within the broker, but
the downstream sink must be transactional or idempotent too.
- See `domains/messaging/delivery-semantics.md` for the full
exactly-once-via-idempotency discussion.
```python
# Transactional exactly-once producer (P4 exactly-once via
# transactions, P3 idempotent produce). The consume-process-
# produce loop is one transaction: the input offset commit and
# the output produce are atomic. A crash mid-transaction rolls
# both back; the consumer re-reads and re-processes.
def consume_transform_produce(stream, group, txn_producer):
# Begin a transaction. All produces and the offset commit in
# this block are atomic (P4).
with txn_producer.transaction() as txn:
for partition in stream.assigned_partitions(group, ME):
offset = stream.committed_offset(group, partition)
for message in stream.read(partition, from_offset=offset):
output = transform(message["payload"])
# P3: idempotent produce. The txn producer
# dedups by an epoch+sequence so a retried
# transaction does not double-produce.
txn.produce(
topic="enriched-events",
key=message["key"],
value=output,
idempotencyKey=message["idempotencyKey"],
)
# Commit the input offset within the same
# transaction (P4 atomicity). A crash before
# txn.commit() rolls this back; the consumer
# re-reads from the prior offset.
txn.commit_offset(group, partition, message["offset"])
# txn.commit() makes the produces and the offset commit
# visible atomically. A crash before this point aborts
# both; a crash after is safe (idempotent produce — P3).
```
- The transactional producer's idempotency (P3) is the defense
against a retried transaction: the broker dedups the output by
the producer's epoch and sequence so a re-commit does not
double-produce. The transaction (P4) is the defense against a
partial failure: the offset and the output commit together.
## Stream Schema (P1, P9, cross-link data/schema-design)
- A stream's messages carry a versioned schema (P1). The schema
evolves compatibly (P9): a new field the old consumer ignores is
backward-compatible; a renamed field the old consumer parses as
`undefined` is a P1 violation compounded as an evolution defect.
- Stream schemas are often registered in a schema registry
(Confluent, Apicurio) that enforces compatibility on produce.
A producer that tries to publish an incompatible schema is
rejected; the registry is the P1/P9 enforcement point.
- See `domains/data/schema-design` for the generic schema-design
discipline (Avro, Protobuf, JSON Schema); messaging owns the
stream-specific instance — the registry, the per-topic
compatibility mode, the consumer-side routing by schema
version.
## Stream-Platform Comparison (D-069)
| Axis | Apache Kafka | AWS Kinesis | Apache Pulsar | NATS JetStream |
|------|--------------|-------------|---------------|----------------|
| **Ordering** | Per-partition strict (P2); global only via single-partition topic | Per-shard strict; global only via single shard | Per-partition strict; global via single partition; also supports shared (out-of-order) subscriptions | Per-stream strict; per-subject ordering; global via single stream |
| **Partitioning model** | Partitions (immutable count post-creation; increase requires recreate); key→partition by hash | Shards (reshardable: split/merge at runtime); key→shard by hash | Partitions (resizable; Pulsar's layered architecture separates compute from storage); key→partition by hash | Streams (subject-based; republish to resize); key→stream by subject |
| **Replay / retention** | Time- or size-based retention; compaction (latest-per-key); replay from offset or timestamp | Time-based retention (24h365d); replay from sequence number or timestamp; no compaction | Time- or size-based; compaction; replay from offset or timestamp; tiered storage (hot/warm/cold) | Time- or size-based; per-stream max-age; replay from sequence; no native compaction |
| **Consumer groups** | Group-coordinated; offsets stored in an internal topic; eager (legacy) or sticky/cooperative (2.4+) rebalance | Enhanced fan-out consumers (per-shard HTTP/2 push); KCL for group coordination; no native group rebalance (shard is the unit) | Group-coordinated; shared or failover subscription modes; cooperative rebalance | Per-stream consumers; durable cursors; no native group rebalance (stream is the unit) |
| **Exactly-once** | Transactions (KIP-98): consume-process-produce atomic; idempotent producer (KIP-516) | No native exactly-once; at-least-once with consumer-side dedup (P3) | Transactions: produce-ack atomic; idempotent producer | At-least-once by default; dedup window per stream (P3 idempotency) |
| **Use-case fit** | High-throughput durable logs, stream processing (Kafka Streams, Flink), event sourcing, multi-consumer replay | AWS-native streaming, log ingestion, simple ETL within AWS; low operational burden | Cloud-native, geo-replication, tiered storage, mixed pub/sub + streaming; multi-tenant | Lightweight, low-latency, edge-friendly; NATS ecosystem; simpler ops than Kafka |
| **Watch out for** | Partition count is fixed at creation (resize requires recreate + republish); rebalance cost on large groups; operational complexity | Shard limits per account; no compaction; retention cap at 365 days; AWS lock-in | Two-arch (BookKeeper + Brokers) operational complexity; smaller ecosystem | Smaller ecosystem; no native compaction; fewer stream-processing libraries |
- The default for high-throughput durable logs with multi-consumer
replay is **Kafka**; for AWS-native streaming, **Kinesis**; for
geo-replicated multi-tenant or mixed pub/sub + streaming,
**Pulsar**; for lightweight low-latency edge-friendly streaming,
**NATS JetStream** (which also cross-links `domains/edge/iot`
via MQTT parallels — see the edge↔messaging bidirectionality in
`domains/messaging/first-principles.md` §4).
- The ordering column is the P2 check: every platform provides
per-partition/per-shard strict order; none provides global order
across partitions except by single-partition. The replay/
retention column is the P8 check: every platform retains for a
configured window; replay is from offset or timestamp. The
exactly-once column is the P4 check: Kafka and Pulsar provide
transactions; Kinesis and JetStream rely on at-least-once plus
consumer-side idempotency (P3).
## Cross-Partition Traces (P10, cross-link observability/tracing)
- A stream-processing pipeline that fans out across partitions
must propagate a trace context per event: the trace ID follows
the event from source to processed output, even as the event
crosses partition boundaries. Without cross-partion traces, a
downstream error cannot be traced back to its source event.
- See `domains/observability/tracing` for the generic distributed-
tracing discipline (trace context propagation, span
correlation). Messaging owns the stream-specific instance: the
trace context is a message header, the span boundary is the
consume-process-produce edge, and the cross-partition
correlation is by trace ID (not by partition — partitions are
independent logs, P2).
- A stream processor with no trace propagation is a P10
violation: the operator cannot trace a processed event back to
its source. Wire the trace context into every produce and every
consume.
## What Violates Stream Discipline
| Violation | Principle |
|-----------|-----------|
| Stream with no retention (pipe, not log; no replay) | P8 Replay and Retention are Configured |
| Infinite retention (grows until disk exhaustion) | P8, P6 |
| Consumer group with eager rebalance at scale (full stop-the-world per deploy) | P6, IDEATE-41 (use sticky/cooperative) |
| Non-idempotent stream consumer under at-least-once (redelivery re-processes the window) | P3 Consumers are Idempotent |
| Default partition key (no rationale; hotspot or wrong-order) | P7 Partitioning is Intentional |
| Partition count too low (caps parallelism) or too high (overhead) | P7 |
| Cross-partition order assumption (no global order guarantee) | P2 Ordering is a Property, Not an Assumption |
| Exactly-once claimed without transactional consume-process-produce (P4 violation) | P4 Delivery Semantics are Explicit |
| Stream schema with no registry / no compatibility enforcement (silent shape break) | P1, P9, `domains/data/schema-design` |
| No per-partition lag metric (consumer falls behind invisibly) | P10, `domains/observability/metrics` |
| No cross-partition trace propagation (downstream error untraceable) | P10, `domains/observability/tracing` |
| Compacted topic treated as a full log (old values already discarded) | P8, C1 (compaction is a retention mode, not a full log) |