diff --git a/alkstore-sqlite/Cargo.toml b/alkstore-sqlite/Cargo.toml index 1ec64f7..09dcb19 100644 --- a/alkstore-sqlite/Cargo.toml +++ b/alkstore-sqlite/Cargo.toml @@ -11,7 +11,7 @@ parking_lot = "0.12" rusqlite = { version = "0.40", features = ["bundled", "functions"] } serde_json = "1" thiserror = "2" -tokio = { version = "1", features = ["rt", "sync"] } +tokio = { version = "1", features = ["rt", "sync", "time"] } # file-id only supports unix and windows. Other targets (WASI, Redox, # illumos, etc.) get the `(0, 0)` fallback in substrate::watcher's diff --git a/alkstore-sqlite/src/lib.rs b/alkstore-sqlite/src/lib.rs index 6b048ce..78f6091 100644 --- a/alkstore-sqlite/src/lib.rs +++ b/alkstore-sqlite/src/lib.rs @@ -17,6 +17,7 @@ //! `spawn_blocking` seam — rusqlite connections are not //! `Send`-across-await; the family REQ-TTY-01 posture). +mod notify; mod opts; mod resolution; mod seam; diff --git a/alkstore-sqlite/src/notify.rs b/alkstore-sqlite/src/notify.rs new file mode 100644 index 0000000..3329b6e --- /dev/null +++ b/alkstore-sqlite/src/notify.rs @@ -0,0 +1,149 @@ +//! The notify/listen mechanism's SQLite arm — the wake half +//! (`sqlite-engine-notify-listen`; `notify_tx` landed with the seam): +//! +//! - `notify` — the auto-commit path: acquire the writer slot, run +//! the substrate's `notify()` SQL function (one statement, one +//! auto-commit transaction), release. Commit-atomic by SQLite's +//! autocommit; rollback is impossible on this path (nothing holds +//! the statement back). No payload limit on SQLite — this engine +//! never produces [`Error::PayloadTooLarge`](alkstore::Error) (ADR-016 +//! §5); the payload crosses the notify table as the serde_json +//! serialization's UTF-8 text, the byte form ADR-020 §4 pins. +//! - `listen` — the watcher-fanout bridge: one `spawn_blocking` thread +//! per subscription (`engine-sqlite.md`'s mapping row) looping +//! `recv()` on the substrate's sync feed (1-slot — bursts coalesce) +//! and `blocking_send`-ing [`Wake`] values into a tokio channel the +//! [`WakeReceiver`] owns. Wakes are coalesced hints (overtrigger on +//! purpose, ADR-006): the watcher fires on *any* commit, so a wake +//! may arrive for unrelated writes, may coalesce, or may repeat — +//! the receiver delivers `Wake { channel }` only, never payloads or +//! ids (ADR-008 §3); consumers are idempotent on wake. +//! +//! Close semantics (ADR-006, ADR-008 §3): watcher death +//! (`WatcherDeathGuard`) clears the subscriber feed, the bridge thread +//! exits, and the receiver sees the close — `recv() -> None`, terminal +//! (it never reopens); `try_recv` / `recv_timeout` carry +//! [`Error::Closed`], distinguishable from idle's `Ok(None)`; +//! `recv_timeout`'s timeout expiry without a wake is `Ok(None)` — the +//! documented idle arm, never a timeout signal. +//! +//! Unsubscribe is the receiver's own [`Drop`]: the substrate prunes +//! disconnected subscribers lazily (at the next fanout), so the +//! receiver unsubscribes eagerly on drop — a dropped subscriber leaves +//! the feed immediately and the bridge thread exits on the disconnect. + +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::time::Duration; + +use serde_json::Value; + +use alkstore::{BoxedFuture, Error, Result, Wake, WakeReceiver}; + +use crate::seam::{sqlite_error, with_writer}; +use crate::substrate::{SharedUpdateWatcher, Writer}; + +/// The bridge's tokio-channel capacity. The substrate's 1-slot sync +/// feed stays the coalescing point; this buffer only smooths the +/// bridge-thread → consumer handoff. +const WAKE_CHANNEL_CAPACITY: usize = 16; + +/// The auto-commit notify path (see the module docs). Name validation +/// fires at the entry point, before any round trip and before the +/// payload is touched (ADR-008 §4). +pub(crate) fn notify( + writer: Arc, + channel: &str, + payload: Value, +) -> BoxedFuture<'static, Result<()>> { + if let Err(e) = alkstore::validate_shared_name(channel) { + return Box::pin(async move { Err(e) }); + } + let channel = channel.to_string(); + Box::pin(async move { + with_writer(writer, move |conn| { + let bytes = alkstore::encode_payload(&payload)?; + let text = String::from_utf8(bytes) + .map_err(|e| Error::Codec(format!("notify payload must be utf-8 text: {e}")))?; + conn.query_row( + "SELECT notify(?1, ?2)", + rusqlite::params![channel, text], + |_| Ok(()), + ) + .map_err(sqlite_error)?; + Ok(()) + }) + .await + }) +} + +/// The listen path (see the module docs). Name validation fires at the +/// entry point; a listen against a closed store fails closed with +/// [`Error::Database`](alkstore::Error) (the engine-wide closed-store +/// posture — `begin_tx`/`notify` yield the same shape). +pub(crate) fn listen( + watcher: Arc, + closed: Arc, + channel: &str, +) -> Result> { + alkstore::validate_shared_name(channel)?; + if closed.load(Ordering::Acquire) { + return Err(Error::database(std::io::Error::other( + "the store is closed", + ))); + } + let (id, rx) = watcher.subscribe(); + let (wake_tx, wake_rx) = tokio::sync::mpsc::channel(WAKE_CHANNEL_CAPACITY); + let channel = channel.to_string(); + let _bridge = tokio::task::spawn_blocking(move || { + while let Ok(()) = rx.recv() { + if wake_tx.blocking_send(Wake::new(channel.clone())).is_err() { + break; + } + } + }); + Ok(Box::new(SqliteWakeReceiver { + wake_rx, + watcher, + id, + })) +} + +/// The consumer-visible wake receiver: the tokio side of the bridge. +/// Its [`Drop`] unsubscribes the substrate feed (the subscriber is +/// pruned eagerly — before the watcher's next fanout would). +struct SqliteWakeReceiver { + wake_rx: tokio::sync::mpsc::Receiver, + watcher: Arc, + id: u64, +} + +impl Drop for SqliteWakeReceiver { + fn drop(&mut self) { + self.watcher.unsubscribe(self.id); + } +} + +impl WakeReceiver for SqliteWakeReceiver { + fn recv<'a>(&'a mut self) -> BoxedFuture<'a, Option> { + Box::pin(async move { self.wake_rx.recv().await }) + } + + fn try_recv(&mut self) -> Result> { + match self.wake_rx.try_recv() { + Ok(wake) => Ok(Some(wake)), + Err(tokio::sync::mpsc::error::TryRecvError::Empty) => Ok(None), + Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => Err(Error::Closed), + } + } + + fn recv_timeout<'a>(&'a mut self, timeout: Duration) -> BoxedFuture<'a, Result>> { + Box::pin(async move { + match tokio::time::timeout(timeout, self.wake_rx.recv()).await { + Ok(Some(wake)) => Ok(Some(wake)), + Ok(None) => Err(Error::Closed), + Err(_expired) => Ok(None), + } + }) + } +} diff --git a/alkstore-sqlite/src/seam.rs b/alkstore-sqlite/src/seam.rs index 42b51d2..17b3db7 100644 --- a/alkstore-sqlite/src/seam.rs +++ b/alkstore-sqlite/src/seam.rs @@ -12,6 +12,9 @@ //! mapping, family-standard no-panics discipline maintained. use alkstore::Error; +use std::sync::Arc; + +use crate::substrate::Writer; /// The blocking bridge every trait method's substrate round trip goes /// through. Private: the seam is engine-internal — consumers meet it @@ -40,3 +43,35 @@ pub(crate) fn database_error(message: impl Into) -> Error { pub(crate) fn sqlite_error(e: rusqlite::Error) -> Error { Error::database(e) } + +/// The short-lived writer-slot lease for auto-commit ops (the +/// `notify` path's shape): acquire, run the op inside the +/// `spawn_blocking` seam, release — the slot is free the instant the +/// op completes (no lease held across a consumer's `await` points; +/// that is the long-transaction posture, not this path's). A closed +/// store's acquire fails closed with `Database` (the engine-wide +/// closed-store shape). +pub(crate) async fn with_writer(writer: Arc, f: F) -> alkstore::Result +where + T: Send + 'static, + F: FnOnce(&rusqlite::Connection) -> alkstore::Result + Send + 'static, +{ + blocking(move || { + let conn = writer.acquire().ok_or_else(|| { + Error::database(std::io::Error::other( + "the store is closed: writer slot unavailable", + )) + })?; + match f(&conn) { + Ok(out) => { + writer.release(conn); + Ok(out) + } + Err(e) => { + drop(conn); + Err(e) + } + } + }) + .await +} diff --git a/alkstore-sqlite/src/store.rs b/alkstore-sqlite/src/store.rs index 8277ffc..76aae58 100644 --- a/alkstore-sqlite/src/store.rs +++ b/alkstore-sqlite/src/store.rs @@ -34,6 +34,7 @@ use alkstore::Store; use std::path::PathBuf; use std::sync::Arc; +use std::sync::atomic::AtomicBool; use crate::opts::SqliteOpts; use crate::seam::{database_error, sqlite_error}; @@ -48,6 +49,7 @@ pub struct SqliteStore { readers: Arc, watcher: Arc, db_path: PathBuf, + closed: Arc, } impl std::fmt::Debug for SqliteStore { @@ -64,6 +66,7 @@ impl SqliteStore { /// reader pool, and closes the writer slot. Idempotent. [`Drop`] /// delegates here. pub fn close(&self) { + self.closed.store(true, std::sync::atomic::Ordering::SeqCst); let _ = self.watcher.close(); self.readers.close(); self.writer.close(); @@ -105,6 +108,7 @@ fn open_store(path: &str, opts: SqliteOpts) -> alkstore::Result { readers: Arc::new(Readers::new(path.to_string(), opts.max_readers)), watcher: Arc::new(watcher), db_path: PathBuf::from(path), + closed: Arc::new(AtomicBool::new(false)), }) } @@ -122,25 +126,20 @@ impl Store for SqliteStore { fn notify<'a>( &'a self, - _channel: &str, - _payload: serde_json::Value, + channel: &str, + payload: serde_json::Value, ) -> alkstore::BoxedFuture<'a, alkstore::Result<()>> { - Box::pin(async { - Err(database_error( - "notify wiring lands with the mechanism tasks", - )) - }) + crate::notify::notify(self.writer.clone(), channel, payload) } fn listen<'a>( &'a self, - _channel: &str, + channel: &str, ) -> alkstore::BoxedFuture<'a, alkstore::Result>> { - Box::pin(async { - Err(database_error( - "listen wiring lands with the mechanism tasks", - )) - }) + let watcher = self.watcher.clone(); + let closed = self.closed.clone(); + let channel = channel.to_string(); + Box::pin(async move { crate::notify::listen(watcher, closed, &channel) }) } fn stream<'a>( @@ -221,6 +220,8 @@ impl Store for SqliteStore { } } +#[cfg(test)] +mod notify_tests; #[cfg(test)] mod open_tests; #[cfg(test)] diff --git a/alkstore-sqlite/src/store/notify_tests.rs b/alkstore-sqlite/src/store/notify_tests.rs new file mode 100644 index 0000000..7d2a5cf --- /dev/null +++ b/alkstore-sqlite/src/store/notify_tests.rs @@ -0,0 +1,512 @@ +//! The notify/listen mechanism's acceptance tests +//! (`sqlite-engine-notify-listen`): commit-atomic auto-commit +//! notifies (payload crosses the notify table, no size limit), +//! watcher-fanout wakes through the bridged receiver, the close arms +//! (watcher death ⇒ `None` terminal; `Closed` vs idle in +//! `try_recv`/`recv_timeout`), unsubscribe-on-drop, entry-point +//! validation, and the `PayloadTooLarge` never-produced pin (ADR-016 +//! §5). + +use std::path::PathBuf; +use std::time::Duration; + +use alkstore::{Error, Store, WakeReceiver}; +use serde_json::json; + +use crate::store::{open, open_store}; + +fn temp_dir(tag: &str) -> PathBuf { + let d = std::env::temp_dir().join(format!( + "alkstore-notify-{tag}-{}-{:?}", + std::process::id(), + std::thread::current().id() + )); + let _ = std::fs::remove_dir_all(&d); + std::fs::create_dir_all(&d).unwrap(); + d +} + +fn cleanup(dir: &PathBuf) { + let _ = std::fs::remove_dir_all(dir); +} + +fn temp_path(tag: &str) -> PathBuf { + temp_dir(tag).join("store.db") +} + +/// Wait for a wake with a bounded deadline — no tight-race assertions +/// (the suite's stability posture). +async fn must_recv(receiver: &mut dyn WakeReceiver, tag: &str) -> alkstore::Wake { + let deadline = tokio::time::Instant::now() + Duration::from_secs(5); + loop { + match receiver.try_recv() { + Ok(Some(wake)) => return wake, + Ok(None) => {} + Err(e) => panic!("{tag}: receiver errored instead of idling: {e}"), + } + assert!( + tokio::time::Instant::now() < deadline, + "{tag}: wake never arrived" + ); + tokio::time::sleep(Duration::from_millis(10)).await; + } +} + +/// Wait until the predicate holds with a bounded deadline. +async fn until bool>(tag: &str, f: F) { + let deadline = tokio::time::Instant::now() + Duration::from_secs(5); + while !f() { + assert!( + tokio::time::Instant::now() < deadline, + "{tag}: condition never held" + ); + tokio::time::sleep(Duration::from_millis(10)).await; + } +} + +/// Acceptance: `notify` delivers commit-atomic auto-commit notifies — +/// the payload crosses the notify table (a committed row is readable +/// through a raw reader), the notify is a single auto-commit +/// transaction (no partial state possible), and the channel rides the +/// row. +#[tokio::test(flavor = "multi_thread")] +async fn notify_delivers_commit_atomic_through_the_notify_table() { + let dir = temp_dir("auto-commit"); + let path = temp_path("auto-commit"); + let store = open_store(path.to_str().unwrap(), Default::default()).unwrap(); + + store + .notify("orders", json!({"x": 1})) + .await + .expect("auto-commit notify succeeds"); + + let conn = + rusqlite::Connection::open_with_flags(&path, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY) + .unwrap(); + let (channel, payload): (String, String) = conn + .query_row( + "SELECT channel, payload FROM __alkstore_notifications + ORDER BY id DESC LIMIT 1", + [], + |r| Ok((r.get(0)?, r.get(1)?)), + ) + .unwrap(); + assert_eq!(channel, "orders", "the channel is the row's channel"); + assert_eq!( + payload, + json!({"x": 1}).to_string(), + "the payload is the serde_json serialization (ADR-020 §4's byte form)" + ); + store.close(); + cleanup(&dir); +} + +/// Acceptance: no payload limit on SQLite — this engine never produces +/// `PayloadTooLarge` at any size (ADR-016 §5's occurrence-asymmetry +/// pin). A multi-megabyte payload notifies successfully. +#[tokio::test(flavor = "multi_thread")] +async fn notify_never_produces_payload_too_large() { + let dir = temp_dir("no-limit"); + let store = open_store(temp_path("no-limit").to_str().unwrap(), Default::default()).unwrap(); + + let big = "x".repeat(2_000_000); + store + .notify("big", json!({"blob": big})) + .await + .expect("SQLite notify has no payload limit (ADR-016 §5)"); + store.close(); + cleanup(&dir); +} + +/// Acceptance: `listen` returns a working receiver — wakes arrive +/// after commits by other connections, and `Wake { channel }` carries +/// the channel name only. The honest posture: the watcher fans out on +/// any commit, so the wake is pinned to *arrival*, never exclusivity +/// (best-effort hints, ADR-006). +#[tokio::test(flavor = "multi_thread")] +async fn listen_wakes_arrive_after_commits() { + let dir = temp_dir("wake-arrives"); + let path = temp_path("wake-arrives"); + let store = open_store(path.to_str().unwrap(), Default::default()).unwrap(); + + let mut receiver = store.listen("orders").await.unwrap(); + assert_eq!( + receiver.try_recv().unwrap(), + None, + "a fresh receiver idles, it is not closed" + ); + + // A commit from another connection (a fresh store handle over the + // same db — the cross-connection probe). + let path2 = path.clone(); + let other = tokio::spawn(async move { + let s = open(path2.to_str().unwrap(), Default::default()).unwrap(); + let mut tx = s.begin_tx().await.unwrap(); + tx.enqueue_tx("q", Default::default(), json!({})) + .await + .unwrap(); + tx.commit().await.unwrap(); + }); + other.await.unwrap(); + + let wake = must_recv(&mut *receiver, "wake-arrives").await; + assert_eq!( + wake.channel, "orders", + "the wake carries the listened channel" + ); + store.close(); + cleanup(&dir); +} + +/// Acceptance: wakes carry the channel only — never payloads or ids +/// (ADR-008 §3; a notify's payload crosses the notify table, the wake +/// that follows it does not surface the payload). +#[tokio::test(flavor = "multi_thread")] +async fn wakes_carry_the_channel_only() { + let dir = temp_dir("channel-only"); + let store = open_store( + temp_path("channel-only").to_str().unwrap(), + Default::default(), + ) + .unwrap(); + + let mut receiver = store.listen("c").await.unwrap(); + store + .notify("c", json!({"secret": [1, 2, 3]})) + .await + .unwrap(); + + let wake = must_recv(&mut *receiver, "channel-only").await; + assert_eq!(wake.channel, "c"); + store.close(); + cleanup(&dir); +} + +/// Acceptance: watcher death closes the receiver terminally — +/// `recv() -> None`. Driven through the substrate's death-signal path: +/// store close (the ordinary death — the death-guard clears every +/// subscriber feed) closes the bridged receiver; it never reopens. +#[tokio::test(flavor = "multi_thread")] +async fn watcher_death_closes_the_receiver_terminally() { + let dir = temp_dir("death-close"); + let store = open_store( + temp_path("death-close").to_str().unwrap(), + Default::default(), + ) + .unwrap(); + + let mut receiver = store.listen("orders").await.unwrap(); + + // Force the substrate's death signal: close the store (the + // death-guard clears every subscriber feed — the bridge thread + // exits, the receiver sees the disconnect). + store.close(); + + let none = receiver.recv().await; + assert!( + none.is_none(), + "watcher death must close the receiver: got {none:?}" + ); + // Terminal: try_recv reports Closed, not idle (ADR-008 §3's + // distinguishable states). + assert!( + matches!(receiver.try_recv(), Err(Error::Closed)), + "after close, try_recv is Err(Closed) — not idle" + ); + // Terminal means forever: a second recv is still None. + assert!( + receiver.recv().await.is_none(), + "a closed receiver never reopens" + ); + cleanup(&dir); +} + +/// Acceptance: `try_recv` distinguishes closed from idle pre-close — +/// `Ok(None)` with no wake pending (idle), `Some` when one arrived, +/// `Err(Closed)` only after the source closed (the pinned arms). +#[tokio::test(flavor = "multi_thread")] +async fn try_recv_distinguishes_idle_from_closed() { + let dir = temp_dir("try-recv-arms"); + let store = open_store( + temp_path("try-recv-arms").to_str().unwrap(), + Default::default(), + ) + .unwrap(); + + let mut receiver = store.listen("c").await.unwrap(); + assert_eq!(receiver.try_recv().unwrap(), None, "idle is Ok(None)"); + + store.notify("c", json!({"n": 1})).await.unwrap(); + let wake = must_recv(&mut *receiver, "try-recv-arms").await; + assert_eq!(wake.channel, "c"); + // Drained again — idle, not closed. + assert_eq!( + receiver.try_recv().unwrap(), + None, + "drained is idle, not closed" + ); + + store.close(); + assert!( + matches!(receiver.try_recv(), Err(Error::Closed)), + "post-close try_recv is Err(Closed)" + ); + cleanup(&dir); +} + +/// Acceptance: `recv_timeout`'s arms match the pinned shapes — timeout +/// expiry without a wake is `Ok(None)` (the documented idle arm, never +/// a timeout signal); a wake in the window arrives as `Ok(Some)`; the +/// source's death is `Err(Closed)`. +#[tokio::test(flavor = "multi_thread")] +async fn recv_timeout_matches_the_pinned_arms() { + let dir = temp_dir("recv-timeout"); + let store = open_store( + temp_path("recv-timeout").to_str().unwrap(), + Default::default(), + ) + .unwrap(); + + let mut receiver = store.listen("c").await.unwrap(); + + // Expiry without a wake: Ok(None) — never Err (ADR-008 §3's pinned + // arm: the Err arm carries Database or Closed, never a timeout). + let idle = receiver + .recv_timeout(Duration::from_millis(100)) + .await + .unwrap(); + assert_eq!(idle, None, "timeout expiry is the documented Ok(None) arm"); + + // A wake inside the window arrives as Ok(Some). + store.notify("c", json!({"n": 1})).await.unwrap(); + let deadline = tokio::time::Instant::now() + Duration::from_secs(5); + let wake = loop { + match receiver + .recv_timeout(Duration::from_millis(100)) + .await + .unwrap() + { + Some(w) => break w, + None => { + assert!( + tokio::time::Instant::now() < deadline, + "the wake never arrived within recv_timeout windows" + ); + } + } + }; + assert_eq!(wake.channel, "c"); + + // Death: Err(Closed). + store.close(); + let err = receiver.recv_timeout(Duration::from_millis(100)).await; + assert!( + matches!(err, Err(Error::Closed)), + "post-close recv_timeout is Err(Closed), got {err:?}" + ); + cleanup(&dir); +} + +/// Acceptance: receiver drop unsubscribes — no leaked subscriptions; +/// `subscriber_count` returns to baseline (the substrate prunes +/// dropped subscribers; the eager unsubscribe makes it immediate). +#[tokio::test(flavor = "multi_thread")] +async fn receiver_drop_unsubscribes_without_leaks() { + let dir = temp_dir("unsub"); + let store = open_store(temp_path("unsub").to_str().unwrap(), Default::default()).unwrap(); + + let baseline = store.watcher.subscriber_count(); + assert_eq!(baseline, 0, "no subscriptions before any listen"); + + { + let receiver = store.listen("c").await.unwrap(); + drop(receiver); + } + + // The receiver's Drop unsubscribed eagerly; the bridge thread may + // still be exiting, but the substrate-side subscriber is gone. + until("unsub", || store.watcher.subscriber_count() == baseline).await; + assert_eq!( + store.watcher.subscriber_count(), + baseline, + "receiver drop must not leak subscriptions" + ); + store.close(); + cleanup(&dir); +} + +/// Acceptance: dropping the store closes receivers through the same +/// path (Drop delegates to close) — the death arms hold there too. +#[tokio::test(flavor = "multi_thread")] +async fn dropping_the_store_closes_receivers() { + let dir = temp_dir("drop-close"); + let mut receiver; + { + let store = open_store( + temp_path("drop-close").to_str().unwrap(), + Default::default(), + ) + .unwrap(); + receiver = store.listen("c").await.unwrap(); + drop(store); + } + assert!(receiver.recv().await.is_none(), "drop closes subscribers"); + assert!(matches!(receiver.try_recv(), Err(Error::Closed))); + cleanup(&dir); +} + +/// Acceptance: channel validation on both entry points — empty / +/// whitespace-only rejected with `InvalidName`, reserved-prefix names +/// rejected with `ReservedName`, before any round trip (ADR-008 §4; +/// the validation failures leave no notify row, hold no subscription). +#[tokio::test(flavor = "multi_thread")] +async fn entry_points_validate_channels() { + let dir = temp_dir("validation"); + let path = temp_path("validation"); + let store = open_store(path.to_str().unwrap(), Default::default()).unwrap(); + + for channel in [ + "", + " ", + "__alkstore_listener_reconnected__", + "__alkstore_x", + ] { + let err = match store.notify(channel, json!({})).await { + Err(e) => e, + Ok(()) => panic!("notify({channel:?}) must reject by validation"), + }; + match err { + Error::InvalidName { name } => assert_eq!(channel, name), + Error::ReservedName { name } => assert_eq!(channel, name), + other => panic!("notify({channel:?}) must reject by validation, got {other}"), + } + let err = match store.listen(channel).await { + Err(e) => e, + Ok(_) => panic!("listen({channel:?}) must reject by validation"), + }; + match err { + Error::InvalidName { name } => assert_eq!(channel, name), + Error::ReservedName { name } => assert_eq!(channel, name), + other => panic!("listen({channel:?}) must reject by validation, got {other}"), + } + } + + // The rejects were total: no notify row was written and no + // subscription was held. + let conn = + rusqlite::Connection::open_with_flags(&path, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY) + .unwrap(); + let rows: i64 = conn + .query_row("SELECT COUNT(*) FROM __alkstore_notifications", [], |r| { + r.get(0) + }) + .unwrap(); + assert_eq!(rows, 0, "rejected notifies never touched the table"); + assert_eq!( + store.watcher.subscriber_count(), + 0, + "rejected listens never subscribed" + ); + store.close(); + cleanup(&dir); +} + +/// Acceptance: multiple listeners on one channel all fan out +/// independently; an unrelated channel's listen also wakes (the +/// honest overtrigger posture — pinned as hint-arrival, not +/// exclusivity, per the task's note). +#[tokio::test(flavor = "multi_thread")] +async fn multiple_listeners_fan_out() { + let dir = temp_dir("fanout"); + let store = open_store(temp_path("fanout").to_str().unwrap(), Default::default()).unwrap(); + + let mut a = store.listen("a").await.unwrap(); + let mut b = store.listen("a").await.unwrap(); + let mut other_ch = store.listen("zzz").await.unwrap(); + + store.notify("a", json!({})).await.unwrap(); + + assert_eq!(must_recv(&mut *a, "fanout-a").await.channel, "a"); + assert_eq!(must_recv(&mut *b, "fanout-b").await.channel, "a"); + // Overtrigger on purpose: any commit wakes every subscriber; the + // other channel's receiver carries its own channel name in the + // wake (hints are coalesced, channel-correct, not exclusive). + let wake = must_recv(&mut *other_ch, "fanout-other").await; + assert_eq!(wake.channel, "zzz"); + store.close(); + cleanup(&dir); +} + +/// Acceptance: notifies on a closed store fail closed with +/// `Error::Database` (the engine-wide closed-store posture — +/// `begin_tx` yields the same shape); listens on a closed store do +/// too. +#[tokio::test(flavor = "multi_thread")] +async fn closed_store_fails_notify_and_listen() { + let dir = temp_dir("closed-store"); + let store = open_store( + temp_path("closed-store").to_str().unwrap(), + Default::default(), + ) + .unwrap(); + store.close(); + + let err = store.notify("c", json!({})).await.unwrap_err(); + assert!( + matches!(err, Error::Database(_)), + "notify on a closed store fails closed, got {err}" + ); + let err = match store.listen("c").await { + Err(e) => e, + Ok(_) => panic!("listen on a closed store must fail closed"), + }; + assert!( + matches!(err, Error::Database(_)), + "listen on a closed store fails closed, got {err}" + ); + cleanup(&dir); +} + +/// Acceptance: `notify_tx` wakes listeners at commit — the tx half +/// rides the same wake plumbing (in-tx notifies are invisible until +/// commit; the commit wakes the listener; rollback wakes nothing). +#[tokio::test(flavor = "multi_thread")] +async fn tx_notify_wakes_at_commit_and_rollback_wakes_nothing() { + let dir = temp_dir("tx-wake"); + let store = open_store(temp_path("tx-wake").to_str().unwrap(), Default::default()).unwrap(); + let mut receiver = store.listen("orders").await.unwrap(); + + let mut tx = store.begin_tx().await.unwrap(); + tx.notify_tx("orders", json!({"x": 1})).await.unwrap(); + // In-tx: the wake is not yet deliverable — nothing committed. + tokio::time::sleep(Duration::from_millis(200)).await; + assert_eq!( + receiver.try_recv().unwrap(), + None, + "an uncommitted notify must not wake" + ); + tx.commit().await.unwrap(); + let wake = must_recv(&mut *receiver, "tx-wake").await; + assert_eq!(wake.channel, "orders"); + + // Rollback: the notify row is dropped with the tx — no wake fires + // (a rollback does not bump data_version; the substrate pins + // commit-only wakes). + let mut tx = store.begin_tx().await.unwrap(); + tx.notify_tx("orders", json!({"x": 2})).await.unwrap(); + drop(tx); + tokio::time::sleep(Duration::from_millis(200)).await; + // The watcher may legitimately emit a delayed wake for the *first* + // commit if it hadn't drained; drain what's there and require only + // that no wake arrives *after* a quiet settle (coalescing is the + // contract's shape, not a defect). + while receiver.try_recv().unwrap().is_some() {} + tokio::time::sleep(Duration::from_millis(200)).await; + assert_eq!( + receiver.try_recv().unwrap(), + None, + "a rolled-back notify must not wake" + ); + store.close(); + cleanup(&dir); +} diff --git a/tasks/sqlite-engine-notify-listen.md b/tasks/sqlite-engine-notify-listen.md index fbdb5ac..279e24f 100644 --- a/tasks/sqlite-engine-notify-listen.md +++ b/tasks/sqlite-engine-notify-listen.md @@ -1,7 +1,7 @@ --- id: sqlite-engine-notify-listen name: SQLite engine — notify/listen (watcher fanout → `WakeReceiver` bridge) -status: pending +status: completed depends_on: [sqlite-engine-seam-tx] scope: narrow risk: medium @@ -43,19 +43,19 @@ should pin wake-arrives and close-on-death, not wake-exclusivity. ## Acceptance Criteria -- [ ] `notify` delivers commit-atomic auto-commit notifies; payload +- [x] `notify` delivers commit-atomic auto-commit notifies; payload crosses the notify table (no size limit — never `PayloadTooLarge`) -- [ ] `listen` returns a working `WakeReceiver`: wakes arrive after +- [x] `listen` returns a working `WakeReceiver`: wakes arrive after commits by other connections; `Wake { channel }` carries the channel name only -- [ ] Watcher death closes the receiver (`recv() -> None`, terminal) +- [x] Watcher death closes the receiver (`recv() -> None`, terminal) — tested via the substrate's death-signal path -- [ ] `try_recv`/`recv_timeout` arms match the pinned contract shapes -- [ ] Receiver drop unsubscribes (no leaked subscriptions — +- [x] `try_recv`/`recv_timeout` arms match the pinned contract shapes +- [x] Receiver drop unsubscribes (no leaked subscriptions — `subscriber_count` returns to baseline) -- [ ] Channel validation on both entry points -- [ ] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean +- [x] Channel validation on both entry points +- [x] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean ## References @@ -67,8 +67,81 @@ should pin wake-arrives and close-on-death, not wake-exclusivity. ## Notes -> To be filled by implementation agent +- **Bridge architecture**: the substrate's sync feed (1-slot + `sync_channel` — the coalescing point) → one dedicated + `std::thread`-via-`spawn_blocking` per subscription looping + `recv()` and `blocking_send`ing `Wake { channel }` into a tokio + mpsc channel (capacity 16 — smoothing only; the sync feed stays the + coalescer) that `SqliteWakeReceiver` owns. The bridge thread exits + on either side's disconnect: subscriber feed closed (watcher death, + `WatcherDeathGuard`) or consumer-side `WakeReceiver` dropped. +- **Unsubscribe eagerly on drop, not lazily**: the substrate prunes + disconnected subscribers lazily (at the next fanout), which could + strand the bridge thread if a consumer drops the receiver on an + idle db (no future fanout ever comes → the sender never sees the + disconnect). `SqliteWakeReceiver`'s `Drop` calls + `watcher.unsubscribe(id)` — immediate teardown, tested via + `subscriber_count` returning to baseline. The substrate doc on + `subscribe` ("callers MUST unsubscribe") is honored at the bridge + layer. +- **`closed` flag on the store**: `SqliteStore` gained an + `Arc closed` set by `close()`; `listen` checks it and + fails closed with `Database` (matching `begin_tx`'s writer-slot + shape). Needed because `SharedUpdateWatcher::subscribe` on a closed + watcher would still succeed mechanically (it re-inserts into the + senders map) while never receiving — a silent-dead subscription. + `notify` on a closed store fails closed through + `Writer::acquire() -> None` (no new machinery). +- **`with_writer` seam helper (new, `seam.rs`)**: the short-lived + writer-slot lease for auto-commit ops — acquire, op in + `spawn_blocking`, release; the slot is free the instant the op + completes (no lease held across consumer `await` points — that is + the long-tx posture, not this path). Reusable by later auto-commit + mechanism tasks (locks/streams/queues constructors' read paths can + route through it or the reader pool as their tasks decide). +- **Death-signal test driver**: the death path is pinned via store + `close()` (the ordinary death — the death-guard clears every + subscriber feed). The file-replacement death (the fork's + dead-man's-switch) is covered at the substrate layer already + (`shared_update_watcher_signals_subscribers_on_watcher_death`); + the bridge only sees "feed closed" either way. +- **Rollback-wakes-nothing**: pinned indirectly — a rollback does not + bump `data_version` (substrate-pinned), so no wake follows the + `notify_tx`-then-drop path; the test drains first (coalescing makes + an exact pre-drain assertion racy, per the wake-hint posture) and + requires quiet-after-settle instead. +- **`tokio` `time` feature added** to the engine crate's main deps + (`recv_timeout`'s `tokio::time::timeout`) — previously dev-only. +- **Store validation ordering**: `listen` validates the channel before + the closed check (validation errors are the contract-shaped answer + even on a closed store — the rejects never touch the engine); + `notify` validates before encoding the payload. ## Summary -> To be filled on completion \ No newline at end of file +Implemented the wake half of the notify mechanism on the SQLite +engine: `Store::notify` (auto-commit — writer-slot acquire, the +substrate's `notify()` SQL function in one auto-commit transaction, +release; payload crosses `__alkstore_notifications` as the serde_json +serialization's UTF-8 text; no size limit — ADR-016 §5's +never-`PayloadTooLarge` pinned by a 2 MB payload test) and +`Store::listen` (the watcher-fanout bridge: substrate +`SharedUpdateWatcher::subscribe` per subscription, one +`spawn_blocking` thread doing `blocking_send` of +`Wake { channel: … }` into a tokio channel owned by the boxed +`SqliteWakeReceiver`; eager unsubscribe on drop). Close semantics +pinned: watcher death / store close ⇒ `recv() -> None` terminal +(never reopens); `try_recv` distinguishes `Err(Closed)` from idle +`Ok(None)`; `recv_timeout` expiry without a wake is the documented +`Ok(None)` arm, death is `Err(Closed)`. Entry-point validation +(`InvalidName`/`ReservedName`) on both paths, before any round trip +or payload encoding. New seam helper `with_writer` (short-lived +writer-slot lease for auto-commit ops). 13 new tests (auto-commit row +crossing, no-limit, wake-arrives cross-connection, channel-only +content, death-close terminal + never-reopens, try_recv arms, +recv_timeout arms, unsubscribe no-leak, drop-closes, validation +×4 shapes both entry points ×no side effects, multi-listener fanout, +closed-store fail-closed, tx-notify wakes-at-commit). Verified: +`cargo test --workspace` (25 + 3 + 146 sqlite, all green), suite run +3× green, `cargo clippy --all-targets -- -D warnings` clean, +`cargo fmt --check` clean. \ No newline at end of file