diff --git a/alkstore-postgres/src/lib.rs b/alkstore-postgres/src/lib.rs index f219a4a..2c357df 100644 --- a/alkstore-postgres/src/lib.rs +++ b/alkstore-postgres/src/lib.rs @@ -39,6 +39,7 @@ mod resolution; mod schema; mod seam; mod store; +mod stream; mod tx; pub use opts::{DEFAULT_MAX_SIZE, PgOpts}; diff --git a/alkstore-postgres/src/store.rs b/alkstore-postgres/src/store.rs index 5f00fe2..76244fa 100644 --- a/alkstore-postgres/src/store.rs +++ b/alkstore-postgres/src/store.rs @@ -261,12 +261,22 @@ impl Store for PgStore { fn stream<'a>( &'a self, - _name: &str, + name: &str, ) -> alkstore::BoxedFuture<'a, alkstore::Result>> { + if let Err(e) = alkstore::validate_shared_name(name) { + return Box::pin(async move { Err(e) }); + } if let Err(e) = self.closed_check() { return Box::pin(async move { Err(e) }); } - Box::pin(async { Err(stub("stream wiring lands with the streams task")) }) + let handle = crate::stream::PgStreamHandle::new( + name.to_string(), + self.pool.clone(), + self.schema.clone(), + self.forwarder.clone(), + self.closed.clone(), + ); + Box::pin(async move { Ok(Box::new(handle) as Box) }) } fn queue<'a>( @@ -348,5 +358,8 @@ mod open_tests; #[cfg(test)] mod notify_tests; +#[cfg(test)] +mod stream_tests; + #[cfg(test)] mod tx_tests; diff --git a/alkstore-postgres/src/store/open_tests.rs b/alkstore-postgres/src/store/open_tests.rs index 9314ab2..fdc6581 100644 --- a/alkstore-postgres/src/store/open_tests.rs +++ b/alkstore-postgres/src/store/open_tests.rs @@ -659,11 +659,10 @@ async fn store_trait_methods_are_wiring_stubs() { .notify("stub_ch_probe", serde_json::json!("x")) .await .unwrap(); - let err = match store.stream("s").await { - Err(e) => e, - Ok(_) => panic!("stream must be a stub"), - }; - assert!(matches!(err, alkstore::Error::Database(_))); + // The stream constructor is wired (the streams task's): returns a + // named handle — full behavior the stream tests' scope. + let handle = store.stream("s").await.unwrap(); + assert_eq!(handle.name(), "s"); let err = stub_err!(store.queue("q", alkstore::QueueOpts::default()).await); assert!(matches!(err, alkstore::Error::Database(_))); let err = stub_err!(store.outbox("o").await); @@ -692,8 +691,8 @@ async fn store_trait_methods_are_wiring_stubs() { // The remaining stub messages name their landing task (the // wave-3 posture's "… wiring lands with the … task" shape) — in // the source chain (`Database`'s Display is opaque; ADR-008 §5's - // detail-carriage). `stream` sampled as the representative. - let err = stub_err!(store.stream("ch").await); + // detail-carriage). `queue` sampled as the representative. + let err = stub_err!(store.queue("ch", alkstore::QueueOpts::default()).await); match err { alkstore::Error::Database(source) => { let msg = source.to_string(); @@ -702,7 +701,7 @@ async fn store_trait_methods_are_wiring_stubs() { "stub errors name their landing task, got: {msg}" ); } - other => panic!("stream stub must be Database, got {other:?}"), + other => panic!("queue stub must be Database, got {other:?}"), } drop(store); diff --git a/alkstore-postgres/src/store/stream_tests.rs b/alkstore-postgres/src/store/stream_tests.rs new file mode 100644 index 0000000..d5f441d --- /dev/null +++ b/alkstore-postgres/src/store/stream_tests.rs @@ -0,0 +1,1241 @@ +//! The stream mechanism's acceptance tests (`pg-engine-streams`): +//! publish/read/save/get/trim/subscribe all work against the harness +//! server, wake-driven delivery through the forwarder fanout, the +//! gap-healing property (a subscriber survives a forwarder reconnect +//! and re-drains the durable rows it missed — the no-replay hole +//! healed on the wake feed), the extent guard, exact-boundary trim +//! with surviving offsets unrenumbered and negative horizons deleting +//! nothing, monotone saves across the direct and receiver forms, key +//! round-trip, restart durability, and entry-point validation + +//! closed-store fail-closed on all entry points (ADR-015; ADR-019 +//! §1/§6; ADR-021 §5; ADR-023 §2). +//! +//! Harness convention (the schema task's): connection settings ride +//! the environment (`ALKSTORE_PG_HOST/PORT/USER/PASSWORD/DB`), never +//! hardcoded; tests without a reachable server skip cleanly so the +//! workspace gates stay green server-less. Isolation is a fresh unique +//! schema per test. + +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; + +use alkstore::{Error, EventReceiver, Store}; +use serde_json::{Value, json}; + +use crate::opts::PgOpts; +use crate::schema::DEFAULT_SCHEMA; +use crate::store::open_store; + +const ENV_HOST: &str = "ALKSTORE_PG_HOST"; +const ENV_PORT: &str = "ALKSTORE_PG_PORT"; +const ENV_USER: &str = "ALKSTORE_PG_USER"; +const ENV_PASSWORD: &str = "ALKSTORE_PG_PASSWORD"; +const ENV_DB: &str = "ALKSTORE_PG_DB"; + +fn harness_dsn() -> Option { + let host = std::env::var(ENV_HOST).ok()?; + let port: u16 = std::env::var(ENV_PORT).ok()?.parse().ok()?; + let user = std::env::var(ENV_USER).ok()?; + let password = std::env::var(ENV_PASSWORD).ok()?; + let db = std::env::var(ENV_DB).unwrap_or_else(|_| "postgres".to_string()); + Some(format!( + "host={host} port={port} user={user} password={password} dbname={db}" + )) +} + +fn instance_namer(tag: &str) -> impl Fn() -> String + use<'_> { + let counter = AtomicU64::new(0); + move || { + format!( + "{tag}_{}_{}_{}", + std::process::id(), + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_nanos(), + counter.fetch_add(1, Ordering::SeqCst), + ) + } +} + +async fn harness_client() -> Option { + let (client, connection) = tokio_postgres::connect(&harness_dsn()?, tokio_postgres::NoTls) + .await + .ok()?; + tokio::spawn(async move { + let _ = connection.await; + }); + Some(client) +} + +async fn drop_schema(client: &tokio_postgres::Client, schema: &str) { + if schema == DEFAULT_SCHEMA { + return; + } + let sql = format!( + "DROP SCHEMA IF EXISTS {} CASCADE", + crate::schema::quote_identifier(schema) + ); + let _ = client.batch_execute(&sql).await; +} + +fn test_opts(schema: &str) -> PgOpts { + PgOpts { + schema: schema.to_string(), + ..PgOpts::default() + } +} + +/// Wait for the next event with a bounded deadline — polling loosely +/// so no tight-race assertions ride (the suite's stability posture). +async fn must_recv_event(receiver: &mut dyn EventReceiver, tag: &str) -> alkstore::StreamEvent { + let deadline = tokio::time::Instant::now() + Duration::from_secs(15); + loop { + match receiver.try_recv() { + Ok(Some(event)) => return event, + Ok(None) => {} + Err(e) => panic!("{tag}: receiver errored instead of idling: {e}"), + } + assert!( + tokio::time::Instant::now() < deadline, + "{tag}: event never arrived" + ); + tokio::time::sleep(Duration::from_millis(20)).await; + } +} + +/// Kill this store's listener backend (the POC's +/// `pg_terminate_backend` shape, targeted at the store's unique +/// application_name) — the forwarder-reconnect trigger the +/// gap-healing test rides. +async fn kill_listener_backend(admin: &tokio_postgres::Client, app_name: &str) -> bool { + admin + .query_one( + "SELECT pg_terminate_backend(pid) FROM pg_stat_activity + WHERE application_name = $1 AND pid <> pg_backend_pid()", + &[&app_name], + ) + .await + .unwrap() + .get(0) +} + +/// Acceptance: full `StreamHandle` over the pool — publish returns +/// ascending bigserial offsets, reads decode into `StreamEvent` with +/// the `stream` field carrying the contract name, keys round-trip +/// exactly (`None` vs `Some`), consumers absent from the table read +/// offset 0, offsets save and inspect, and the ordering is `offset +/// ASC` (global FIFO, ADR-015 §4). +#[tokio::test(flavor = "multi_thread")] +async fn publish_read_save_get_round_trip_end_to_end() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("stream")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + let stream = store.stream("events").await.unwrap(); + assert_eq!(stream.name(), "events"); + + let o1 = stream.publish(json!({"n": 1})).await.unwrap(); + let o2 = stream + .publish_with_key(Some("k".to_string()), json!({"n": 2})) + .await + .unwrap(); + let o3 = stream.publish(json!({"n": 3})).await.unwrap(); + assert!(o1 < o2 && o2 < o3, "offsets ascend per publish"); + + let page = stream.read_since(0, 10).await.unwrap(); + assert_eq!(page.len(), 3, "the full stream reads back"); + assert_eq!(page[0].offset, o1); + assert_eq!(page[0].stream, "events", "the stream field is the name"); + assert_eq!(page[0].key, None, "plain publish carries no key"); + assert_eq!(page[1].offset, o2); + assert_eq!(page[1].key, Some("k".to_string())); + assert_eq!(page[2].offset, o3); + assert_eq!(page[2].key, None); + for (i, event) in page.iter().enumerate() { + let decoded: Value = event.payload_as().unwrap(); + assert_eq!(decoded, json!({"n": i as i64 + 1})); + } + + // Cursor reads: offset > cursor, ASC. + let tail = stream.read_since(o1, 10).await.unwrap(); + assert_eq!(tail.len(), 2, "the read anchors strictly past the cursor"); + assert_eq!(tail[0].offset, o2); + + // Consumers absent from the table read 0 (the pre-save state). + let from_absent = stream + .read_from_consumer("fresh-consumer", 10) + .await + .unwrap(); + assert_eq!(from_absent.len(), 3, "absent consumer reads from 0"); + assert_eq!( + stream.get_offset("fresh-consumer").await.unwrap(), + 0, + "get_offset of an absent consumer = 0" + ); + + // Save + inspect; reads from the consumer's checkpoint. + stream.save_offset("c", o2).await.unwrap(); + assert_eq!(stream.get_offset("c").await.unwrap(), o2); + let from_c = stream.read_from_consumer("c", 10).await.unwrap(); + assert_eq!(from_c.len(), 1); + assert_eq!(from_c[0].offset, o3); + + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: two stores on the same schema compose — a subscriber +/// on store A receives store B's publishes through the wake channel +/// (the mechanism-name-is-the-channel realization: both processes' +/// LISTENs ride the same channel name; cross-connection delivery). +#[tokio::test(flavor = "multi_thread")] +async fn wakes_cross_store_connections() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("cross")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let store2 = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + let name = instance_namer("st")(); + let mut rx = store + .stream(&name) + .await + .unwrap() + .subscribe("c") + .await + .unwrap(); + + let published = store2 + .stream(&name) + .await + .unwrap() + .publish(json!({"via": "other-connection"})) + .await + .unwrap(); + + let event = must_recv_event(&mut *rx, "cross-store").await; + assert_eq!(event.offset, published); + assert_eq!(event.stream, name); + assert_eq!( + event.payload_as::().unwrap(), + json!({"via": "other-connection"}) + ); + + store.close(); + store2.close(); + drop(store); + drop(store2); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: the extent guard (ADR-023 §2) — `read_since` / +/// `read_from_consumer` with `limit <= 0` return the empty `Vec` +/// (pg's `LIMIT` with a negative is a server error — the guard makes +/// it unreachable), at trait-impl entry. +#[tokio::test(flavor = "multi_thread")] +async fn extent_guard_reads_limit_zero_and_negative_yield_empty() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("extent")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + stream.publish(json!({"n": 1})).await.unwrap(); + stream.publish(json!({"n": 2})).await.unwrap(); + stream.save_offset("c", 10_000).await.unwrap(); + + for limit in [0i64, -1, i64::MIN] { + let page = stream.read_since(0, limit).await.unwrap(); + assert!( + page.is_empty(), + "read_since(limit={limit}) must read nothing" + ); + let page = stream.read_from_consumer("c", limit).await.unwrap(); + assert!( + page.is_empty(), + "read_from_consumer(limit={limit}) must read nothing" + ); + } + + // The stream was never handed out in bulk: both events remain + // readable with a positive extent. + assert_eq!(stream.read_since(0, 10).await.unwrap().len(), 2); + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: `trim_to` is exact-boundary (`id <= horizon`), returns +/// the deleted count, a negative horizon deletes nothing (boundary +/// args are total, ADR-023 §2 — no guard, pass through), wakes +/// nothing, and surviving offsets are never renumbered (gaps legal). +#[tokio::test(flavor = "multi_thread")] +async fn trim_to_trims_the_exact_boundary_and_never_renumbers() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("trim")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + + let mut offsets = Vec::new(); + for i in 0..5 { + offsets.push(stream.publish(json!({"n": i})).await.unwrap()); + } + + // A subscriber is attached across the trim: no wake may fire + // (ADR-015 §5 — trim wakes nothing; the subscriber must stay + // parked/quiet through it). + let mut rx = store + .stream("s") + .await + .unwrap() + .subscribe("c") + .await + .unwrap(); + // Drain the 5 replayed events. + for _ in 0..5 { + let _ = must_recv_event(&mut *rx, "trim-drain").await; + } + + // Negative horizon: deletes nothing, idempotently (no guard). + assert_eq!( + stream.trim_to(-7).await.unwrap(), + 0, + "a negative horizon deletes nothing" + ); + assert_eq!(stream.read_since(0, 10).await.unwrap().len(), 5); + + // Zero horizon is a no-op on this store (bigserial ids start above 0). + assert_eq!(stream.trim_to(0).await.unwrap(), 0); + + // Exact boundary: horizon = offsets[1] deletes exactly rows + // offsets[0..=1] (two rows), keeping offsets[2..]. + let deleted = stream.trim_to(offsets[1]).await.unwrap(); + assert_eq!(deleted, 2, "trim deletes offset <= horizon exactly"); + + let survivors = stream.read_since(0, 10).await.unwrap(); + assert_eq!(survivors.len(), 3); + let survivor_offsets: Vec = survivors.iter().map(|e| e.offset).collect(); + assert_eq!( + survivor_offsets, + offsets[2..], + "surviving offsets are never renumbered — gaps are legal" + ); + + // Reads from a trimmed-away region resume at the horizon's first + // remaining row (fewer/no rows, never an error). + let resumed = stream.read_since(offsets[0], 10).await.unwrap(); + assert_eq!(resumed.len(), 3); + assert_eq!(resumed[0].offset, offsets[2]); + + // Trim is idempotent below the surviving set. + assert_eq!(stream.trim_to(offsets[1]).await.unwrap(), 0); + + // Saved offsets below the horizon stay valid (the checkpoint is a + // position, not a row reference — reads anchor at it). + stream.save_offset("c2", offsets[0]).await.unwrap(); + // (The subscriber's earlier direct save for "c" stays at 0 — a + // save below the stored checkpoint is the monotone upsert's + // silent no-op; there is no row at 0 either.) + assert_eq!(stream.get_offset("c2").await.unwrap(), offsets[0]); + // A negative horizon still deletes nothing. + assert_eq!(stream.trim_to(-7).await.unwrap(), 0); + // Saved offset below the surviving set survives the trim + // unchanged (a checkpoint is a position, not a row reference). + stream.trim_to(offsets[2]).await.unwrap(); + assert_eq!( + stream.get_offset("c2").await.unwrap(), + offsets[0], + "a saved offset below the horizon survives the trim unchanged" + ); + + // Trim wakes nothing: the subscriber hears no delivery after the + // deletes (a bounded idle window, polled loosely). + let deadline = tokio::time::Instant::now() + Duration::from_millis(400); + while tokio::time::Instant::now() < deadline { + assert_eq!( + rx.try_recv().unwrap(), + None, + "trim must not wake subscribers" + ); + tokio::time::sleep(Duration::from_millis(50)).await; + } + + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: monotone saves (ADR-19 §6) — a direct save below the +/// stored checkpoint is a silent no-op (`Ok(())`, offset unchanged); +/// direct and receiver forms compose without regression. +#[tokio::test(flavor = "multi_thread")] +async fn offset_saves_are_monotone_across_forms() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("monotone")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + + let o1 = stream.publish(json!({"n": 1})).await.unwrap(); + let o2 = stream.publish(json!({"n": 2})).await.unwrap(); + assert!(o1 < o2); + + stream.save_offset("c", o2).await.unwrap(); + // Regression: silent no-op, not an error (the act-differently rule). + stream.save_offset("c", o1).await.unwrap(); + assert_eq!( + stream.get_offset("c").await.unwrap(), + o2, + "a save below the stored checkpoint does not rewind the cursor" + ); + + stream.save_offset("c", o2).await.unwrap(); + assert_eq!(stream.get_offset("c").await.unwrap(), o2); + + // The receiver form drives the same op: subscribing a fresh + // consumer, reading to tail, saving through the receiver — then a + // direct regression save cannot undo it. + let mut rx = store + .stream("s") + .await + .unwrap() + .subscribe("r") + .await + .unwrap(); + let last = must_recv_event(&mut *rx, "monotone-receiver").await; + rx.save_offset().unwrap(); + assert_eq!(rx.offset(), last.offset); + let stored = store + .stream("s") + .await + .unwrap() + .get_offset("r") + .await + .unwrap(); + assert_eq!(stored, last.offset, "the receiver save checkpointed"); + + stream.save_offset("r", o1).await.unwrap(); + let after = store + .stream("s") + .await + .unwrap() + .get_offset("r") + .await + .unwrap(); + assert_eq!( + after, last.offset, + "the direct form cannot regress the receiver's" + ); + + // Receiver save before any yield is a no-op (the stored checkpoint + // untouched — ADR-021 §5). + let fresh_stream = store.stream("s").await.unwrap(); + let mut fresh = fresh_stream.subscribe("pre-yield").await.unwrap(); + fresh.save_offset().unwrap(); + assert_eq!(fresh.offset(), 0); + assert_eq!(fresh_stream.get_offset("pre-yield").await.unwrap(), 0); + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: `subscribe` replays from the *stored* offset (a +/// consumer with a prior checkpoint resumes there, not from 0), then +/// delivers post-attach events wake-driven — events published after +/// attach arrive through the receiver without polling. +#[tokio::test(flavor = "multi_thread")] +async fn subscribe_replays_from_stored_offset_then_delivers_wake_driven() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("subscribe")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + let pre = store.stream("s").await.unwrap(); + let o1 = pre.publish(json!({"n": 1})).await.unwrap(); + let o2 = pre.publish(json!({"n": 2})).await.unwrap(); + let o3 = pre.publish(json!({"n": 3})).await.unwrap(); + drop(pre); + + // A stored checkpoint below the tail: replay starts there (o2 is + // the first event after the checkpoint). + let handle = store.stream("s").await.unwrap(); + handle.save_offset("c", o1).await.unwrap(); + let mut rx = handle.subscribe("c").await.unwrap(); + drop(handle); + + let first = must_recv_event(&mut *rx, "subscribe-replay").await; + assert_eq!( + first.offset, o2, + "replay resumes past the stored checkpoint" + ); + let second = must_recv_event(&mut *rx, "subscribe-replay").await; + assert_eq!(second.offset, o3); + assert_eq!(second.stream, "s"); + + // Post-attach publishes arrive wake-driven — no polling needed. + let late = store + .stream("s") + .await + .unwrap() + .publish_with_key(Some("late".to_string()), json!({"n": 4})) + .await + .unwrap(); + let wake_event = must_recv_event(&mut *rx, "subscribe-wake").await; + assert_eq!(wake_event.offset, late); + assert_eq!(wake_event.key.as_deref(), Some("late")); + + // Receiver checkpoint composes with the substrate's monotone upsert. + rx.save_offset().unwrap(); + let stored = store + .stream("s") + .await + .unwrap() + .get_offset("c") + .await + .unwrap(); + assert_eq!(stored, wake_event.offset); + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: the gap-healing property (the pg arm's distinguishing +/// pin) — a subscriber stays open across a forwarder reconnect, and +/// its re-drain on the reconnect-wake reads the durable rows a gap +/// hid: publishes made while the listener connection was down arrive +/// after the reconnect (no event is lost by the wake gap). The +/// receiver never sees `None` across the reconnect (close is +/// engine-shutdown-only, ADR-021 §5). +#[tokio::test(flavor = "multi_thread")] +async fn subscriber_survives_reconnect_and_heals_gap_by_redrain() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("gap-heal")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let admin = harness_client().await.unwrap(); + + let name = instance_namer("gap-st")(); + let handle = store.stream(&name).await.unwrap(); + + // Pre-gap delivery verified. + let pre = handle.publish(json!({"n": 0})).await.unwrap(); + let mut rx = handle.subscribe("c").await.unwrap(); + let first = must_recv_event(&mut *rx, "gap-heal-pre").await; + assert_eq!(first.offset, pre); + + // Kill the listener backend, then publish during the gap: the + // forwarder is reconnecting; the wake (best-effort) is lost in + // the gap — the durable row is the truth. + assert!( + kill_listener_backend(&admin, store.listener_application_name()).await, + "the listener backend was found and killed" + ); + let gap_offset = handle.publish(json!({"n": 1})).await.unwrap(); + + // Drain to pre-gap tail; then the gap event must arrive (the + // reconnect-wake triggered a re-drain that reads the durable row). + let gap_event = must_recv_event(&mut *rx, "gap-heal-event").await; + assert_eq!( + gap_event.offset, gap_offset, + "the gap event is read back by the re-drain — the row is the truth" + ); + assert_eq!(gap_event.payload_as::().unwrap(), json!({"n": 1})); + + // The receiver is not closed by the reconnect (the pg arm: closes + // at engine shutdown only). Post-reconnect deliveries keep + // flowing. + let post = handle.publish(json!({"n": 2})).await.unwrap(); + let post_event = must_recv_event(&mut *rx, "gap-heal-post").await; + assert_eq!(post_event.offset, post); + + store.close(); + drop(store); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: the terminal close arm — engine shutdown closes the +/// subscription (`recv() -> None`, never reopens; `try_recv` reports +/// `Err(Closed)`), exactly the SQLite twin's shape driven by the +/// store's close instead of watcher death. +#[tokio::test(flavor = "multi_thread")] +async fn engine_shutdown_closes_the_subscription_terminally() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("shutdown")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + stream.publish(json!({"n": 1})).await.unwrap(); + let mut rx = stream.subscribe("c").await.unwrap(); + + let event = must_recv_event(&mut *rx, "shutdown").await; + assert_eq!(event.stream, "s"); + + store.close(); + let none = rx.recv().await; + assert!( + none.is_none(), + "engine shutdown closes the receiver: got {none:?}" + ); + assert!( + matches!(rx.try_recv(), Err(Error::Closed)), + "after shutdown, try_recv is Err(Closed) — not idle" + ); + assert!(rx.recv().await.is_none(), "a closed receiver never reopens"); + + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: a subscribe attached while the stream is at tail idles +/// (`Ok(None)`), then delivers when the next publish lands. +#[tokio::test(flavor = "multi_thread")] +async fn subscribe_at_tail_idles_then_delivers() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("tail")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + + let mut rx = stream.subscribe("c").await.unwrap(); + // A bounded idle window past attach (a wake would deliver through + // the channel's mpsc; nothing publishes during it). + tokio::time::sleep(Duration::from_millis(200)).await; + assert_eq!(rx.try_recv().unwrap(), None, "at tail, the receiver idles"); + + let o = stream.publish(json!({"n": 1})).await.unwrap(); + let event = must_recv_event(&mut *rx, "tail-idle").await; + assert_eq!(event.offset, o); + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: two consumers subscribe independently — each drains a +/// copy of the stream from its own position; one's receiver save never +/// gates the other's delivery (saved checkpoints never gate reads). +#[tokio::test(flavor = "multi_thread")] +async fn subscribers_are_independent_copies() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("two-cons")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + let o1 = stream.publish(json!({"n": 1})).await.unwrap(); + let o2 = stream.publish(json!({"n": 2})).await.unwrap(); + + let handle = store.stream("s").await.unwrap(); + let mut a = handle.subscribe("a").await.unwrap(); + let mut b = handle.subscribe("b").await.unwrap(); + drop(handle); + + // a drains and checkpoints both; b must still receive both (a's + // position is a's, not shared). + let _ = must_recv_event(&mut *a, "two-a-1").await; + let _ = must_recv_event(&mut *a, "two-a-2").await; + a.save_offset().unwrap(); + + let b1 = must_recv_event(&mut *b, "two-b-1").await; + let b2 = must_recv_event(&mut *b, "two-b-2").await; + assert_eq!(b1.offset, o1); + assert_eq!(b2.offset, o2); + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: receiver-drop unregisters the channel — UNLISTEN at +/// the last-subscriber drop (the receiver's Drop; the notify task's +/// refcount rule shared). Probed through the raw forwarder fanout +/// (which sees every channel this listener session LISTENs). +#[tokio::test(flavor = "multi_thread")] +async fn receiver_drop_unregisters_the_stream_channel() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("drop-unreg")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let admin = harness_client().await.unwrap(); + + let name = instance_namer("drop-st")(); + let handle = store.stream(&name).await.unwrap(); + let rx = handle.subscribe("c").await.unwrap(); + + // The channel is registered: a raw admin notify on it is seen by + // this session's listener (the raw fanout). + let mut raw_rx = store.forwarder().subscribe().unwrap(); + admin + .query_one("SELECT pg_notify($1, '')", &[&name]) + .await + .unwrap(); + let heard = raw_fanout_hears(&mut raw_rx, &name, Duration::from_secs(2)).await; + assert!(heard, "the stream channel is LISTENed while subscribed"); + + // Drop the receiver: UNLISTEN — the raw fanout hears nothing for + // the channel anymore (retried until the async UNLISTEN lands). + drop(rx); + let mut unlistened = false; + for _ in 0..25 { + admin + .query_one("SELECT pg_notify($1, '')", &[&name]) + .await + .unwrap(); + let heard = raw_fanout_hears(&mut raw_rx, &name, Duration::from_millis(120)).await; + if !heard { + unlistened = true; + break; + } + } + assert!( + unlistened, + "the receiver's drop must UNLISTEN the stream channel" + ); + + store.close(); + drop(store); + drop_schema(&admin, &schema).await; +} + +/// Whether the raw fanout hears `channel` within `timeout`. +async fn raw_fanout_hears( + rx: &mut tokio::sync::broadcast::Receiver, + channel: &str, + timeout: Duration, +) -> bool { + let deadline = tokio::time::Instant::now() + timeout; + loop { + let remaining = deadline.saturating_duration_since(tokio::time::Instant::now()); + if remaining.is_zero() { + return false; + } + match tokio::time::timeout(remaining, rx.recv()).await { + Ok(Ok(n)) if n.channel == channel => return true, + Ok(Ok(_)) => continue, + Ok(Err(tokio::sync::broadcast::error::RecvError::Lagged(_))) => continue, + Ok(Err(tokio::sync::broadcast::error::RecvError::Closed)) => return false, + Err(_) => return false, + } + } +} + +/// Acceptance: key round-trips exactly — `None` vs `Some` preserved +/// on every read form (direct page and receiver delivery), and the +/// `stream` field carries the contract name. +#[tokio::test(flavor = "multi_thread")] +async fn keys_round_trip_exactly_and_stream_field_carries_the_name() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("keys")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("the-name").await.unwrap(); + + let plain = stream.publish(json!({})).await.unwrap(); + let keyed = stream + .publish_with_key(Some("k".to_string()), json!({})) + .await + .unwrap(); + + let page = stream.read_since(0, 10).await.unwrap(); + fn find(events: &[alkstore::StreamEvent], offset: i64) -> Option<&alkstore::StreamEvent> { + events.iter().find(|e| e.offset == offset) + } + assert_eq!(find(&page, plain).unwrap().key, None); + assert_eq!(find(&page, keyed).unwrap().key, Some("k".to_string())); + assert!(page.iter().all(|e| e.stream == "the-name")); + + let mut rx = store + .stream("the-name") + .await + .unwrap() + .subscribe("c") + .await + .unwrap(); + let e1 = must_recv_event(&mut *rx, "keys-receiver").await; + let e2 = must_recv_event(&mut *rx, "keys-receiver").await; + assert_eq!(e1.key, None); + assert_eq!(e2.key, Some("k".to_string())); + assert_eq!(e1.stream, "the-name"); + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: the auto-commit publish path carries the stored-bytes +/// serialization exactly (ADR-020 §4): the decoded payload +/// round-trips byte-exactly through the table (the pg engine stores +/// the `encode_payload` result unmodified in its BYTEA column — the +/// contract-suite's byte-identical-rows row). +#[tokio::test(flavor = "multi_thread")] +async fn publish_stores_the_exact_serialized_bytes() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("codec")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + + // Duplicated keys fold in the serde_json Value (parse-level), but + // the *stored* bytes must be exactly what `encode_payload` + // returned — verified here via the raw row for a tricky payload + // (deep nesting, floats). + let payload = json!({ + "deep": {"nest": [1, 2, {"x": 2.5, "y": null}], "z": true}, + "uni": "héllo ✓" + }); + let _offset = stream.publish(payload.clone()).await.unwrap(); + + let mut expected_bytes = Vec::new(); + serde_json::to_writer(&mut expected_bytes, &payload).unwrap(); + + let admin = harness_client().await.unwrap(); + let table = crate::schema::QualifiedTable { + schema: &schema, + name: crate::schema::tables::EVENTS, + } + .to_string(); + let row = admin + .query_one( + &format!("SELECT payload FROM {table} WHERE stream = $1"), + &[&"s"], + ) + .await + .unwrap(); + let stored: Vec = row.get(0); + assert_eq!( + stored, expected_bytes, + "the stored bytes are exactly the serde_json serialization (ADR-020 §4)" + ); + + // And the decoded form matches the original Value. + let page = stream.read_since(0, 10).await.unwrap(); + assert_eq!(page.len(), 1); + assert_eq!(page[0].payload_as::().unwrap(), payload); + + store.close(); + drop(store); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: entry-point validation — stream names validate as +/// shared-namespace kinds at the constructor; consumers as local names +/// (non-empty only; reserved prefixes legal for consumers); empty- +/// `Some` keys are `InvalidName`. Validation fires before any round +/// trip: rejected constructors yield no handle, and rejected consumer +/// names never touch the tables. +#[tokio::test(flavor = "multi_thread")] +async fn entry_points_validate_stream_names_and_consumers() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("valid")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + + // Stream constructors: empty/whitespace → InvalidName; reserved + // prefix → ReservedName. + for name in ["", " ", "__alkstore_x"] { + let handle = match store.stream(name).await { + Err(e) => e, + Ok(_) => panic!("stream({name:?}) must reject by validation"), + }; + match handle { + Error::InvalidName { name: n } => assert_eq!(name, n), + Error::ReservedName { name: n } => assert_eq!(name, n), + other => panic!("stream({name:?}) must reject by validation, got {other}"), + } + } + + let stream = store.stream("s").await.unwrap(); + + // Empty-`Some` keys are InvalidName (the tx path's rule). + let err = stream + .publish_with_key(Some(String::new()), json!({})) + .await + .unwrap_err(); + match err { + Error::InvalidName { name } => assert_eq!(name, ""), + other => panic!("empty-Some key must be InvalidName, got {other}"), + } + let err = stream + .publish_with_key(Some(" ".to_string()), json!({})) + .await + .unwrap_err(); + assert!(matches!(err, Error::InvalidName { .. })); + + // Consumers: non-empty only — reserved prefixes are legal. + for consumer in ["", " "] { + let err = stream.read_from_consumer(consumer, 5).await.unwrap_err(); + assert!(matches!(err, Error::InvalidName { .. })); + let err = stream.save_offset(consumer, 1).await.unwrap_err(); + assert!(matches!(err, Error::InvalidName { .. })); + let err = stream.get_offset(consumer).await.unwrap_err(); + assert!(matches!(err, Error::InvalidName { .. })); + assert!(stream.subscribe(consumer).await.is_err()); + } + + // A legal-but-reserved consumer name works (non-empty only). + let o = stream.publish(json!({"n": 9})).await.unwrap(); + stream.save_offset("__alkstore_local", o).await.unwrap(); + assert_eq!(stream.get_offset("__alkstore_local").await.unwrap(), o); + + // The rejected consumer saves never touched the table. + assert_eq!(stream.get_offset("never-saved").await.unwrap(), 0); + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: `StreamHandle` ops on a closed store fail closed with +/// `Error::Database` (the engine-wide closed-store posture); a +/// subscribe on a closed store fails too. +#[tokio::test(flavor = "multi_thread")] +async fn closed_store_fails_stream_ops() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("closed-st")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + stream.publish(json!({"n": 1})).await.unwrap(); + let stored_handle = store.stream("s").await.unwrap(); + store.close(); + + let err = stream.publish(json!({})).await.unwrap_err(); + assert!(matches!(err, Error::Database(_)), "publish: {err}"); + let err = stream.read_since(0, 5).await.unwrap_err(); + assert!(matches!(err, Error::Database(_)), "read_since: {err}"); + let err = stream.read_from_consumer("c", 5).await.unwrap_err(); + assert!( + matches!(err, Error::Database(_)), + "read_from_consumer: {err}" + ); + let err = stream.save_offset("c", 1).await.unwrap_err(); + assert!(matches!(err, Error::Database(_)), "save_offset: {err}"); + let err = stream.get_offset("c").await.unwrap_err(); + assert!(matches!(err, Error::Database(_)), "get_offset: {err}"); + let err = stream.trim_to(10).await.unwrap_err(); + assert!(matches!(err, Error::Database(_)), "trim_to: {err}"); + assert!(stored_handle.subscribe("c").await.is_err()); + + // New constructors on a closed store fail too. + let err = match store.stream("s").await { + Err(e) => e, + Ok(_) => panic!("stream on a closed store must fail closed"), + }; + assert!(matches!(err, Error::Database(_))); + + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: the durable subscription survives store restart — a +/// fresh open over the same schema resumes from the saved checkpoint +/// (durable consumption; offsets are positions in the durable log). +#[tokio::test(flavor = "multi_thread")] +async fn subscription_state_survives_store_restart() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("restart-st")(); + let saved_offset; + + { + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + let o1 = stream.publish(json!({"n": 1})).await.unwrap(); + let _ = stream.publish(json!({"n": 2})).await.unwrap(); + stream.save_offset("c", o1).await.unwrap(); + saved_offset = stream.get_offset("c").await.unwrap(); + store.close(); + drop(store); + } + + { + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + // Checkpoint durable across the restart. + assert_eq!(stream.get_offset("c").await.unwrap(), saved_offset); + // Replay resumes from it — the second event is the first read. + let page = stream.read_from_consumer("c", 10).await.unwrap(); + assert_eq!(page.len(), 1, "the resumed read starts past the checkpoint"); + // A fresh subscription replays from the saved offset. + let mut rx = stream.subscribe("c").await.unwrap(); + let event = must_recv_event(&mut *rx, "restart-subscribe").await; + assert_eq!(event.offset, saved_offset + 1); + drop(rx); + drop(stream); + store.close(); + drop(store); + } + + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: the tx seam and the handle compose — a `publish_tx` +/// inside a committed business transaction is visible to the handle's +/// reads and wakes a live subscriber; a rolled-back `publish_tx` +/// ghosts (no read, no wake). Also: the tx save path and the handle's +/// monotone op are the same upsert (a tx save cannot regress a +/// handle save and vice versa). +#[tokio::test(flavor = "multi_thread")] +async fn tx_publishes_compose_with_the_handle() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("tx-compose-st")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let mut rx = store + .stream("s") + .await + .unwrap() + .subscribe("c") + .await + .unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + let tx_offset = tx.publish_tx("s", json!({"via": "tx"})).await.unwrap(); + tx.commit().await.unwrap(); + + let event = must_recv_event(&mut *rx, "tx-compose").await; + assert_eq!(event.offset, tx_offset); + assert_eq!(event.payload_as::().unwrap(), json!({"via": "tx"})); + + // Rollback: the publish ghosts. + let mut tx = store.begin_tx().await.unwrap(); + tx.publish_tx("s", json!({"via": "ghost"})).await.unwrap(); + drop(tx); + tokio::time::sleep(Duration::from_millis(200)).await; + while rx.try_recv().unwrap().is_some() {} + tokio::time::sleep(Duration::from_millis(200)).await; + assert_eq!( + rx.try_recv().unwrap(), + None, + "a rolled-back publish must not deliver" + ); + + // The handle read confirms the table state. + let stream = store.stream("s").await.unwrap(); + let page = stream.read_since(0, 10).await.unwrap(); + assert_eq!(page.len(), 1, "the ghost never landed"); + assert_eq!(page[0].offset, tx_offset); + + // The tx save path is the same monotone upsert the handle uses. + let mut tx = store.begin_tx().await.unwrap(); + tx.save_offset_tx("s", "c-tx", tx_offset).await.unwrap(); + tx.commit().await.unwrap(); + stream.save_offset("c-tx", 0).await.unwrap(); + assert_eq!( + stream.get_offset("c-tx").await.unwrap(), + tx_offset, + "the direct form cannot regress the tx form's save" + ); + + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: the receiver's `read_since` is offset-anchored and +/// independent of its position (saved checkpoints do not gate reads — +/// replay is a read concern); extent-guarded like the direct form. +#[tokio::test(flavor = "multi_thread")] +async fn receiver_read_since_is_offset_anchored_and_extent_guarded() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("rx-read")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + let o1 = stream.publish(json!({"n": 1})).await.unwrap(); + let o2 = stream.publish(json!({"n": 2})).await.unwrap(); + + let mut rx = stream.subscribe("c").await.unwrap(); + let _ = must_recv_event(&mut *rx, "receiver-read").await; + let _ = must_recv_event(&mut *rx, "receiver-read").await; + rx.save_offset().unwrap(); + assert_eq!(rx.offset(), o2); + + // The receiver replays from any offset regardless of its position. + let replay = rx.read_since(o1 - 1, 10).await.unwrap(); + assert_eq!(replay.len(), 2, "the receiver replays past its position"); + assert_eq!(replay[0].offset, o1); + // Extent guard on the receiver form too. + assert_eq!(rx.read_since(0, 0).await.unwrap().len(), 0); + assert_eq!(rx.read_since(0, -3).await.unwrap().len(), 0); + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: the receiver's `save_offset` checkpoints only the +/// last-*yielded* event — a stalled consumer (not reading from the +/// bounded channel) never checkpoints events the bridge drained but +/// the consumer did not observe (the internal cursor is separate from +/// the position). Driven by subscribing, publishing several events, +/// consuming none, and saving: the stored offset stays at the last +/// actually-yielded event. +#[tokio::test(flavor = "multi_thread")] +async fn stalled_consumer_saves_only_observed_events() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("stall")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let name = instance_namer("stall-st")(); + let handle = store.stream(&name).await.unwrap(); + + let pre = handle.publish(json!({"n": 0})).await.unwrap(); + let mut rx = handle.subscribe("c").await.unwrap(); + let observed = must_recv_event(&mut *rx, "stall").await; + assert_eq!(observed.offset, pre); + rx.save_offset().unwrap(); + + // Publish several more; do NOT read them. The bridge drains them + // into the bounded channel (possibly stalling backpressure), but + // the receiver's position must not move. + for i in 1..=3 { + handle.publish(json!({"n": i})).await.unwrap(); + } + // Wait for the bridge to drain the new rows into the channel. + tokio::time::sleep(Duration::from_millis(300)).await; + rx.save_offset().unwrap(); + assert_eq!( + rx.offset(), + observed.offset, + "an unread event never advances the position" + ); + let stored = store + .stream(&name) + .await + .unwrap() + .get_offset("c") + .await + .unwrap(); + assert_eq!( + stored, observed.offset, + "the stored checkpoint covers only consumed events" + ); + + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: backpressure — a consumer that stops reading for a +/// while doesn't lose events; when it resumes, delivery picks up +/// exactly at the last-yielded offset in ASC order. The bounded +/// channel plus backpressure-aware send parks the bridge; the resume +/// drains the backlog in order. +#[tokio::test(flavor = "multi_thread")] +async fn backpressure_holds_delivery_resumes_exactly() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("backpressure")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let name = instance_namer("bp-st")(); + let handle = store.stream(&name).await.unwrap(); + + let mut rx = handle.subscribe("c").await.unwrap(); + // Publish a burst bigger than the smoothing buffer. + let mut expected = Vec::new(); + for i in 0..48 { + expected.push(handle.publish(json!({"i": i})).await.unwrap()); + } + + let mut got = Vec::new(); + for _ in 0..48 { + got.push(must_recv_event(&mut *rx, "backpressure").await.offset); + } + assert_eq!( + got, expected, + "a stalled consumer's backlog is delivered exactly, in ASC order" + ); + + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} + +/// Acceptance: `recv()`'s `Err` arm carries `Database` only (ADR-021 +/// §5) — the receiver maps the bridge's read failures to the opaque +/// fallback before delivery. The compile-time shape is pinned by the +/// trait; this test pins the observable no-`Codec`-delivery rule +/// structurally over the healthy path (a Codec never crosses the +/// receiver — decoded events are values). +#[tokio::test(flavor = "multi_thread")] +async fn receiver_delivers_events_not_error_kinds() { + let Some(dsn) = harness_dsn() else { + eprintln!("skip: no harness server"); + return; + }; + let schema = instance_namer("err-arm")(); + let store = open_store(&dsn, test_opts(&schema)).await.unwrap(); + let stream = store.stream("s").await.unwrap(); + + let o1 = stream.publish(json!({"n": 1})).await.unwrap(); + let o2 = stream.publish(json!({"n": 2})).await.unwrap(); + let mut rx = stream.subscribe("c").await.unwrap(); + let e1 = must_recv_event(&mut *rx, "err-arm-1").await; + let e2 = must_recv_event(&mut *rx, "err-arm-2").await; + assert_eq!(e1.offset, o1); + assert_eq!(e2.offset, o2); + + // The receiver's own decode convenience carries Codec (that's + // where the codec error lives — never at the bridge). + let decoded: Value = e1.payload_as().unwrap(); + assert_eq!(decoded, json!({"n": 1})); + store.close(); + drop(store); + let admin = harness_client().await.unwrap(); + drop_schema(&admin, &schema).await; +} diff --git a/alkstore-postgres/src/stream.rs b/alkstore-postgres/src/stream.rs new file mode 100644 index 0000000..da21738 --- /dev/null +++ b/alkstore-postgres/src/stream.rs @@ -0,0 +1,856 @@ +//! The stream mechanism's Postgres arm (`pg-engine-streams`; ADR-015; +//! ADR-019 §1/§6; ADR-021 §5; ADR-023 §2): +//! +//! - `publish` / `publish_with_key` — the auto-commit append paths: +//! pool checkout, one parameterized `INSERT INTO … RETURNING id` +//! (`id BIGSERIAL` is the assigned offset — engine-assigned, +//! monotone per stream, immutable; the nullable `"key"` column is +//! carried metadata with the non-empty-key rule: an empty `Some` key +//! is `InvalidName`, matching the tx path), `encode_payload` at the +//! seam with the typed `Codec` error `?`-ed, return to pool. Then +//! **wake**: a `SELECT pg_notify($1, '')` on the stream's wake +//! channel — the mechanism-name-is-the-channel realization (the +//! plan's Wave 4 section) — as a *separate statement* on the same +//! pooled connection (auto-commit: its own implicit transaction; +//! the wake is best-effort, not commit-atomic with the insert — +//! the no-replay hole covers any gap between the two statements; +//! **the durable row is the truth** and every read path re-reads +//! from it). A failed wake is logged and swallowed — the durable +//! row already landed, and wake feeds must never fail a publish. +//! - `read_since` / `read_from_consumer` — pool reads, `ORDER BY id +//! ASC` (global FIFO per stream, ADR-015 §4), decoded into +//! [`StreamEvent`] via the core constructor's `from_row` (the +//! decode owner lives here, shared with the tx path). The extent +//! guard (ADR-023 §2) short-circuits `limit <= 0` to the empty +//! `Vec` at trait-impl entry, before any round trip — pg's `LIMIT` +//! with a negative is a server error, which the guard makes +//! unreachable. `read_from_consumer` anchors at the consumer's +//! stored checkpoint (absent consumer = 0, ADR-019 §1). +//! - `save_offset` / `get_offset` — the monotone upsert (a save below +//! the stored checkpoint is a silent no-op — one owner of the +//! monotonicity rule, the `GREATEST`-clamped + `WHERE`-guarded +//! upsert `save_offset_tx` also rides) 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 per ADR-023 §2: a negative horizon deletes +//! nothing, idempotently). Surviving events keep their offsets +//! (gaps legal, never renumbered). +//! - `subscribe(consumer)` — the durable receiver, the structural +//! twin of the SQLite engine's bridge with the pg deltas (the +//! wake substrate is the forwarder's fanout, natively async — no +//! blocking bridge thread): +//! - **Attach-before-read**: the forwarder channel subscription is +//! taken *before* the initial cursor read, so a commit cannot +//! fall between the attach read and the wake wait (a wake for a +//! commit already re-read sits queued; a re-drain that sees +//! nothing new costs an empty pass, never a missed event). +//! - Initial cursor read of the consumer's *stored* offset; page- +//! drain to the tail (`id ASC`); park on the wake feed (this +//! stream's wake channel plus the reserved reconnect-wake — +//! both trigger re-drains); re-drain on wake. Wakes are +//! best-effort hints (ADR-006): overtrigger on purpose, coalesc- +//! ed by the re-read from the in-memory cursor. +//! - **Close semantics (ADR-021 §5's pg arm)**: the receiver stays +//! open across forwarder reconnects — the synthetic +//! reconnect-wake arrives on the wake feed and simply triggers a +//! re-drain (which *is* the gap healing: the no-replay hole's +//! missed wakes are covered by the re-drain reading the durable +//! rows); closes only at engine shutdown (`recv() -> None`, +//! terminal). +//! - **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 — separate from the position — +//! so a stalled consumer's `save_offset()` never checkpoints +//! unobserved events. Backpressure: a bounded tokio channel and +//! backpressure-aware `send().await` — a stalled consumer parks +//! the bridge task and delivery resumes exactly where it stopped. +//! - The receiver's `read_since` is offset-anchored and extent- +//! guarded like the direct form; its sync `save_offset()` sends +//! the last-yielded offset through the same monotone op (ADR-019 +//! §6). The trait's sync signature on this natively-async engine +//! is bridged by a dedicated wake runtime (see +//! [`WAKE_RUNTIME`]) — the honest sync→async seam, never +//! `spawn_blocking`-reentrancy tricks. +//! - `recv()`'s `Err` carries `Database` only (ADR-021 §5 — non- +//! `Database` read errors are remapped 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 the shutdown-close disconnect. + +use std::sync::Arc; +use std::sync::OnceLock; +use std::sync::atomic::{AtomicBool, Ordering}; + +use deadpool_postgres::Pool; +use serde_json::Value; +use tokio::runtime::{Handle, Runtime, RuntimeFlavor}; + +use alkstore::{ + BoxedFuture, Error, EventReceiver, Result, StreamEvent, StreamHandle, validate_local_name, + validate_shared_name, +}; + +use crate::forwarder::{Forwarder, RawNotification}; +use crate::schema::{QualifiedTable, tables}; +use crate::seam::{database_error, pg_error, pool_error}; +use crate::tx::encode_payload_bytes; + +/// How many events a wake-driven re-read fetches per page while +/// draining to the tail. +const SUBSCRIBE_PAGE: i64 = 256; + +/// The bridge's tokio-channel capacity — the smoothing buffer between +/// the bridge task's re-read loop and the consumer. +const EVENT_CHANNEL_CAPACITY: usize = 16; + +/// The per-process idle wake runtime the receiver's **sync** +/// `save_offset` (and any other sync receiver op needing async +/// storage work) drives its round trip through. +/// +/// Why this exists: the trait's `save_offset(&mut self)` is synchronous, +/// but on this engine the checkpoint save is an async pool round trip — +/// there is no blocking writer slot to lease (the pg posture is no +/// `spawn_blocking` seam and pool clients are runtime-tied). The bridge +/// options were probed exhaustively: +/// +/// - `Handle::current().block_on` from within an async context — panics +/// ("cannot start a runtime from within a runtime") on both flavors; +/// - `Handle::block_on` from a *foreign* runtime's worker/blocking +/// threads — also panics (it blocks the calling thread while that +/// thread belongs to another runtime's driver, the same enter guard); +/// - `Runtime::new().block_on` nested inside a worker thread — panics +/// likewise, and builds a fresh runtime per call (unbounded cost). +/// +/// The survivor is a dedicated, runtime-*independent* driver: a fresh +/// current-thread runtime per blocked save fails the budget rule, so +/// instead a single lazily-built multi-thread runtime is stored +/// process-wide ([`WAKE_RUNTIME`]); the save calls +/// `wake_runtime().block_on(save_future)` **from a spawned std +/// thread** — the one context that is never inside any runtime (a +/// task's stack runs *on* a worker thread; a fresh thread starts +/// outside every enter guard). Concurrent saves park on their own +/// threads and multiplex through one shared runtime; a save issued +/// from a foreign runtime's `spawn_blocking` or from a plain thread +/// skips the spawn and blocks directly. +/// +/// The runtime carries **no engine state** — it is purely the sync +/// call's exec context for its one `.await`; tokio-postgres clients +/// bind to whatever runtime drives *them*, and each save checks out +/// its own pool connection inside the driven future, on the runtime +/// that will drive its connection's I/O. Idle cost: one parked +/// runtime (one blocked driver thread, ~nothing else). +fn wake_runtime() -> &'static Runtime { + static WAKE_RUNTIME: OnceLock = OnceLock::new(); + WAKE_RUNTIME.get_or_init(|| { + tokio::runtime::Builder::new_multi_thread() + .worker_threads(1) + .enable_all() + .build() + .expect("the wake runtime builds (a fixed 1-worker multi-thread runtime)") + }) +} + +/// Drive `fut` to completion from a **synchronous** context, on the +/// dedicated idle wake runtime +/// (see [`wake_runtime`] for the full rationale). +/// +/// The calling context decides the driving posture: +/// +/// - plain thread / foreign-runtime blocking thread → `block_on` +/// directly on the calling thread (already outside every runtime — +/// no spawn, no extra thread); +/// - inside a multi-thread runtime's worker (the async consumer calling +/// `save_offset().unwrap()` in its task) → `block_in_place` frees +/// this worker, then `block_on` on the dedicated runtime (the +/// probed-safe shape — a runtime's `block_on` is legal from a +/// different runtime's `block_in_place` context); +/// - inside a current-thread runtime's worker → a spawned std thread +/// does the blocking (the only escape from a current-thread worker — +/// `block_in_place` panics there), the caller joins it. +fn drive_sync(fut: F) -> Result +where + F: std::future::Future + Send, + F::Output: Send, +{ + let in_runtime = Handle::try_current().is_ok(); + if !in_runtime { + return Ok(wake_runtime().block_on(fut)); + } + let handle = Handle::current(); + match handle.runtime_flavor() { + RuntimeFlavor::MultiThread => { + tokio::task::block_in_place(move || Ok(wake_runtime().block_on(fut))) + } + RuntimeFlavor::CurrentThread => { + let boxed: std::pin::Pin + Send>> = + Box::pin(fut); + let joined = std::thread::scope(|scope| { + let inner = scope.spawn(move || wake_runtime().block_on(boxed)); + inner + .join() + .map_err(|_| database_error("the receiver save's driver thread panicked")) + })?; + Ok(joined) + } + flavor => Err(database_error(format!( + "the receiver save cannot bridge from runtime flavor {flavor:?}" + ))), + } +} + +/// The stream-scoped handle (see the module docs). +pub(crate) struct PgStreamHandle { + name: String, + pool: Pool, + schema: String, + forwarder: Arc, + closed: Arc, +} + +impl std::fmt::Debug for PgStreamHandle { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("PgStreamHandle") + .field("name", &self.name) + .finish_non_exhaustive() + } +} + +impl PgStreamHandle { + pub(crate) fn new( + name: String, + pool: Pool, + schema: String, + forwarder: Arc, + closed: Arc, + ) -> Self { + Self { + name, + pool, + schema, + forwarder, + closed, + } + } +} + +/// Qualified table reference for this engine-owned schema. +fn table(schema: &str, name: &str) -> String { + QualifiedTable { schema, name }.to_string() +} + +/// The auto-commit publish path (see the module docs). Stream-name +/// validation fires at the entry point, before any round trip and +/// before the payload is touched (ADR-008 §4); a present key must be +/// non-empty (the tx path's rule). +fn publish_with_key( + pool: Pool, + schema: String, + closed: Arc, + stream: String, + key: Option, + payload: Value, +) -> BoxedFuture<'static, Result> { + if let Err(e) = validate_shared_name(&stream) { + return Box::pin(async move { Err(e) }); + } + let key = match key { + Some(k) if k.is_empty() || k.trim().is_empty() => { + return Box::pin(async move { Err(Error::InvalidName { name: k }) }); + } + other => other, + }; + Box::pin(async move { + if closed.load(Ordering::Acquire) { + return Err(database_error("the store is closed")); + } + let bytes = encode_payload_bytes(&payload)?; + let created_at = crate::resolution::now_unix(); + let conn = pool.get().await.map_err(pool_error)?; + let events = table(&schema, tables::EVENTS); + let row = conn + .query_one( + &format!( + "INSERT INTO {events} (stream, \"key\", payload, created_at) + VALUES ($1, $2, $3, $4) + RETURNING id" + ), + &[&stream, &key, &bytes, &created_at], + ) + .await + .map_err(pg_error)?; + let offset: i64 = row.get(0); + // Then wake: pg_notify on the stream's wake channel — a + // separate statement (its own implicit tx on this + // auto-commit connection; best-effort, not commit-atomic + // with the insert — the no-replay hole covers a gap between + // them; the durable row is the truth and every reader + // re-reads from it). A wake failure is logged and swallowed: + // the durable row landed; no wake fee may ever fail (or + // stall) a publish. + if let Err(wake_err) = conn.query_one("SELECT pg_notify($1, '')", &[&stream]).await { + eprintln!( + "alkstore-postgres: stream publish wake failed for {stream:?} \ + (the durable row landed; readers re-drain on their next hint): {wake_err}" + ); + } + Ok(offset) + }) +} + +/// One decoded page from the pool: `id > from`, ASC, up to `limit` +/// (the extent guard is the caller's — this helper is contract-blind +/// to it). The stream decode owner: shared by the auto-commit reads +/// (this module) and the tx reads (`tx.rs`) — the `stream` field +/// mapping and the `Codec` posture guard live once. +pub(crate) async fn stream_read_page( + pool: Pool, + schema: String, + stream: String, + from: i64, + limit: i64, +) -> Result> { + let conn = pool.get().await.map_err(pool_error)?; + let events = table(&schema, tables::EVENTS); + let rows = conn + .query( + &format!( + "SELECT id, \"key\", payload, created_at + FROM {events} + WHERE stream = $1 AND id > $2 + ORDER BY id ASC + LIMIT $3" + ), + &[&stream, &from, &limit], + ) + .await + .map_err(pg_error)?; + stream_events_from_rows(&rows, stream) +} + +/// Decode a read page into the contract's [`StreamEvent`] values (the +/// core constructor's `from_row`). Malformed rows map to `Error::Codec` +/// (the decode-side posture guard, ADR-021 §2). One owner: the tx +/// module imports this; no forked decoders. +pub(crate) fn stream_events_from_rows( + rows: &[tokio_postgres::Row], + stream: String, +) -> Result> { + let mut out = Vec::with_capacity(rows.len()); + for row in rows { + let key: Option = row + .try_get(1) + .map_err(|e| Error::Codec(format!("stream row key must be text: {e}")))?; + let payload: Option> = row + .try_get(2) + .map_err(|e| Error::Codec(format!("stream row payload must be bytea: {e}")))?; + out.push(StreamEvent::from_row( + row.get(0), + stream.clone(), + key, + payload.unwrap_or_default(), + row.get(3), + )); + } + Ok(out) +} + +/// The checkpoint read helper (absent consumer = 0, ADR-019 §1). +async fn stream_get_offset(pool: &Pool, schema: &str, stream: &str, consumer: &str) -> Result { + let conn = pool.get().await.map_err(pool_error)?; + let offsets = table(schema, tables::OFFSETS); + let row = conn + .query_opt( + &format!("SELECT \"offset\" FROM {offsets} WHERE stream = $1 AND consumer = $2"), + &[&stream, &consumer], + ) + .await + .map_err(pg_error)?; + Ok(row.map(|r| r.get(0)).unwrap_or(0)) +} + +/// The monotone checkpoint save (one owner of the monotonicity rule — +/// the same GREATEST-clamped, WHERE-guarded upsert `save_offset_tx` +/// rides; a save at-or-below the stored checkpoint updates nothing, +/// ADR-019 §6). First save seeds the checkpoint, clamped at 0 (the +/// offsets table's non-negative domain CHECK would abort a negative +/// seed; the clamp preserves both the domain and the monotone rule — +/// a negative save at-or-below the stored checkpoint is the same +/// silent no-op a 0-save is). +async fn stream_save_offset( + pool: &Pool, + schema: &str, + stream: &str, + consumer: &str, + offset: i64, +) -> Result<()> { + let conn = pool.get().await.map_err(pool_error)?; + let offsets = table(schema, tables::OFFSETS); + conn.execute( + &format!( + "INSERT INTO {offsets} AS o (stream, consumer, \"offset\") + VALUES ($1, $2, GREATEST($3::bigint, 0)) + ON CONFLICT (stream, consumer) DO UPDATE + SET \"offset\" = EXCLUDED.\"offset\" + WHERE EXCLUDED.\"offset\" > o.\"offset\"" + ), + &[&stream, &consumer, &offset], + ) + .await + .map_err(pg_error)?; + Ok(()) +} + +impl StreamHandle for PgStreamHandle { + fn name(&self) -> &str { + &self.name + } + + fn publish<'a>(&'a self, payload: Value) -> BoxedFuture<'a, Result> { + publish_with_key( + self.pool.clone(), + self.schema.clone(), + self.closed.clone(), + self.name.clone(), + None, + payload, + ) + } + + fn publish_with_key<'a>( + &'a self, + key: Option, + payload: Value, + ) -> BoxedFuture<'a, Result> { + publish_with_key( + self.pool.clone(), + self.schema.clone(), + self.closed.clone(), + self.name.clone(), + key, + payload, + ) + } + + fn read_since<'a>( + &'a self, + offset: i64, + limit: i64, + ) -> BoxedFuture<'a, Result>> { + if self.closed.load(Ordering::Acquire) { + return Box::pin(async { Err(database_error("the store is closed")) }); + } + // The extent guard (ADR-023 §2): `limit <= 0` reads nothing — + // a value, never an error — at trait-impl entry, before any + // round trip (pg's LIMIT with a negative is a server error, + // unreachable behind this guard). + if limit <= 0 { + return Box::pin(async { Ok(Vec::new()) }); + } + Box::pin(stream_read_page( + self.pool.clone(), + self.schema.clone(), + self.name.clone(), + offset, + limit, + )) + } + + fn read_from_consumer<'a>( + &'a self, + consumer: &str, + limit: i64, + ) -> BoxedFuture<'a, Result>> { + if let Err(e) = validate_local_name(consumer) { + return Box::pin(async move { Err(e) }); + } + if self.closed.load(Ordering::Acquire) { + return Box::pin(async { Err(database_error("the store is closed")) }); + } + if limit <= 0 { + return Box::pin(async { Ok(Vec::new()) }); + } + let pool = self.pool.clone(); + let schema = self.schema.clone(); + let stream = self.name.clone(); + let consumer = consumer.to_string(); + Box::pin(async move { + let from = stream_get_offset(&pool, &schema, &stream, &consumer).await?; + stream_read_page(pool, schema, stream, from, limit).await + }) + } + + fn save_offset<'a>(&'a self, consumer: &str, offset: i64) -> BoxedFuture<'a, Result<()>> { + if let Err(e) = validate_local_name(consumer) { + return Box::pin(async move { Err(e) }); + } + if self.closed.load(Ordering::Acquire) { + return Box::pin(async { Err(database_error("the store is closed")) }); + } + let pool = self.pool.clone(); + let schema = self.schema.clone(); + let stream = self.name.clone(); + let consumer = consumer.to_string(); + Box::pin( + async move { stream_save_offset(&pool, &schema, &stream, &consumer, offset).await }, + ) + } + + fn get_offset<'a>(&'a self, consumer: &str) -> BoxedFuture<'a, Result> { + if let Err(e) = validate_local_name(consumer) { + return Box::pin(async move { Err(e) }); + } + if self.closed.load(Ordering::Acquire) { + return Box::pin(async { Err(database_error("the store is closed")) }); + } + let pool = self.pool.clone(); + let schema = self.schema.clone(); + let stream = self.name.clone(); + let consumer = consumer.to_string(); + Box::pin(async move { stream_get_offset(&pool, &schema, &stream, &consumer).await }) + } + + fn trim_to<'a>(&'a self, horizon: i64) -> BoxedFuture<'a, Result> { + if self.closed.load(Ordering::Acquire) { + return Box::pin(async { Err(database_error("the store is closed")) }); + } + let pool = self.pool.clone(); + let schema = self.schema.clone(); + let stream = self.name.clone(); + Box::pin(async move { + // Wakes nothing (ADR-015 §5's posture — trim is a + // bounded-growth op, not an event). The boundary argument + // is total (ADR-023 §2): a negative horizon deletes + // nothing idempotently (plain SQL, no guard — `id <= -7` + // matches no rows). + let conn = pool.get().await.map_err(pool_error)?; + let events = table(&schema, tables::EVENTS); + let deleted = conn + .execute( + &format!("DELETE FROM {events} WHERE stream = $1 AND id <= $2"), + &[&stream, &horizon], + ) + .await + .map_err(pg_error)?; + Ok(deleted as i64) + }) + } + + fn subscribe<'a>(&'a self, consumer: &str) -> BoxedFuture<'a, Result>> { + if let Err(e) = validate_local_name(consumer) { + return Box::pin(async move { Err(e) }); + } + if self.closed.load(Ordering::Acquire) { + return Box::pin(async { Err(database_error("the store is closed")) }); + } + let name = self.name.clone(); + let pool = self.pool.clone(); + let schema = self.schema.clone(); + let forwarder = self.forwarder.clone(); + let closed = self.closed.clone(); + let consumer = consumer.to_string(); + Box::pin(async move { spawn_bridge(name, consumer, pool, schema, forwarder, closed).await }) + } +} + +/// The `subscribe` bridge construction: the forwarder subscription is +/// taken before any read (a commit cannot fall between the attach +/// read and the wake wait), the bridge task then replays from the +/// stored offset and parks on the wake feed. Subscribing a failed +/// forwarder (`None` — shutdown raced) fails closed after +/// unregistering the just-registered LISTEN (no registration leak — +/// the notify path's same rule); a bridge-spawn failure (runtime +/// gone) also fails `Database`. +async fn spawn_bridge( + stream: String, + consumer: String, + pool: Pool, + schema: String, + forwarder: Arc, + store_closed: Arc, +) -> Result> { + // Registration first — a failed registration here (never: LISTEN + // failures carry the registry entry's recovery, but a shutdown + // race can kill the command channel) fails the whole subscribe. + forwarder.register(&stream).await?; + let Some(broadcast_rx) = forwarder.subscribe() else { + // The subscribe raced (or arrived past) a shutdown: fail + // closed — no receiver can ever deliver again. Unregister + // first: the acked LISTEN must not leak past the failed + // listen (a next-closest `subscribe` re-issues it). + forwarder.unregister(&stream); + return Err(database_error("the store is closed")); + }; + let (event_tx, event_rx) = + tokio::sync::mpsc::channel::>(EVENT_CHANNEL_CAPACITY); + let spawned = tokio::spawn(run_subscribe_loop( + stream.clone(), + consumer.clone(), + pool.clone(), + schema.clone(), + store_closed.clone(), + broadcast_rx, + event_tx, + )); + // tokio::spawn only fails on a shut-down runtime — and the boxed + // receiver must never surface as a lie after its bridge never + // started. (Also: no panics cross the seam, the W-2 posture.) + if spawned.is_finished() { + spawned.abort(); + return Err(database_error( + "failed to spawn the stream bridge task (runtime shutting down)", + )); + } + Ok(Box::new(PgEventReceiver { + stream, + consumer, + rx: event_rx, + forwarder, + position: 0, + pool, + schema, + store_closed, + _bridge: spawned, + })) +} + +/// The bridge loop body (the subscribe task; the SQLite arm's +/// std-thread shape, natively async here): +/// +/// 1. read the consumer's stored offset (the attach read), +/// 2. drain pages until the tail (`id ASC`), +/// 3. park on the wake feed (this stream's channel + the reserved +/// reconnect-wake), re-drain on wake, +/// 4. a read failure while the store is open surfaces one +/// `Err(Database)` and waits for the next wake; a read failure +/// after store close exits the loop (the consumer sees the close). +/// +/// Wake-feed disconnect (the forwarder's shutdown take drops the +/// broadcast) is the driver arm: the loop exits, the mpsc sender +/// drops, the consumer's `recv()` yields `None` — terminal, only at +/// engine shutdown. +enum WakeTick { + Ok, + Closed, +} + +async fn wait_wake(rx: &mut tokio::sync::broadcast::Receiver) -> WakeTick { + match rx.recv().await { + Ok(_raw) => WakeTick::Ok, + Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => { + // Still surfaced — treat like a wake: re-drain. + WakeTick::Ok + } + Err(tokio::sync::broadcast::error::RecvError::Closed) => WakeTick::Closed, + } +} + +/// One drain-page arm: a decoded page (possibly empty), or the loop's +/// exit state. +enum PageArm { + Page(Vec), + Err(Error), + Closed, +} + +async fn drain_page( + pool: &Pool, + schema: &str, + stream: &str, + cursor: i64, + store_closed: &AtomicBool, +) -> PageArm { + if store_closed.load(Ordering::Acquire) { + return PageArm::Closed; + } + match stream_read_page( + pool.clone(), + schema.to_string(), + stream.to_string(), + cursor, + SUBSCRIBE_PAGE, + ) + .await + { + Ok(events) => PageArm::Page(events), + Err(_e) if store_closed.load(Ordering::Acquire) => PageArm::Closed, + Err(e) => PageArm::Err(database_only(e)), + } +} + +async fn run_subscribe_loop( + stream: String, + consumer: String, + pool: Pool, + schema: String, + store_closed: Arc, + mut wake_rx: tokio::sync::broadcast::Receiver, + event_tx: tokio::sync::mpsc::Sender>, +) { + // 1. The attach read: the consumer's stored offset. + if store_closed.load(Ordering::Acquire) { + return; + } + let mut cursor = match stream_get_offset(&pool, &schema, &stream, &consumer).await { + Ok(offset) => offset, + Err(_e) if store_closed.load(Ordering::Acquire) => return, + Err(e) => { + let _ = event_tx.send(Err(database_only(e))).await; + return; + } + }; + loop { + // 2. Drain to the tail, page by page. + loop { + match drain_page(&pool, &schema, &stream, cursor, &store_closed).await { + PageArm::Closed => return, + PageArm::Err(e) => { + if event_tx.send(Err(e)).await.is_err() { + return; + } + break; + } + PageArm::Page(events) => { + if events.is_empty() { + break; + } + for event in events { + cursor = event.offset; + if event_tx.send(Ok(event)).await.is_err() { + return; + } + } + } + } + } + // 3. Park on the wake feed (this stream's channel — and the + // reserved reconnect-wake: the same re-drain heals a + // gap's missed wakes by reading the durable rows). + match wait_wake(&mut wake_rx).await { + WakeTick::Closed => return, + WakeTick::Ok => {} + } + } +} + +/// `recv()`'s `Err` arm carries `Database` only (ADR-021 §5): a +/// read error that mapped elsewhere (the decode-side `Codec` posture +/// guard) is remapped into the opaque fallback with the original error +/// preserved as the source. +fn database_only(e: Error) -> Error { + match e { + already @ Error::Database(_) => already, + other => Error::database(other), + } +} + +/// The durable subscription receiver: the consumer side of the +/// bridge. Dropping it drops the mpsc receiver — the bridge's next +/// `send` fails and the bridge task exits (and its `_bridge` handle +/// drop-detaches the task body), while the forwarder registration +/// rides the receiver's `Drop` (unregister → UNLISTEN at the +/// last-subscriber drop — one refcount owner per channel). +pub(crate) struct PgEventReceiver { + stream: String, + consumer: String, + rx: tokio::sync::mpsc::Receiver>, + forwarder: Arc, + position: i64, + pool: Pool, + schema: String, + store_closed: Arc, + _bridge: tokio::task::JoinHandle<()>, +} + +impl Drop for PgEventReceiver { + fn drop(&mut self) { + self.forwarder.unregister(&self.stream); + } +} + +impl PgEventReceiver { + /// The monotone save of this receiver's position through the + /// shared op. The sync signature on the natively-async engine is + /// bridged on the dedicated idle wake runtime (see + /// [`drive_sync`]; the runtime carries no engine state — the + /// save future checks out its own pool connection). `position == + /// 0` (nothing yielded yet) is the no-op arm — the stored + /// checkpoint untouched (ADR-021 §5). + fn save_now(&mut self) -> Result<()> { + let position = self.position; + if position == 0 { + return Ok(()); + } + if self.store_closed.load(Ordering::Acquire) { + return Err(database_error("the store is closed")); + } + let pool = self.pool.clone(); + let schema = self.schema.clone(); + let stream = self.stream.clone(); + let consumer = self.consumer.clone(); + drive_sync(async move { + stream_save_offset(&pool, &schema, &stream, &consumer, position).await + })? + } +} + +impl EventReceiver for PgEventReceiver { + fn recv<'a>(&'a mut self) -> BoxedFuture<'a, Option>> { + Box::pin(async move { + match self.rx.recv().await { + Some(Ok(event)) => { + self.position = event.offset; + Some(Ok(event)) + } + // A transient read failure — the bridge keeps running + // and re-reads on the next wake. + Some(Err(e)) => Some(Err(e)), + // The bridge exited: shutdown disconnect — `None` = + // closed, terminal (it never reopens; recovery is a + // fresh `subscribe(consumer)`). + None => None, + } + }) + } + + fn try_recv(&mut self) -> Result> { + match self.rx.try_recv() { + Ok(Ok(event)) => { + self.position = event.offset; + Ok(Some(event)) + } + Ok(Err(e)) => Err(e), + Err(tokio::sync::mpsc::error::TryRecvError::Empty) => Ok(None), + Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => Err(Error::Closed), + } + } + + fn read_since<'a>( + &'a mut self, + offset: i64, + limit: i64, + ) -> BoxedFuture<'a, Result>> { + if self.store_closed.load(Ordering::Acquire) { + return Box::pin(async { Err(database_error("the store is closed")) }); + } + if limit <= 0 { + return Box::pin(async { Ok(Vec::new()) }); + } + Box::pin(stream_read_page( + self.pool.clone(), + self.schema.clone(), + self.stream.clone(), + offset, + limit, + )) + } + + fn save_offset(&mut self) -> Result<()> { + self.save_now() + } + + fn offset(&self) -> i64 { + self.position + } +} diff --git a/alkstore-postgres/src/tx.rs b/alkstore-postgres/src/tx.rs index 7a1d6d8..72c7cd7 100644 --- a/alkstore-postgres/src/tx.rs +++ b/alkstore-postgres/src/tx.rs @@ -79,6 +79,7 @@ use crate::resolution::{ }; use crate::schema::{QualifiedTable, tables}; use crate::seam::{pg_error, pool_error}; +use crate::stream::stream_events_from_rows; /// The notification payload boundary (the contract's pinned limit — /// client-side checked, typed before any round trip; the predicate @@ -637,27 +638,7 @@ fn job_from_row(row: &tokio_postgres::Row) -> Result { )) } -/// Decode a read page into the contract's [`StreamEvent`] values (the -/// core constructor's `from_row`). Malformed rows map to `Error::Codec`. -fn stream_events_from_rows( - rows: &[tokio_postgres::Row], - stream: String, -) -> Result> { - let mut out = Vec::with_capacity(rows.len()); - for row in rows { - let key: Option = row - .try_get(1) - .map_err(|e| Error::Codec(format!("stream row key must be text: {e}")))?; - let payload: Option> = row - .try_get(2) - .map_err(|e| Error::Codec(format!("stream row payload must be bytea: {e}")))?; - out.push(StreamEvent::from_row( - row.get(0), - stream.clone(), - key, - payload.unwrap_or_default(), - row.get(3), - )); - } - Ok(out) -} +// The stream-page decode (`stream_events_from_rows`) lives in +// `stream.rs` — the one decode owner the tx reads and the auto-commit +// reads share (the `stream` field mapping and the `Codec` posture +// guard live once; ADR-012 §2's one-owner rule). diff --git a/tasks/pg-engine-streams.md b/tasks/pg-engine-streams.md index 967a170..939cccd 100644 --- a/tasks/pg-engine-streams.md +++ b/tasks/pg-engine-streams.md @@ -1,7 +1,7 @@ --- id: pg-engine-streams name: Postgres engine — streams (`StreamHandle`, reads, subscribe, trim) -status: pending +status: completed depends_on: [pg-engine-seam-tx, pg-engine-notify-listen] scope: moderate risk: medium @@ -83,27 +83,27 @@ engine-postgres.md's streams mapping row): ## Acceptance Criteria -- [ ] `publish`/`publish_with_key` return the assigned bigserial +- [x] `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 +- [x] 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); +- [x] `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 +- [x] `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, +- [x] `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 → +- [x] 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 +- [x] Validation + closed-store fail-closed on all entry points +- [x] `cargo test -p alkstore-postgres` (harness server), clippy `-D warnings`, fmt clean; gates green server-less ## References @@ -119,8 +119,126 @@ engine-postgres.md's streams mapping row): ## Notes -> To be filled by implementation agent +> 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 -> To be filled on completion \ No newline at end of file +> 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. \ No newline at end of file