13 KiB
13 KiB
id, name, status, depends_on, scope, risk, impact, level, tags
| id | name | status | depends_on | scope | risk | impact | level | tags | ||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|
| pg-engine-streams | Postgres engine — streams (`StreamHandle`, reads, subscribe, trim) | completed |
|
moderate | medium | component | implementation |
|
Description
Implement the stream mechanism's pg arm — the full StreamHandle trait
and the durable subscribe receiver, over the schema task's events/
offsets tables and the forwarder's wake substrate (ADR-015's depth,
engine-postgres.md's streams mapping row):
publish/publish_with_key(auto-commit): pool checkout,INSERTinto the events table (bigserial id = the assigned offset; nullable key column — carried metadata, empty-Some-key →InvalidNameper the tx path's rule),encode_payloadat the seam (Codecpropagation), return to pool. Then wake:pg_notifyon the stream's wake channel (the mechanism-name-is-the-channel realization pinned in the plan's Wave 4 section) — a separate statement on the pooled connection (auto-commit: its own implicit tx; the wake is best-effort, not commit-atomic with the insert — the no-replay hole covers a gap between them; the durable row is the truth).read_since/read_from_consumer: pool reads,ORDER BY id ASC(global FIFO — bigserial offsets, immutable, never renumbered; ADR-015 §4's ordering row), the extent guard (limit <= 0→ emptyVecat trait-impl entry, ADR-023 §2 — neverLIMIT -1-class dialect semantics; pg'sLIMITwith a negative is an error, which the guard makes unreachable).read_from_consumeranchors at the consumer's stored offset.save_offset/get_offset: the monotone upsert (a save below the stored checkpoint is a silent no-op —GREATEST-guarded upsert orWHERE offset > stored; one owner of the monotonicity rule) and the checkpoint read. Both save forms'Errarm carriesDatabaseonly (ADR-021 §5).trim_to(horizon):DELETE … WHERE id <= horizonon a pool connection (wakes nothing — ADR-015 §5's posture; the boundary argument is total: a negative horizon deletes nothing, idempotently), returning the deleted count. Surviving events keep their offsets (gaps legal).subscribe(consumer)— the durable receiver, the structural twin of the SQLite engine's bridge with the pg deltas:- Attach-before-read ordering (subscribe to the stream's wake channel before the initial cursor read — a commit cannot fall between attach and wait).
- Initial cursor read of the consumer's stored offset; page-drain
to the tail (
offset ASC); park on the wake feed (the forwarder fanout for this stream's channel); re-drain on wake (wake-driven with table re-read — events never require polling to become visible, core-contract.md's streams section). - Close semantics (the pg arm, ADR-021 §5): the receiver stays
open across forwarder reconnects — reconnect-wakes arrive on the
wake feed and simply trigger re-drains (which is the recovery:
the no-replay hole's gap is healed by the re-drain reading the
durable rows); closes only at engine shutdown (terminal
None). - 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); the bridge's internal cursor advances on drained-from-storage events — a stalled consumer'ssave_offset()never checkpoints unobserved events. Backpressure: bounded channel + backpressure-aware send. recv()'sErrcarriesDatabaseonly (ADR-021 §5 — remap non-Databaseread errors with the original as source, the SQLitedatabase_onlyhelper's shape);try_recvarms mirror the wake receiver (idleOk(None),Err(Closed)after shutdown-close).- The receiver's
read_since(anchored, extent-guarded) and syncsave_offset(the last-yielded offset through the same monotone op) per ADR-019 §6.
- Validation: stream names via
validate_shared_nameatstore.stream()and (defensively) the handle's reads; consumers viavalidate_local_name(non-empty only — reserved prefixes legal for consumers). - Closed-store posture: every handle op checks the store's closed
flag first and fails closed with
Database(the wave-3 shape).
Acceptance Criteria
publish/publish_with_keyreturn the assigned bigserial offset; key round-trips (NonevsSome);encode_payload'sCodecpropagates; empty-Some-key rejectedInvalidName- Reads yield
offset ASC(global FIFO); extent guard (limit <= 0→ emptyVec) at trait-impl entry; offsets immutable (trim leaves gaps, never renumbers — tested) save_offsetmonotone (regression save = silent no-op);get_offsetreads the checkpoint; saveErr=Databaseonlytrim_toexact-boundary (<=), returns deleted count, wakes nothing, negative horizon deletes nothing; reads resume at the horizon's first remaining row; saved offsets below the horizon stay validsubscribe: replay from stored offset, wake-driven delivery, position = last-yielded offset, backpressure holds, receiver stays open across forwarder reconnects (reconnect-wake triggers re-drain — the gap-healing property tested), terminal close at engine shutdown only- Restart durability: publish → close store → reopen →
subscriberesumes from the saved offset - Validation + closed-store fail-closed on all entry points
cargo test -p alkstore-postgres(harness server), clippy-D warnings, fmt clean; gates green server-less
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/021-tx-reads-and-value-shape-fixes.md §5
- docs/architecture/decisions/023-fourth-review-round.md §2 (extent guard)
- docs/architecture/engine-postgres.md (Mapping the contract: streams)
- alkstore-sqlite/src/stream.rs (the structural twin)
- alkstore/src/stream_handle.rs, event_receiver.rs (the traits)
Notes
Decisions of record the implementation made that the description didn't pin:
- The wake statement is
SELECT pg_notify($1, '')on the same pooled connection, posted-insert, best-effort: auto-commit — the notify fires at its own implicit tx's commit, a separate statement per the description (the gap between insert-commit and wake-commit is the no-replay hole; the durable row is the truth). A wake failure is logged to stderr and swallowed — the durable row already landed, and a wake-fee must never fail (or stall) a publish. - The subscribe receiver is a natively-async bridge task (the
structural twin of the SQLite bridge, without the std thread): a
tokio task holds a broadcast subscription (taken before the attach
read — attach-before-read ordering), drains pages to the tail,
parks on
broadcast::recv, re-drains on wake (including the lagged arm — lagged is woken-state too, surfaced-by-drain). Channel capacity 16 (the SQLite twin's smoothing-buffer shape);send().awaitis backpressure-aware — a stalled consumer parks the bridge, delivery resumes exactly where it stopped. EventReceiver::save_offset's sync signature needed a real sync→async bridge (this engine has no blocking writer slot to lease — the SQLite twin's sync posture does not exist here). Probed exhaustively:Handle::block_onpanics on a runtime worker (both flavors),Runtime::new().block_onnested panics and would build a runtime per call,block_in_place+block_onof a different runtime is safe only on multi-thread workers, and a plain (non-runtime) thread can block directly. Landed shape: a lazily built, process-wide, stateless one-worker idle "wake runtime" (stream.rs'sWAKE_RUNTIME);drive_syncruns the save futureblock_onit — directly on the caller's thread when outside any runtime or inside a foreignspawn_blocking, viablock_in_placeon a multi-thread worker, via a scoped std thread on a current-thread worker (the only escape there). The runtime carries no engine state; each save checks out its own pool connection on the runtime that drives it. Idle cost: one parked runtime (one blocked driver thread).- The bridge spawn failure arm is
is_finished()-detected (+ abort/Database):tokio::spawncannot fail directly — a shut-down runtime drops the task before it runs, so the spawn site checks finished-ness and fails the subscribe closed (no boxed receiver surfacing as a lie, the W-2 posture's pg shape). - Decode ownership moved:
stream_events_from_rowsmoved fromtx.rsintostream.rsas the onepub(crate)decode owner (the tx reads import it — the SQLite twin's decode-sharing shape, thestreamfield mapping and theCodecposture guard live once).NOTIFY_PAYLOAD_LIMITis not needed here — the wake payload is the empty string (stream wakes carry no payload; the durable row is the truth; the 8000-byte budget never binds). trim_tois a plain parameterized DELETE on a pool connection (wakes nothing — no notify statement at all; ADR-015 §5). The boundary argument needs no guard on SQL:id <= -7matches no rows server-side, idempotently (the totality is native here — contrast the SQLite arm's writer-lease).- The wake-feed arm shape in
run_subscribe_loop: store-close is checked at every read arm (drain, attach) and exits the loop (PageArm::Closed) — the consumer sees the terminal close via the channel disconnect; a mid-drain error while open surfaces oneErr(Database)and waits for the next wake (the SQLite twin's transient arm, verbatim).database_onlyremaps non-Databaseerrors (the decode-sideCodecguard) into the opaque fallback with the source preserved. - The
closed_checkhelper from the store is not reused on the handle (the handle checks its ownclosedarc inline at each op entry — the ops return from sync context into boxed futures, so the check happens before theBox::pin, matching the SQLite twin's placement). - Test-infra note: the "wakes nothing" and gap-heal pins use the
kill_listener_backend shape from the notify tests (the
reconnect-wake then re-drain cycle is what heals the gap); the
drop-unregister pin probes the raw forwarder fanout (no
cross-session LISTEN catalog view exists — the notify tests'
finding reused). The trimmed-horizon reads-resume assertion pins
the first remaining row explicitly (
offsets[2]).
Summary
What landed, verified how:
alkstore-postgres/src/stream.rs(new): the stream mechanism's pg arm —PgStreamHandle(the fullStreamHandletrait:publish/publish_with_keyauto-commit INSERT .. RETURNING id on a pool connection withencode_payloadat the seam and the best-effortpg_notifywake on the stream's channel after the insert (separate statement, failure logged-and-swallowed),read_since/read_from_consumerwith the ADR-023 §2 extent guard andid ASCordering, monotonesave_offset/get_offset(the shared upsert owner),trim_to's pool-connection DELETE) and the durablesubscribereceiver (PgEventReceiver+ the async bridge task: attach-before-read, replay from the stored offset, wake- driven re-drains — reconnect-wakes included, terminal close at engine shutdown only, position = last-yielded offset, backpressure- aware delivery, syncsave_offsetbridged through the stateless idle wake runtime,recv()Err =Databaseonly).tx.rs: now importsstream_events_from_rowsfromstream.rs(the one decode owner; the forked copy deleted).store.rs:Store::streamwired (validation + closed-store fail-closed + handle construction); the open task's stub-surface test updated (streamasserted wired,queuethe remaining-stub representative).- Tests (
store/stream_tests.rs, 21 new): the round-trip (offsets/stream-field/keys/decode), cross-store wake delivery, the extent guard (0/-1/MIN on both read forms), trim (negative/zero/ exact-boundary/idempotent/no-renumbering/reads-resume/saved-offset- survives/wakes-nothing), monotone saves across forms (incl. the receiver pre-yield no-op), replay + wake-driven delivery, gap-heal across a real backend-kill reconnect (the gap event arrives via the re-drain; the receiver never closes), terminal shutdown close (None+Err(Closed)+ never reopens), tail-idle, independent subscribers, receiver-drop UNLISTEN (raw-fanout probe), key round-trip on every read form, byte-exact stored payload (the raw BYTEA row vs the serde_json serialization), tx-compose (commit/ghost/monotone-tx-save), receiverread_sinceanchoring + extent guard, stalled-consumer saves-only-observed, backpressure (48-event burst delivered exactly in ASC order), entry-point validation (shared kinds + empty-Some keys + reserved-prefix-local consumers + residue-free rejects), and closed-store fail-closed on every op + constructors. - Verification:
cargo test -p alkstore-postgresagainst the harness server (pglo-poc :15432): 72 lib + 9 schema tests green, repeated ×4 (the stream subset 4× stable); workspacecargo testgreen server-less (pg tests skip per convention, 10 result lines all ok);cargo clippy --all-targets -- -D warningsclean;cargo fmt --checkclean.