diff --git a/domains/messaging/delivery-semantics.md b/domains/messaging/delivery-semantics.md new file mode 100644 index 0000000..17f90c3 --- /dev/null +++ b/domains/messaging/delivery-semantics.md @@ -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 | \ No newline at end of file diff --git a/domains/messaging/first-principles.md b/domains/messaging/first-principles.md new file mode 100644 index 0000000..6e0ae80 --- /dev/null +++ b/domains/messaging/first-principles.md @@ -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 +(C1–C8). 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). \ No newline at end of file diff --git a/domains/messaging/pubsub.md b/domains/messaging/pubsub.md new file mode 100644 index 0000000..b83bccd --- /dev/null +++ b/domains/messaging/pubsub.md @@ -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 | \ No newline at end of file diff --git a/domains/messaging/queues.md b/domains/messaging/queues.md new file mode 100644 index 0000000..add02e9 --- /dev/null +++ b/domains/messaging/queues.md @@ -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` | \ No newline at end of file diff --git a/domains/messaging/streams.md b/domains/messaging/streams.md new file mode 100644 index 0000000..79f016e --- /dev/null +++ b/domains/messaging/streams.md @@ -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 (24h–365d); 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) | \ No newline at end of file