--- id: sqlite-engine-streams name: SQLite engine — streams (`StreamHandle`, reads, subscribe, trim) status: completed depends_on: [sqlite-engine-seam-tx] scope: moderate risk: medium impact: component level: implementation tags: [wave-3, sqlite-engine] --- ## Description Implement `Store::stream` and the full `StreamHandle` trait over the substrate's stream ops (`stream_publish`, `stream_read_since`, `stream_save_offset`, `stream_get_offset` — inherited near-verbatim; the `topic` column carries the contract's `stream` field, the one name delta, mapped at this layer): - `stream(name)` — validated constructor returning the boxed handle. - `publish` / `publish_with_key` — auto-commit publishes through the writer slot; `encode_payload` at the seam (`Codec` propagates); return the assigned offset. - `read_since(offset, limit)` / `read_from_consumer(consumer, limit)` — reader-pool reads, offset ASC (global FIFO); map rows into `StreamEvent` via the core constructor task's `from_row`. **Extent guard (ADR-023 §2)**: `limit <= 0` yields the empty `Vec` at the trait-impl entry, before the substrate — the `LIMIT -1` dialect artifact dies here (the substrate passes `limit` through contract-blind; the guard is the engine's obligation). - `save_offset(consumer, offset)` / `get_offset(consumer)` — monotone saves (substrate's upsert already refuses regressions silently); `get_offset` of an absent consumer = 0. - `trim_to(horizon)` — `DELETE FROM __alkstore_stream WHERE topic = ? AND offset <= ?` inside the writer-slot lease; boundary argument is **total** (ADR-023 §2: a negative horizon deletes nothing, idempotently — no guard, pass through); returns the deleted count. - `subscribe(consumer)` — the durable receiver: attach, read to tail from the stored offset, then wake-driven re-reads off the watcher fanout (same bridge shape as `listen`, but yielding `Result` per the `EventReceiver` trait — `Err` carries `Database` only; `None` = closed on watcher death, terminal). Receiver-form `save_offset()` saves the last-yielded event's offset through the same monotone op (ADR-019 §6); `offset()` reports the receiver's position. - Consumer identifiers (`consumer`) take the non-empty rule only; stream names the shared-namespace rules. ## Acceptance Criteria - [x] Full `StreamHandle` impl; publish/read/save/get/trim/subscribe all work against a temp-file store - [x] Extent guard: `read_since`/`read_from_consumer` with `limit <= 0` return the empty `Vec` (never the whole stream) — pinned per ADR-023 §2's backlog row - [x] `trim_to` exact-boundary (`<=`), negative horizon deletes nothing, surviving offsets never renumbered - [x] Monotone saves: regression saves are silent no-ops; direct and receiver forms compose without regression - [x] `subscribe` replays from the stored offset, then delivers post-attach events wake-driven; watcher death closes it (`None`, terminal); `recv()`'s `Err` carries `Database` only - [x] Key round-trips exactly (`None` vs `Some`); `stream` field carries the contract name (not `topic`) - [x] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean ## References - docs/architecture/core-contract.md (streams section) - docs/architecture/decisions/015-streams-depth.md - docs/architecture/decisions/019-mechanism-handle-surfaces.md §6 - docs/architecture/decisions/023-fourth-review-round.md §2 - docs/architecture/engine-sqlite.md (Mapping the contract: streams) - alkstore-sqlite/src/substrate/ops.rs (stream_* fns) ## Notes - **Stream module (stream.rs)**: `SqliteStreamHandle { name, writer, readers, watcher, closed }` — the handle carries the store's machinery arcs (writer slot, reader pool, watcher, closed flag). Publishes (both forms) ride `with_writer` (the seam's short-lived slot lease); reads ride a new `with_reader` bridge (pooled reader, `spawn_blocking`, release on both arms); `trim_to`'s DELETE rides `with_writer` inside the writer-slot lease, straight SQL over the substrate (the substrate has no trim op — engine-side statement, `deleted as i64` count). - **Decode shared**: `stream_events_from_json_page` moved from `tx.rs` into `stream.rs` as the one `pub(crate)` decode owner (both the tx reads and the auto-commit reads call it; the `topic` → `stream` name delta and the `Codec` posture guard live there once). - **Extent guard placement**: `read_since`/`read_from_consumer` (and the receiver's `read_since`) check `limit <= 0` at trait-impl entry, before validation of *any* round trip into the substrate — the empty-`Vec` is the whole result, `LIMIT -1` can never fire. The closed-store check precedes the guard, so a closed store still fails closed even for `limit <= 0` (the guard is about what an *open* store's limit means). - **Subscribe bridge**: a dedicated std thread (one per subscription, the `listen` bridge shape) — (1) `watcher.subscribe()` is taken *before* any read so a commit cannot fall between the attach read and the wake wait; (2) initial cursor read of the consumer's *stored* offset; (3) page-drain to the tail (`SUBSCRIBE_PAGE = 256`, `offset ASC`); (4) park on the wake feed, re-drain on wake (overtrigger-coalesced hints); exit on wake-feed disconnect (watcher death / store close) or failed send. Read failures while open surface one `Err(Database)` item and wait for the next wake (the transient arm); attach-read failure yields one `Err` and closes. A thread-spawn failure unsubscribes and fails `Database` (W-2 posture, no panics cross the seam). - **Receiver position semantics**: `position` = the last-*yielded* event's offset (0 before any yield ⇒ `save_offset()` before any yield is a no-op, ADR-021 §5); advances only on delivered events — the bridge's internal cursor (which advances on *drained-from-storage* events) is separate, so a stalled consumer's `save_offset()` never checkpoints unobserved events. Backpressure: bounded tokio channel + `blocking_send` parks the bridge; delivery resumes exactly where it stopped. - **`EventReceiver::save_offset` is sync in the trait** — the receiver's save takes the writer slot on the calling thread (the slot is free between ops; a long open transaction parks this call — the honest lease posture, surfaced in the doc comment). `offset()` reports the in-memory position, not a db read. - **`recv()`'s `Err` carries `Database` only (ADR-021 §5)**: a `database_only` helper remaps any non-`Database` read error (e.g. the decode-side `Codec` posture guard) into `Database` with the original error preserved as the source. `try_recv` arms mirror the notify receiver: idle `Ok(None)`, `Err(Closed)` after the channel disconnected. - **Closed-store postures**: every handle op (and the receiver's `read_since`/`save_offset`) checks the store's closed flag first and fails closed with `Database`; the reader pool's closed error (`Database is closed` MISUSE shape) is recognized and remapped to the same closed-store error. New `stream(name)` constructors on a closed store fail `Database` (like `listen`). The bridge exits on the closed reads — receiver sees `None`/`Err(Closed)` terminal. - **Validation**: stream names via `validate_shared_name` at `store. stream()` and (defensively) the handle's reads; consumers via `validate_local_name` (non-empty only — reserved prefixes legal for consumers, tested); empty-`Some` keys rejected `InvalidName` per the tx path's rule. - **Test infra**: `store/stream_tests.rs`, 15 tests over temp-dir stores (two-store-open cross-connection probes where wake arrival needs a second writer; raw-reader assertions only where the task demands table-level honesty — trim count/durability read through the handle, not raw SQL — the substrate's ops.rs already probes the raw rows). Stability: suite ran 4× green. ## Summary Implemented the stream mechanism's SQLite arm: `SqliteStreamHandle` (full `StreamHandle` trait — `publish`/`publish_with_key` auto-commit through the writer slot with `encode_payload` at the seam and `Codec` propagation, `read_since`/`read_from_consumer` reader-pool reads with the ADR-023 §2 extent guard and ASC ordering, monotone `save_offset`/`get_offset`, `trim_to`'s writer-lease DELETE with the total boundary argument) and the durable `subscribe` receiver (a std-thread bridge: attach-before-read, replay from the stored offset, wake-driven re-reads off the watcher fanout, terminal close on watcher death, position = last-yielded offset, sync receiver-form `save_offset` through the same monotone op; `recv()`'s `Err` remapped to `Database`-only). Shared the stream-page decoder between the tx and auto-commit paths (`stream.rs` owns it, `tx.rs` imports). Wired `Store::stream` in `store.rs` (entry-point validation + closed-store fail-closed). 15 new tests (round-trip + stream-field + key, extent guard, trim boundary/negative/no-renumbering, monotone saves across forms, replay + wake-driven delivery, terminal close, tail-idle, independent subscribers, receiver-drop unsubscribe, entry-point validation incl. reserved-prefix-local consumers, closed-store fail closures, restart durability, tx-compose with ghost verification, receiver `read_since` anchoring + extent guard). Verified: `cargo test` (workspace: 25 + 3 suite + 161 sqlite), `cargo clippy --all-targets -- -D warnings`, `cargo fmt --check` all clean; suite run 4× green.