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
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
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.