SQLite engine: notify/listen — auto-commit notify, watcher-fanout WakeReceiver bridge (task sqlite-engine-notify-listen)

This commit is contained in:
glm-5.3-flash committed 2026-10-08 12:12:34 +00:00
1 parent c38033db1a
commit 8efe0c2e78
7 files changed
+795 -24

No files matched your search

+1 -1
View File
@@ -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
+1
View File
@@ -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;
+149
View File
@@ -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<Writer>,
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<SharedUpdateWatcher>,
closed: Arc<AtomicBool>,
channel: &str,
) -> Result<Box<dyn WakeReceiver>> {
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<Wake>,
watcher: Arc<SharedUpdateWatcher>,
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<Wake>> {
Box::pin(async move { self.wake_rx.recv().await })
}
fn try_recv(&mut self) -> Result<Option<Wake>> {
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<Option<Wake>>> {
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),
}
})
}
}
+35
View File
@@ -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<String>) -> 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<T, F>(writer: Arc<Writer>, f: F) -> alkstore::Result<T>
where
T: Send + 'static,
F: FnOnce(&rusqlite::Connection) -> alkstore::Result<T> + 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
}
+14 -13
View File
@@ -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<Readers>,
watcher: Arc<SharedUpdateWatcher>,
db_path: PathBuf,
closed: Arc<AtomicBool>,
}
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<SqliteStore> {
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<dyn alkstore::WakeReceiver>>> {
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)]
+512
View File
@@ -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<F: Fn() -> 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);
}