Files
alkstore/tasks/sqlite-engine-streams.md

9.3 KiB
Raw Permalink Blame History

id, name, status, depends_on, scope, risk, impact, level, tags
id name status depends_on scope risk impact level tags
sqlite-engine-streams SQLite engine — streams (`StreamHandle`, reads, subscribe, trim) completed
sqlite-engine-seam-tx
moderate medium component implementation
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<StreamEvent> 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

  • Full StreamHandle impl; publish/read/save/get/trim/subscribe all work against a temp-file store
  • 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
  • trim_to exact-boundary (<=), negative horizon deletes nothing, surviving offsets never renumbered
  • Monotone saves: regression saves are silent no-ops; direct and receiver forms compose without regression
  • 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
  • Key round-trips exactly (None vs Some); stream field carries the contract name (not topic)
  • 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.