diff --git a/alkstore-contract-suite/src/lib.rs b/alkstore-contract-suite/src/lib.rs index 4029558..689ee0e 100644 --- a/alkstore-contract-suite/src/lib.rs +++ b/alkstore-contract-suite/src/lib.rs @@ -46,5 +46,6 @@ pub use properties::{ payload_too_large_produced_on_pg, publish_with_key_tx_commit_atomicity, queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms, scheduler_boundary_fires, scheduler_bounded_catchup, scheduler_leadership_discipline, stream_ordering_equivalence, - sweep_no_stranded_rows, trim_to_semantics, with_tx_panicking_closure_rolls_back, + sweep_no_stranded_rows, trim_to_semantics, wake_receiver_shapes, + with_tx_panicking_closure_rolls_back, }; diff --git a/alkstore-contract-suite/src/properties.rs b/alkstore-contract-suite/src/properties.rs index 2a651b6..d2b8a55 100644 --- a/alkstore-contract-suite/src/properties.rs +++ b/alkstore-contract-suite/src/properties.rs @@ -1117,6 +1117,244 @@ pub async fn receiver_close_and_save_arms(factory: &dyn StoreFactory) { ); } +/// **Wake semantics under the pinned `WakeReceiver` shapes** — the +/// wake contract's pinnable core, tolerance-bounded (wakes are +/// best-effort hints — the row asserts state outcomes, never delivery +/// counts or latencies). A pre-attached listener receives a wake after +/// a committed `notify` on its channel, with all three recv forms +/// exercised (`recv`, `try_recv`, `recv_timeout`) and every wake +/// observed carrying the listened channel — the one piece of semantic +/// content a wake carries (ADR-008 §3). A burst of notifies yields +/// *at least one* wake, never an exact per-notify count (coalescing is +/// the documented engine asymmetry — possibly coalesced on SQLite's +/// 1-slot feed, per-notify on pg; the row pins the floor, not the +/// shape). A listener attached *after* a commit never sees that +/// commit's notify (`recv_timeout` idles through a bounded absence +/// window — no replay; the notify commits before the attach, and pg's +/// no-replay hole — a commit during a forwarder gap is never +/// re-delivered — and SQLite's burst coalescing are both legal under +/// this pin). Wakes from a different channel never arrive on this +/// listener: after a foreign-channel notify, the bounded window +/// surfaces no wake carrying the foreign consumer channel (SQLite's +/// documented overtrigger may deliver a wake on the same commit, but +/// it carries the *listened* channel only — nothing about what was +/// notified; the pg fanout skips non-LISTENed channels, and the +/// reserved reconnect-wake carries its own reserved name — that +/// reserved straggler is a tolerated engine arm, not a violation), and +/// a same-channel positive control still delivers — the silence was +/// channel scoping, never a dead listener. +/// +/// Cross-reference: the failure-surface close arms — SQLite's +/// watcher-death `recv() -> None` and the pg synthetic reconnect-wake +/// — are the engine-differing close arms, already pinned engine-side +/// and in `receiver_close_and_save_arms` (the stream receiver's +/// disposal leg); this row does not re-pin them — the reserved +/// reconnect-wake appears here only as a tolerated straggler. +/// +/// Contract stamp: ADR-006 (the wake contract — opaque wake + +/// re-read, best-effort hints not per-commit delivery, coalescing / +/// overtrigger asymmetry, the channel-scoped fanout); ADR-008 §3 (the +/// pinned `Wake { channel }` / recv forms — the wake's channel field +/// matches the listened channel; `recv_timeout`'s expiry-without-wake +/// is the documented `Ok(None)` idle arm). +pub async fn wake_receiver_shapes(factory: &dyn StoreFactory) { + let store = factory.open().await.expect("factory opens a store"); + + // Delivery + channel identity: a pre-attached listener, a + // committed notify on its channel, every recv form once. + let mut rx = store.listen("wake-a").await.expect("listen attaches"); + + store + .notify("wake-a", json!({"n": 1})) + .await + .expect("notify commits"); + let wake = tokio::time::timeout(std::time::Duration::from_secs(10), rx.recv()) + .await + .expect("the committed notify delivers a wake — the bounded wait is enough") + .expect("the receiver is open through a healthy listen"); + assert_eq!( + wake.channel, "wake-a", + "the recv wake's channel field matches the listened channel" + ); + + store + .notify("wake-a", json!({"n": 2})) + .await + .expect("notify commits"); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10); + let wake = loop { + match rx.try_recv() { + Ok(Some(w)) => break w, + Ok(None) => {} + Err(e) => panic!("the receiver is open through a healthy listen: {e}"), + } + assert!( + std::time::Instant::now() < deadline, + "the committed notify delivered no wake within the bounded wait" + ); + tokio::time::sleep(std::time::Duration::from_millis(25)).await; + }; + assert_eq!( + wake.channel, "wake-a", + "the try_recv wake's channel field matches the listened channel" + ); + + store + .notify("wake-a", json!({"n": 3})) + .await + .expect("notify commits"); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10); + let wake = loop { + match rx + .recv_timeout(std::time::Duration::from_millis(50)) + .await + .expect("recv_timeout's Err is closed/storage-failure only, idle is Ok(None)") + { + Some(w) => break w, + None => { + assert!( + std::time::Instant::now() < deadline, + "the committed notify delivered no wake within the bounded wait" + ); + } + } + }; + assert_eq!( + wake.channel, "wake-a", + "the recv_timeout wake's channel field matches the listened channel" + ); + + // The burst floor: three committed notifies, at least one wake. + // The count is never pinned — coalescing is the documented engine + // asymmetry (SQLite's 1-slot feed; pg per-notify), the contract + // leaves the shape engine-owned. + for n in 0..3 { + store + .notify("wake-a", json!({"burst": n})) + .await + .expect("notify commits"); + } + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10); + let mut heard = 0; + while heard == 0 { + match rx + .recv_timeout(std::time::Duration::from_millis(50)) + .await + .expect("the receiver is open through a healthy listen") + { + Some(w) => { + assert_eq!( + w.channel, "wake-a", + "the burst wake's channel field matches the listened channel" + ); + heard += 1; + } + None => { + assert!( + std::time::Instant::now() < deadline, + "a burst of notifies yields at least one wake — never zero \ + within the bounded wait" + ); + } + } + } + + // No replay: a listener attached after a commit never sees that + // commit's notify. The pre-settle sleep closes SQLite's baseline + // race — the watcher's tick past the settle has consumed the + // commit's version change (a past-window spurious wake arriving + // before the late attach is the documented overtrigger; the pin + // starts at the attach). + store + .notify("wake-b", json!({"late": true})) + .await + .expect("notify commits"); + tokio::time::sleep(std::time::Duration::from_millis(300)).await; + let mut late = store + .listen("wake-b") + .await + .expect("the late listen attaches"); + let heard = late + .recv_timeout(std::time::Duration::from_millis(400)) + .await + .expect("the absence window is Ok(None) or a wake — never a hang"); + match heard { + None => {} + Some(w) if w.channel == alkstore::RESERVED_LISTENER_RECONNECTED => {} + Some(w) => panic!( + "a listener attached after a commit never sees that commit's \ + notify — got a wake carrying {}", + w.channel + ), + } + + // Channel-scoped delivery: a foreign-channel notify never wakes + // this listener. The bounded window drains no wake carrying the + // foreign channel; what may arrive is the listened channel (SQLite's + // legal overtrigger) or the reserved reconnect straggler (pg). + let mut scoped = store.listen("scope-a").await.expect("listen attaches"); + store + .notify("scope-b", json!({"foreign": true})) + .await + .expect("notify commits"); + let deadline = std::time::Instant::now() + std::time::Duration::from_millis(400); + loop { + let wake = scoped + .recv_timeout(std::time::Duration::from_millis(100)) + .await + .expect("the scoped listener is open"); + match wake { + None => break, + Some(w) + if w.channel != "scope-a" + && w.channel != alkstore::RESERVED_LISTENER_RECONNECTED => + { + panic!( + "a wake from a different channel never arrives on this \ + listener: {}", + w.channel + ); + } + Some(_) => {} + } + assert!( + std::time::Instant::now() < deadline, + "the scoped silence window overflowed with unnamed wakes" + ); + } + + // Positive control: a same-channel notify still delivers — the + // idle above was channel scoping, not a dead listener. + store + .notify("scope-a", json!({"own": true})) + .await + .expect("notify commits"); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10); + let wake = loop { + match scoped + .recv_timeout(std::time::Duration::from_millis(50)) + .await + .expect("the scoped listener is open") + { + Some(w) => break w, + None => { + assert!( + std::time::Instant::now() < deadline, + "the same-channel positive control delivered no wake — \ + the listener is dead, not idle, and the scoping pin \ + above is vacuous" + ); + } + } + }; + assert_eq!( + wake.channel, "scope-a", + "the control wake's channel field matches the listened channel" + ); + + factory.teardown().await.expect("factory teardown"); +} + /// **`PayloadTooLarge` produced pg-side** — the Postgres engine /// rejects an oversized `notify`/`notify_tx` payload client-side with /// the typed `PayloadTooLarge` carrying the limit, **before any round diff --git a/alkstore-postgres/tests/contract_suite.rs b/alkstore-postgres/tests/contract_suite.rs index 90316be..e188f4a 100644 --- a/alkstore-postgres/tests/contract_suite.rs +++ b/alkstore-postgres/tests/contract_suite.rs @@ -172,7 +172,7 @@ mod rows { queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms, scheduler_boundary_fires, scheduler_bounded_catchup, scheduler_leadership_discipline, stream_ordering_equivalence, sweep_no_stranded_rows, trim_to_semantics, - with_tx_panicking_closure_rolls_back, + wake_receiver_shapes, with_tx_panicking_closure_rolls_back, }; use super::PgFactory; @@ -294,6 +294,11 @@ mod rows { "row-stream-ordering" ); harness_row!(row_trim_to_semantics, trim_to_semantics, "row-trim-to"); + harness_row!( + row_wake_receiver_shapes, + wake_receiver_shapes, + "row-wake-shapes" + ); harness_row!( row_with_tx_panicking_closure_rolls_back, with_tx_panicking_closure_rolls_back, diff --git a/alkstore-sqlite/tests/contract_suite.rs b/alkstore-sqlite/tests/contract_suite.rs index 311d1eb..6f74d16 100644 --- a/alkstore-sqlite/tests/contract_suite.rs +++ b/alkstore-sqlite/tests/contract_suite.rs @@ -76,7 +76,7 @@ mod rows { queue_depth_reclaim_and_dead_letter, receiver_close_and_save_arms, scheduler_boundary_fires, scheduler_bounded_catchup, scheduler_leadership_discipline, stream_ordering_equivalence, sweep_no_stranded_rows, trim_to_semantics, - with_tx_panicking_closure_rolls_back, + wake_receiver_shapes, with_tx_panicking_closure_rolls_back, }; use super::SqliteFactory; @@ -190,6 +190,11 @@ mod rows { trim_to_semantics(&factory("row-trim-to")).await; } + #[tokio::test(flavor = "multi_thread")] + async fn row_wake_receiver_shapes() { + wake_receiver_shapes(&factory("row-wake-shapes")).await; + } + #[tokio::test(flavor = "multi_thread")] async fn row_with_tx_panicking_closure_rolls_back() { with_tx_panicking_closure_rolls_back(&factory("row-panic-rollback")).await; diff --git a/tasks/suite-wake-rows.md b/tasks/suite-wake-rows.md index a6696f0..034ce42 100644 --- a/tasks/suite-wake-rows.md +++ b/tasks/suite-wake-rows.md @@ -1,7 +1,7 @@ --- id: suite-wake-rows name: Contract-suite row — wake semantics under the pinned WakeReceiver shapes -status: pending +status: completed depends_on: [] scope: narrow risk: low @@ -66,8 +66,55 @@ doc comment). ## Notes -> To be filled by implementation agent +> Dispositions of record: + +- **The no-replay leg's pre-settle sleep is load-bearing on SQLite**: + the substrate watcher's `last_version` baseline is captured at store + open, so a commit landing between the watcher's last poll and the + late `listen()` fires one late tick at the new subscriber — + indistinguishable from a replay. The 300 ms settle (bounded far above + the 1 ms default poll cadence) closes the race: the pin starts at + the attach, and a past-window spurious wake arriving before it is + the documented overtrigger. +- **The channel-scoped leg's observable is the wake's `channel` + field, not absence**: on SQLite the watcher fires on *any* commit, + so a foreign-channel notify (may) deliver a wake — but its field + carries only the *listened* channel (cloned at listen). The row + asserts: no wake naming a foreign consumer channel ever surfaces + (`scope-a` and the reserved reconnect name are the two admitted + carriers), and a same-channel positive control still delivers — the + silence is scoping, not a dead listener. On pg the same window is + structurally empty (LISTEN is channel-scoped) — the vacuous column + is itself the scoping proof. +- **The reserved reconnect-wake is admitted as a tolerated straggler** + in both absence windows (a reconnect coincident with a window is a + legal engine arm) — the arms themselves stay pinned engine-side and + in `receiver_close_and_save_arms`, cross-referenced in the row's doc + comment as the task requires. +- **Absence windows are single `recv_timeout` calls** (400 ms / 400 ms) + — the documented `Ok(None)` idle arm doubles as the bounded wait; + no timing-value assertions, single-task drive, no two-task race + windows. +- **The core-contract backlog row stays in the inventory** — same + posture as the other discharged rows (locks, streams, depth): the + backlog list flushes at the review-wave-5 gate that flips the engine + specs to stable. ## Summary -> To be filled on completion \ No newline at end of file +> Landed the `wake_receiver_shapes` contract-suite row +> (`alkstore-contract-suite/src/properties.rs`, version-stamped +> ADR-006 + ADR-008 §3, doc comment cross-references the close arms it +> does not re-pin) and wired it into both engines' suite targets +> (SQLite tokio test, pg `harness_row!`) — 23 → 24 rows per column. +> Legs: pre-attached listener + committed notify through all three +> recv forms with channel-field match; three-notify burst flooring at +> one wake (never an exact count); attach-after-commit never replays +> (`recv_timeout` idles, settled); foreign-channel notify never +> surfaces a foreign wake + same-channel positive control. Verified: +> SQLite column green server-less; pg column green against the harness +> server (postgres/poc@:15432); 3 solo re-runs of the row per engine +> (determinism); `cargo test -p alkstore-sqlite -p alkstore-postgres` +> green (SQLite 189+24, pg 121+24+9); workspace `cargo test` green; +> clippy `-D warnings` clean; fmt clean (fmt reflowed the pg +> `harness_row!` after the first pass). \ No newline at end of file