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

13 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
pg-engine-streams Postgres engine — streams (`StreamHandle`, reads, subscribe, trim) completed
pg-engine-seam-tx
pg-engine-notify-listen
moderate medium component implementation
wave-4
postgres-engine

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, INSERT into the events table (bigserial id = the assigned offset; nullable key column — carried metadata, empty-Some-key → InvalidName per the tx path's rule), encode_payload at the seam (Codec propagation), return to pool. Then wake: pg_notify on 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 → empty Vec at trait-impl entry, ADR-023 §2 — never LIMIT -1-class dialect semantics; pg's LIMIT with a negative is an error, which the guard makes unreachable). read_from_consumer anchors 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 or WHERE offset > stored; one owner of the monotonicity rule) and the checkpoint read. Both save forms' Err arm carries Database only (ADR-021 §5).
  • trim_to(horizon): DELETE … WHERE id <= horizon on 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's save_offset() never checkpoints unobserved events. Backpressure: bounded channel + backpressure-aware send.
    • recv()'s Err carries Database only (ADR-021 §5 — remap non-Database read errors with the original as source, the SQLite database_only helper's shape); try_recv arms mirror the wake receiver (idle Ok(None), Err(Closed) after shutdown-close).
    • The receiver's read_since (anchored, extent-guarded) and sync save_offset (the last-yielded offset through the same monotone op) per ADR-019 §6.
  • 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).
  • 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_key return the assigned bigserial offset; key round-trips (None vs Some); encode_payload's Codec propagates; empty-Some-key rejected InvalidName
  • Reads yield offset ASC (global FIFO); extent guard (limit <= 0 → empty Vec) at trait-impl entry; offsets immutable (trim leaves gaps, never renumbers — tested)
  • save_offset monotone (regression save = silent no-op); get_offset reads the checkpoint; save Err = Database only
  • trim_to exact-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 valid
  • subscribe: 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 → subscribe resumes 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().await is 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_on panics on a runtime worker (both flavors), Runtime::new().block_on nested panics and would build a runtime per call, block_in_place + block_on of 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's WAKE_RUNTIME); drive_sync runs the save future block_on it — directly on the caller's thread when outside any runtime or inside a foreign spawn_blocking, via block_in_place on 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::spawn cannot 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_rows moved from tx.rs into stream.rs as the one pub(crate) decode owner (the tx reads import it — the SQLite twin's decode-sharing shape, the stream field mapping and the Codec posture guard live once). NOTIFY_PAYLOAD_LIMIT is 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_to is 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 <= -7 matches 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 one Err(Database) and waits for the next wake (the SQLite twin's transient arm, verbatim). database_only remaps non-Database errors (the decode-side Codec guard) into the opaque fallback with the source preserved.
  • The closed_check helper from the store is not reused on the handle (the handle checks its own closed arc inline at each op entry — the ops return from sync context into boxed futures, so the check happens before the Box::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 full StreamHandle trait: publish/publish_with_key auto-commit INSERT .. RETURNING id on a pool connection with encode_payload at the seam and the best-effort pg_notify wake on the stream's channel after the insert (separate statement, failure logged-and-swallowed), read_since/read_from_consumer with the ADR-023 §2 extent guard and id ASC ordering, monotone save_offset/get_offset (the shared upsert owner), trim_to's pool-connection DELETE) and the durable subscribe receiver (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, sync save_offset bridged through the stateless idle wake runtime, recv() Err = Database only).
  • tx.rs: now imports stream_events_from_rows from stream.rs (the one decode owner; the forked copy deleted).
  • store.rs: Store::stream wired (validation + closed-store fail-closed + handle construction); the open task's stub-surface test updated (stream asserted wired, queue the 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), receiver read_since anchoring + 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-postgres against the harness server (pglo-poc :15432): 72 lib + 9 schema tests green, repeated ×4 (the stream subset 4× stable); workspace cargo test green server-less (pg tests skip per convention, 10 result lines all ok); cargo clippy --all-targets -- -D warnings clean; cargo fmt --check clean.