From 2a2ad7e11f8fab450d43e2edebc1b640d5cf5079 Mon Sep 17 00:00:00 2001 From: "glm-5.3-flash" Date: Sat, 10 Oct 2026 07:12:31 +0000 Subject: [PATCH] =?UTF-8?q?Contract-suite=20wake=20row=20(task=20suite-wak?= =?UTF-8?q?e-rows):=20wake=5Freceiver=5Fshapes=20=E2=80=94=20the=20wake=20?= =?UTF-8?q?contract's=20pinnable=20core,=20tolerance-bounded=20to=20state?= =?UTF-8?q?=20outcomes=20(never=20delivery=20counts=20or=20latencies):=20a?= =?UTF-8?q?=20pre-attached=20listener=20receives=20a=20wake=20after=20a=20?= =?UTF-8?q?committed=20notify=20on=20its=20channel=20through=20all=20three?= =?UTF-8?q?=20recv=20forms=20(recv/try=5Frecv/recv=5Ftimeout),=20every=20w?= =?UTF-8?q?ake's=20channel=20field=20matching=20the=20listened=20channel?= =?UTF-8?q?=20(the=20one=20piece=20of=20semantic=20content=20a=20wake=20ca?= =?UTF-8?q?rries,=20ADR-008=20=C2=A73);=20a=20three-notify=20burst=20floor?= =?UTF-8?q?s=20at=20one=20wake=20=E2=80=94=20never=20an=20exact=20per-noti?= =?UTF-8?q?fy=20count=20(coalescing=20the=20documented=20engine=20asymmetr?= =?UTF-8?q?y:=20SQLite's=201-slot=20feed=20vs=20pg=20per-notify,=20the=20r?= =?UTF-8?q?ow=20pins=20the=20floor=20not=20the=20shape);=20a=20listener=20?= =?UTF-8?q?attached=20after=20a=20commit=20never=20sees=20that=20commit's?= =?UTF-8?q?=20notify=20=E2=80=94=20recv=5Ftimeout=20idles=20a=20400=20ms?= =?UTF-8?q?=20absence=20window=20(a=20300=20ms=20pre-settle=20sleep=20clos?= =?UTF-8?q?es=20SQLite's=20watcher=20baseline=20race:=20the=20last=5Fversi?= =?UTF-8?q?on=20baseline=20is=20store-open-captured,=20so=20an=20unconsume?= =?UTF-8?q?d=20commit's=20version=20change=20would=20fire=20one=20late=20t?= =?UTF-8?q?ick=20at=20the=20new=20subscriber=20indistinguishable=20from=20?= =?UTF-8?q?replay;=20pg's=20gap-commit=20no-replay=20hole=20and=20SQLite's?= =?UTF-8?q?=20burst=20coalescing=20both=20legal=20under=20the=20pin);=20an?= =?UTF-8?q?d=20the=20channel-scoped=20leg=20=E2=80=94=20a=20foreign-channe?= =?UTF-8?q?l=20notify=20never=20surfaces=20a=20wake=20naming=20it=20(SQLit?= =?UTF-8?q?e's=20same-commit=20overtrigger=20may=20deliver=20but=20carries?= =?UTF-8?q?=20only=20the=20listened=20channel;=20the=20pg=20LISTEN=20fanou?= =?UTF-8?q?t=20skips=20foreign=20channels;=20the=20reserved=20reconnect=20?= =?UTF-8?q?straggler=20tolerated),=20with=20a=20same-channel=20positive=20?= =?UTF-8?q?control=20proving=20the=20silence=20is=20scoping=20not=20a=20de?= =?UTF-8?q?ad=20listener.=20The=20failure-surface=20close=20arms=20(SQLite?= =?UTF-8?q?=20watcher-death=20recv()->None,=20pg=20synthetic=20reconnect-w?= =?UTF-8?q?ake)=20stay=20pinned=20engine-side=20and=20in=20receiver=5Fclos?= =?UTF-8?q?e=5Fand=5Fsave=5Farms's=20disposal=20leg=20=E2=80=94=20cross-re?= =?UTF-8?q?ferenced=20in=20the=20doc=20comment,=20not=20re-pinned.=20Stamp?= =?UTF-8?q?ed=20ADR-006=20+=20ADR-008=20=C2=A73.=20Wired=20into=20both=20e?= =?UTF-8?q?ngines'=20suite=20targets=20(SQLite=20tokio=20test,=20pg=20harn?= =?UTF-8?q?ess=5Frow!),=2024=20rows=20per=20column=20up=20from=2023.=20Dis?= =?UTF-8?q?positions=20in=20Notes:=20the=20backlog=20row=20stays=20in=20co?= =?UTF-8?q?re-contract.md's=20inventory=20(flushes=20at=20review-wave-5's?= =?UTF-8?q?=20stable=20gate,=20like=20the=20other=20discharged=20rows),=20?= =?UTF-8?q?the=20absence=20windows=20are=20single=20recv=5Ftimeout=20calls?= =?UTF-8?q?=20(the=20documented=20Ok(None)=20idle=20arm=20as=20the=20bound?= =?UTF-8?q?ed=20wait),=20and=20the=20scoping=20leg's=20observable=20is=20t?= =?UTF-8?q?he=20channel=20field=20not=20absence=20(SQLite=20overtrigger=20?= =?UTF-8?q?makes=20an=20absence-only=20pin=20vacuous=20there).=20Verified:?= =?UTF-8?q?=20SQLite=20column=20green=20server-less,=20pg=20column=20green?= =?UTF-8?q?=20vs=20harness=20(postgres/poc@:15432),=203=20solo=20re-runs?= =?UTF-8?q?=20of=20the=20row=20per=20engine=20(determinism),=20cargo=20tes?= =?UTF-8?q?t=20-p=20alkstore-sqlite=20-p=20alkstore-postgres=20green=20(SQ?= =?UTF-8?q?Lite=20189+24,=20pg=20121+24+9),=20workspace=20cargo=20test=20g?= =?UTF-8?q?reen=20(12=20binaries),=20clippy=20-D=20warnings,=20fmt=20clean?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- alkstore-contract-suite/src/lib.rs | 3 +- alkstore-contract-suite/src/properties.rs | 238 ++++++++++++++++++++++ alkstore-postgres/tests/contract_suite.rs | 7 +- alkstore-sqlite/tests/contract_suite.rs | 7 +- tasks/suite-wake-rows.md | 53 ++++- 5 files changed, 302 insertions(+), 6 deletions(-) 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