Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| f61dffbb5a | |||
| f7f007dce8 |
@@ -1,14 +1,15 @@
|
||||
{
|
||||
"phase": 1,
|
||||
"stage": "execute",
|
||||
"stage": "complete",
|
||||
"milestone": "v0.4",
|
||||
"phase_role": "execution",
|
||||
"project": "atelier",
|
||||
"attempts": 0,
|
||||
"updated_at": "2026-08-05T06:00:00Z",
|
||||
"updated_at": "2026-08-05T06:15:00Z",
|
||||
"milestone_complete": false,
|
||||
"milestone_branch": "milestone/v0.4-edge-quantum-langs",
|
||||
"phase_branch": "phase/01-edge",
|
||||
"tag_base": "v0.3",
|
||||
"phase_tag": "v0.3.1"
|
||||
"phase_tag": "v0.3.1",
|
||||
"release_id": 479
|
||||
}
|
||||
@@ -131,7 +131,7 @@ NFR milestone: no separate minor tag. The final patch (v0.2.6) IS the v0.3 deliv
|
||||
| Phase | Name | Type | Status | Key Deliverables |
|
||||
|-------|------|------|--------|------------------|
|
||||
| 0 | Pre-Execution | docs | complete | Spec, clarify, research, ideate, plan, PERSONAS.md (adds edge-engineer + languages-engineer phase-specific personas) — shipped v0.3.0 |
|
||||
| 1 | Edge Domain | docs | pending | domains/edge/{first-principles, cdn, offline-first, iot, sync}.md |
|
||||
| 1 | Edge Domain | docs | complete | domains/edge/{first-principles, cdn, offline-first, iot, sync}.md — shipped v0.3.1 |
|
||||
| 2 | Messaging Domain | docs | pending | domains/messaging/{first-principles, queues, pubsub, streams, delivery-semantics}.md |
|
||||
| 3 | Language-Derived Extension | docs | pending | languages/ × 4 → first-principles + 4 derived docs each (16 derived docs) |
|
||||
| 4 | Matrix + Review Integration | docs | pending | matrix/principles-matrix.md (+20 mappings, 170→190), matrix/domain-coverage.md (+ languages/ sub-table), review/{agent-checklist, peer-review-checklist, anti-patterns}.md, MANIFEST.md (languages/ section) |
|
||||
|
||||
@@ -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 |
|
||||
@@ -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).
|
||||
@@ -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 |
|
||||
@@ -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` |
|
||||
@@ -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) |
|
||||
Reference in New Issue
Block a user