ADR-015: streams depth — carried-metadata keys, global-FIFO ordering, StreamEvent, trim_to (OQ-12 resolved)

This commit is contained in:
glm-5.3-flash committed 2026-10-05 14:15:28 +00:00
1 parent 04a04651dc
commit 8c4ec48f92
12 files changed
+562 -110

No files matched your search

@@ -59,7 +59,7 @@ coalesced-but-complete re-read on SQLite).
| Mechanism | Durability | Replay | Atomicity | Guarantee |
|---|---|---|---|---|
| `notify` / listen | none | never | commit-atomic (delivers at commit; rollback drops) | fire-and-forget, at-most-once per listener session |
| streams | durable table row | yes, per-consumer offset, replay-on-attach default | publish is commit-atomic | every committed event readable by every consumer that hasn't passed its offset |
| streams | durable table row | yes, per-consumer offset, replay-on-attach default | publish is commit-atomic | every committed event readable by every consumer that hasn't passed its offset; **global FIFO by offset within each stream** (read paths yield `offset ASC`; offsets immutable, never renumbered — [ADR-015](015-streams-depth.md) §4) |
| queues | durable row | claim/ack model | enqueue is commit-atomic | at-least-once *work* with visibility timeouts |
- The trait **does not promise replay under `listen()`** — on either
@@ -53,7 +53,11 @@ implement identically):
shape, the ordering guarantee row, and log retention are OQ-12's —
the method *names* are pinned here, their depth was left
un-pinned, mirrored in [core-contract.md](../core-contract.md)
streams.)*
streams. **Resolved 2026-10-05 by [ADR-015](015-streams-depth.md)** —
key = carried metadata, global-FIFO ordering row, event shape
pinned, `trim_to` added; `publish_with_key_tx` also joins the
tx seam (§2 there); the depth was completed in place,
pre-implementation.)*
- queues v1 skeleton — `enqueue`, `claim_one`, `claim_batch`,
`ack_batch`, `cancel`, `get_job`, `sweep_expired`; job handle
`ack / retry / fail / heartbeat`; `EnqueueOpts { delay, run_at,
@@ -114,7 +118,8 @@ a **`TxHandle` trait carrying the `*_tx` methods directly**:
```text
trait TxHandle {
enqueue_tx(name, opts, payload) -> job_id
publish_tx(stream, payload) -> event_id
publish_tx(stream, payload) -> offset
publish_with_key_tx(stream, key, payload) -> offset // ADR-015
notify_tx(channel, payload)
save_offset_tx(stream, consumer, offset)
outbox_enqueue_tx(outbox, opts, payload) -> job_id // ADR-014
@@ -277,7 +282,10 @@ v1 variants (guaranteed-matchable on every engine):
`try_recv` / `recv_timeout` form of the wake close in §3).
- `Codec` — payload serialization/deserialization failures
(`payload_as` on stream events and job payloads; wakes carry no
payload per §3).
payload per §3). *(Annotated 2026-10-05: a present key on
`publish_with_key(_tx)` being non-empty is validated as
`InvalidName` ([ADR-015](015-streams-depth.md) §2 — the same
entry-point-validation variant, not a new one).)*
- `Database` — everything else: driver/connection/SQL errors, opaque
to the contract, engine detail preserved via the error source chain.
No engine-specific variants are minted for its contents.
@@ -339,6 +347,7 @@ differs):
| `Lock::heartbeat(ttl)` | `Lock::renew(ttl)` | renew describes the TTL discipline the guarantee row pins; heartbeat is the queue-side renewal term |
| `Queue::claim_waker()` | (not surfaced) | engine-internal consumption wake (§1) |
| `StreamSubscription` (`save_every`, auto-save on drop) | `subscribe(consumer) -> Box<dyn EventReceiver>`, explicit saves only | the explicit-offset obligation (core-contract.md streams; the auto-checkpoint ambiguity is not inherited, [ADR-006](006-wake-and-delivery-contract.md)) |
| `StreamEvent.topic` (field) | `StreamEvent.stream` *(added 2026-10-05 by [ADR-015](015-streams-depth.md) §3)* | one term for the mechanism everywhere — the contract names it streams (the §4 kinds list); the substrate's column stays `topic` per the [ADR-012](012-forked-substrate-design.md) §3 fidelity posture |
| `update_events()` | (not surfaced) | superseded by the wake subscription — the pinned `WakeReceiver` *is* the reactive-event surface; honker's separate raw update-event stream has no contract role |
| `prune_notifications` / `prune_notifications_keep_latest` | (not surfaced) | notifications-table maintenance tooling on the SQLite engine; disposition rides OQ-05's sweep/maintenance design (consumer-visible only if that design surfaces it) |
| `Scheduler` surface | (out of v1, OQ-09 decides) | §1 |
@@ -53,7 +53,7 @@ One method, added to the `TxHandle` trait:
```text
trait TxHandle {
enqueue_tx(queue, opts, payload) -> job_id
publish_tx(stream, payload) -> event_id
publish_tx(stream, payload) -> offset
notify_tx(channel, payload)
save_offset_tx(stream, consumer, offset)
outbox_enqueue_tx(outbox, opts, payload) -> job_id // the addition
@@ -0,0 +1,326 @@
# ADR-015: Streams depth — key semantics, `StreamEvent` shape, ordering row, retention
## Status
Accepted
## Context
[OQ-12](../../architecture/open-questions.md) — found by the contract
v1 Phase 1 architecture review: [ADR-008](008-contract-v1-pinning.md)
§1 pinned streams' *method names* (`publish`, `publish_with_key`,
`read_since`, `read_from_consumer`, `save_offset`, `get_offset`,
`subscribe`) but not their *depth*. Four sub-questions were open:
1. **Key semantics** — what `publish_with_key`'s key means, and what
if any ordering it buys. Honker's realization is a stored nullable
`key` column (`_honker_stream(topic, key, payload)`) with no
per-key read path or ordering enforcement anywhere; the doc comment
says "used for per-key ordering downstream"
(`packages/honker-rs/src/lib.rs:894`) — per-key FIFO is an
*emergent* property of a consumer reading `offset ASC` filtered by
key, not a server guarantee. Three options were on the table: (a)
pin honker's shape (key = carried metadata, ordering stays global
FIFO, document the emergent pattern); (b) pin server-enforced
per-key ordering (real machinery on both engines, gated on a
consumer-inventory row naming the need); (c) cut `publish_with_key`
from v1 (no consumer row names keys).
2. **`StreamEvent` shape** — honker's is
`{ offset: i64, topic: String, key: Option<String>, payload:
Vec<u8>, created_at: i64 }` with `payload_as<T>` decode
(lib.rs:1008). The contract needs its own pinned field list.
3. **Ordering guarantee row** — [ADR-006](006-wake-and-delivery-contract.md)
§2's delivery table pins streams' durability/replay/atomicity but
its ordering cell content was effectively "readable"; the contract
suite cannot pin cross-engine stream equivalence without an
explicit ordering row.
4. **Retention / bounded growth** — the reference read flags the
stream table as unbounded-growth-unless-swept (honker-machinery
§8 defect 5); no consumer row names retention
(replay-deep-history is the *purpose* of the subscription need).
The consumer evidence was completed for this decision by adding the
**alkcall row** to
[consumer-inventory.md](../../research/consumer-inventory.md)
(ecosystem-shape evidence, added per the inventory's own
row-before-assumption convention). Alkcall's pub/sub vocabulary —
`Sub` operations (server→client streaming, ADR-021), `Pub`
(client→server streaming, ADR-046), `channel/resources/subscribe`
(snapshot + change events, ADR-037), and the deferred
broker/fan-out (ADR-046's Gap B, channels OQ-22) — is the shape
downstream crates' reactive ops are expected to take over this store,
and the operator confirmed alkstore is meant as the base storage
layer for that ecosystem. Its ordering posture is per-stream
(ordered, reliable transport streams; correlation by request ID):
**no alkcall mechanism names server-enforced per-key ordering.** The
topic/event-type dimension of its subscriptions routes to *separate
streams* (one stream per topic; e.g. `channel/resources/subscribe` is
its own event source), not to keys within one stream.
## Decision
### 1. Key semantics: option (a) — the key is carried metadata; ordering stays global FIFO
Option (b) is **rejected on the decision rule it named for itself**:
it requires "a consumer row naming the need," and neither the
inventory's original rows nor the alkcall row names server-enforced
per-key ordering anywhere. Real machinery (per-key sequencing,
per-key claim/exhaustion paths, or per-key head-of-line bookkeeping on
both engines) for an unnamed need is exactly the scope discipline
[ADR-002](002-feature-scope.md) exists to enforce.
Option (c) (cut `publish_with_key` from v1) is **rejected**: the
method name is a pinned v1 surface element
([ADR-008](008-contract-v1-pinning.md) §1 — removing it would be a
contract-breaking event for OQ-10 to govern, not a pre-implementation
refinement); the *cost* of carrying the key in the inherited
realizations is one nullable column honker's schema already has; and
the key is genuinely useful as a *grouping token* — the emergent
per-key pattern is real and cheap for consumers (SQL-level filter +
in-order read), and an ecosystem whose ops carry keys (a repo id, an
entity id) deserves a stable way to store that grouping durably.
Cutting the only keyed-write path would push consumers to encode keys
into payloads — an unqueryable convention — the same shape of waste
ADR-010 §7 called out for result storage.
Pinned semantics, exactly honker's:
- `publish_with_key(stream, key, payload)` stores `key` as a nullable
column value on the event row; `publish` stores `NULL`.
- The key is **carried metadata** — it appears on every `StreamEvent`
read back and is `Option<String>`-shaped in the event type; the
engine assigns it no behavioral role (it rides the row the publish
writes; the wake trigger is the publish itself, key or no key).
- **The ordering guarantee stays global FIFO by offset, per stream** —
regardless of key. Per-key FIFO-by-offset is the *documented
emergent pattern*: a consumer that filters a read by key observes
that stream's events for that key in offset order. The contract
does not enforce, promise, or price any cross-subscriber per-key
handoff. (The guarantee row, §3, is written so this remains *true
under trimming* — offset order, never renumbered.)
- The key is **not** a partitioning instruction, not a dedup key, and
not an ordering key in any engine-enforced sense. Nothing in the
contract suite tests keys behaviorally beyond carried-metadata
round-trip and equivalence across engines.
The alkcall evidence makes this the right ecosystem shape too: its
subscriptions distinguish event *sources* (separate streams per
topic), not positions *within* a source (per-key lanes). When the
channels broker (Gap B) arrives, topic matching selects streams;
within a stream, offset order is the whole story.
### 2. `publish_with_key_tx` — the tx seam carries the key
A gap the key decision exposed in the pinned tx seam:
[TxHandle](../../architecture/core-contract.md) carries
`publish_tx(stream, payload)` but no keyed form — so a keyed event
could not be written commit-atomic with a business write, breaking the
no-ghosts property's uniformity for keyed publishes (the same shape of
hole OQ-13 exposed for outbox enqueue, resolved by
[ADR-014](014-outbox-tx-enqueue.md)). Pinned:
`publish_with_key_tx(stream, key, payload) -> offset` joins the
TxHandle trait — one method, symmetric with `publish_tx`, validating
the key the same way the non-tx form does (`None` or non-empty — an
empty `Some` key is `InvalidName`) with no new error variants
([ADR-008](008-contract-v1-pinning.md) §5's rule). Framed as
completing contract v1 in place, pre-implementation, like ADR-014 —
no versioning event, OQ-10 untouched. A plain-`publish_tx` is
`publish_with_key_tx(stream, None, payload)`, not a separate
engine path. *(Naming pinned here too: the publish family returns the
**assigned offset** — honker's `publish` already returns it — and the
tx-seam sketches' `event_id` placeholder name is corrected to
`offset` so one term covers the value every read/cursor/trim method
speaks.)*
### 3. `StreamEvent` shape — honker's, with one rename
Pinned (inherited near-verbatim, per the fork's API fidelity posture,
[ADR-012](012-forked-substrate-design.md) §3):
```text
struct StreamEvent {
offset: i64, // position in the stream's log; immutable
stream: String, // the stream name (honker: "topic")
key: Option<String>, // §1's carried metadata; None unless
// published keyed
payload: Vec<u8>, // core value bytes
created_at: i64, // unix seconds at publish (the engines'
// single clock source)
}
```
- The `topic` field renames to **`stream`** — the contract surface
names this mechanism streams, not topics (honker-rs's `Stream`
handle's `topic()` accessor is likewise the stream name); one term
everywhere, matching the reserved-namespace kinds list
([ADR-008](008-contract-v1-pinning.md) §4).
- `payload` is core value bytes crossed as a non-generic value; the
concrete `StreamEvent` owns the bytes; `payload_as<T>` is the decode
convenience ([ADR-008](008-contract-v1-pinning.md) §8's rule — the
deserialization error is `Codec`).
- `created_at` is **unix seconds at publish**, from the engine's
single clock source (the partial-index/clock-pin rule in
[queues.md](../../architecture/queues.md)); it is *informational*
and NOT an ordering candidate — offset is the only ordering field
(two events can share a second; offset never collides).
- `offset` is the stream-wide position, engine-assigned at publish,
monotonically increasing per stream on both engines. It is
**immutable once assigned** and is the value every read/cursor/
trim method speaks.
### 4. Ordering guarantee row — pinned
ADR-006 §2's table's streams row extends with an explicit ordering
cell (the guarantee cell gains a clause; durability/replay/atomicity
as already pinned):
> **streams** — Guarantee: *global FIFO by offset within each stream
> (`read_since`/`subscribe` yield `offset ASC`); every committed
> event readable by every consumer that hasn't passed its offset;
> offsets immutable and never renumbered.*
This is the row the contract suite pins cross-engine stream
equivalence against: two engines, same publish sequence → same offset
sequence → same `read_since` output order, for direct and subscriber
reads alike. It is deliberately *not* stronger than both engines can
implement by a shared SQL shape (`ORDER BY offset ASC` is already
both engines' claim path) and not weaker than the mechanism's purpose
(replay demands monotone positions).
### 5. Retention: `trim_to` on the stream handle — a consumer-invoked bounded-growth op
The mechanism-shaped answer to defect 5's family (unbounded
`_honker_stream` growth), applying the ADR-010 §6 posture to the one
stream-side table that consumers interact with *as rows* (unlike
notifications, the transport detail whose hygiene stays engine-
internal):
- `store.stream(name)`'s handle gains `trim_to(horizon)` — delete
events with `offset <= horizon`, emitting no dedicated wake and no
notify (deletes are not events; consumers whose cursors move
forward never lose reads). On SQLite the `data_version` watcher may
still fire on the trim commit — any committed write bumps it; that
spurious hint is contract-legal (the wake contract's best-effort +
idempotence posture, ADR-006 §1), not a trim-specific delivery.
`OffsetGone`-style errors are not needed: reads past the trim
horizon simply return fewer/no rows (a consumer reading from a
trimmed region has, by definition, an offset older than the
horizon — that is the API contract; re-attaching consumers at
trimmed-away offsets resume at the horizon's first remaining row).
- **No engine-side retention default, no ambient sweeper** (the
ADR-010 §6 posture; consumer-side, cadenced by the collapse recipe
like every sweep). Trim is *available*, not automatic: the
replay-forever default is the *purpose* for the consumer rows, so
the default is unbounded — with the growth documented squarely,
this time with an in-contract tool to fix it.
- **Caveat documented in-contract**: offsets are never renumbered and
never reused; surviving events keep the offsets they were published
with (gaps after a trim are legal and ordinary). Consumers must
not assume dense offset sequences (the same discipline the
scheduler's catch-up/skip-forward already imposes).
- Trim's effect on cursors: a consumer's *saved* offset below the
trim horizon stays (it remains a valid position marker; the next
read resumes at the horizon's first remaining row); a consumer may
re-save after processing. This keeps `get_offset`/`subscribe`
resume semantics well-defined without engine-side cursor-tracking
machinery.
Rejected retention alternatives (per the OQ's option set): the
consumer-recipe-without-an-API option was self-admittedly "document
the gap" (no deletion path exists today — this ADR *creates* the
deletion path, which is what the recipe needed); documenting-nothing
repeats the notifications-growth mistake ADR-010 §6 just fixed.
## Consequences
**Positive**
- Streams' depth is fully pinned: names (ADR-008 §1) + semantics (this
ADR) make the mechanism implementable on both engines from the same
text, and the contract suite has an ordering row to pin cross-engine
equivalence against.
- The key decision is honest: a zero-machinery feature (carried
metadata) documented for exactly what it is, with option (b)'s
re-entry gate (a consumer-inventory row naming server-enforced
per-key ordering) recorded — the ecosystem-shape evidence says the
topic dimension routes to streams, not keys.
- Keyed publishes are tx-seam-native — keyed and unkeyed events carry
the same commit-atomicity, so an outbox-style keyed event log co-writes
with business data correctly.
- `trim_to` gives the family a real answer to stream growth in the
contract, consistent with the no-ambient-sweeper posture (invoked
deliberately, cadenced from above).
**Negative**
- Per-key *ordering enforcement*, if a consumer ever names it, is a
later contract extension with real machinery on both engines —
carrying metadata now does not prejudge or approximate it, and the
emergent pattern's documented status must not be read as a promise.
- `trim_to` adds one more consumer-invoked maintenance op to the
surface (backlog rows follow).
- `created_at` is second-precision (the shared clock pin) — consumers
needing sub-second publish timestamps must encode them in payloads.
- The `StreamEvent.stream` rename diverges from honker's wire-era name
(`topic`) — the fork keeps the column name `topic` engine-internally
per the ADR-012 §3 fidelity posture; only the contract type renames.
## Verification backlog additions
- **Cross-engine stream equivalence** (the §3/§4 row): publish
sequence → identical offset sequence → identical `read_since`
output order, both engines; keyed and unkeyed interleavings
preserve global FIFO; `key` round-trips exactly (`None` vs
`Some`), `StreamEvent.stream`/`created_at` equal on both engines.
- **`publish_with_key_tx` commit-atomicity on both engines** (§2):
rollback drops the keyed event row with the business write (no
ghost event); commit makes it visible to `read_since`/`subscribe`;
empty-`Some`-key `InvalidName` on both engines' tx paths.
- **`trim_to` semantics on both engines** (§5): exact-boundary trim
(`<=`), surviving rows keep their offsets (gaps legal, never
renumbered), reads from a trimmed horizon resume correctly,
`save_offset` below the horizon stays valid, trim emits no
dedicated wake/notify (on SQLite the `data_version` watcher may
fire spuriously — pinned as contract-legal), and a concurrent
subscriber mid-`subscribe` never loses its place
(offset-anchored reads only move forward).
## References
- [OQ-12](../../architecture/open-questions.md) — this ADR's resolution.
- [core-contract.md](../../architecture/core-contract.md) — the spec
this ADR's §1–§5 pin into place (streams section, tx seam,
guarantee table, backlog).
- The alkcall evidence (reference checkout
`/workspace/@alkdev/alkcall` @ HEAD `main`, 2026-10-05, per the
AGENTS.md §3 posture):
decisions/021 (Sub — the Subscription streaming handler),
decisions/046 (Pub + the `Subscription`→`Sub` rename + Gap B),
decisions/037 (`channel/resources/subscribe`), open-questions.md
OQ-22 (fan-out deferral) — the pub/sub vocabulary this contract's
streams substrate; its no-per-key-ordering posture is the
inventory row added for this decision.
- [consumer-inventory.md](../../research/consumer-inventory.md) — the
streams row (operator-authority record) + the alkcall
ecosystem-shape row; the (b)-gate that rejected server-enforced
per-key ordering.
- [ADR-006](006-wake-and-delivery-contract.md) §2 — the guarantee
table the streams ordering cell extends.
- [ADR-008](008-contract-v1-pinning.md) §1/§4/§5/§8 — the pinned
method names the depth lands on, the namespace kinds, the
act-differently error rule, the `payload_as`/`Codec` rule.
- [ADR-010](010-queue-semantics-depth.md) §6 — the no-ambient-
sweeper posture this ADR's retention decision applies to the
consumer-facing stream table (the notify-table hygiene precedent
is §6's too).
- [ADR-012](012-forked-substrate-design.md) §3 — the fidelity posture
the column-name keep rides; [ADR-014](014-outbox-tx-enqueue.md) —
the amend-in-place framing §2 reuses.
- Honker stream machinery
(`/workspace/honker` `honker-core/src/honker_ops.rs:1906-1996`, `src/lib.rs:417-429`; `packages/honker-rs/src/lib.rs:876-1010, 1099`)
— the inherited realization (key column, offset-ASC reads, monotone
offset saves, event shape); the unbounded-growth defect
(reference-honker-machinery.md §8 defect 5).