SQLite engine: streams — StreamHandle, extent-guarded reads, trim, durable subscribe (task sqlite-engine-streams)

This commit is contained in:
glm-5.3-flash committed 2026-10-08 12:30:32 +00:00
1 parent 8efe0c2e78
commit 4913302b47
6 files changed
+1507 -48

No files matched your search

+1
View File
@@ -22,6 +22,7 @@ mod opts;
mod resolution;
mod seam;
mod store;
mod stream;
mod substrate;
mod tx;
+18 -5
View File
@@ -144,12 +144,23 @@ impl Store for SqliteStore {
fn stream<'a>(
&'a self,
_name: &str,
name: &str,
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::StreamHandle>>> {
Box::pin(async {
Err(database_error(
"stream wiring lands with the mechanism tasks",
))
if let Err(e) = alkstore::validate_shared_name(name) {
return Box::pin(async move { Err(e) });
}
if self.closed.load(std::sync::atomic::Ordering::Acquire) {
return Box::pin(async { Err(crate::seam::database_error("the store is closed")) });
}
let name = name.to_string();
let writer = self.writer.clone();
let readers = self.readers.clone();
let watcher = self.watcher.clone();
let closed = self.closed.clone();
Box::pin(async move {
Ok(Box::new(crate::stream::SqliteStreamHandle::new(
name, writer, readers, watcher, closed,
)) as Box<dyn alkstore::StreamHandle>)
})
}
@@ -225,4 +236,6 @@ mod notify_tests;
#[cfg(test)]
mod open_tests;
#[cfg(test)]
mod stream_tests;
#[cfg(test)]
mod tx_tests;
+741
View File
@@ -0,0 +1,741 @@
//! The stream mechanism's acceptance tests (`sqlite-engine-streams`):
//! publish/read/save/get/trim/subscribe all work against a temp-file
//! store, the extent guard, exact-boundary trim with surviving offsets
//! unrenumbered and negative horizons deleting nothing, monotone
//! saves across the direct and receiver forms, replay + wake-driven
//! delivery + watcher-death close on subscribe, key round-trip, the
//! `stream` field carrying the contract name, and entry-point
//! validation (ADR-015; ADR-019 §1/§6; ADR-021 §5; ADR-023 §2).
use std::path::PathBuf;
use std::time::Duration;
use alkstore::{Error, EventReceiver, Store};
use serde_json::{Value, json};
use crate::store::{open, open_store};
fn temp_dir(tag: &str) -> PathBuf {
let d = std::env::temp_dir().join(format!(
"alkstore-stream-{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")
}
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;
}
}
/// Wait for the next event with a bounded deadline — no tight-race
/// assertions (the suite's stability posture).
async fn must_recv_event(receiver: &mut dyn EventReceiver, tag: &str) -> alkstore::StreamEvent {
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
loop {
match receiver.try_recv() {
Ok(Some(event)) => return event,
Ok(None) => {}
Err(e) => panic!("{tag}: receiver errored instead of idling: {e}"),
}
assert!(
tokio::time::Instant::now() < deadline,
"{tag}: event never arrived"
);
tokio::time::sleep(Duration::from_millis(10)).await;
}
}
/// Acceptance: full `StreamHandle` over a temp-file store — publish
/// returns ascending offsets, reads decode into `StreamEvent` with the
/// `stream` field carrying the contract name (not `topic`), keys
/// round-trip exactly (`None` vs `Some`), consumers absent from the
/// table read offset 0, offsets save and inspect (the monotone save's
/// composition pinned with the direct forms below), and the ordering
/// is `offset ASC` (global FIFO, ADR-015 §4).
#[tokio::test(flavor = "multi_thread")]
async fn publish_read_save_get_round_trip_end_to_end() {
let dir = temp_dir("round-trip");
let path = temp_path("round-trip");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
let stream = store.stream("events").await.unwrap();
assert_eq!(stream.name(), "events");
let o1 = stream.publish(json!({"n": 1})).await.unwrap();
let o2 = stream
.publish_with_key(Some("k".to_string()), json!({"n": 2}))
.await
.unwrap();
let o3 = stream.publish(json!({"n": 3})).await.unwrap();
assert!(o1 < o2 && o2 < o3, "offsets ascend per publish");
let page = stream.read_since(0, 10).await.unwrap();
assert_eq!(page.len(), 3, "the full stream reads back");
assert_eq!(page[0].offset, o1);
assert_eq!(page[0].stream, "events", "the stream field is the name");
assert_eq!(page[0].key, None, "plain publish carries no key");
assert_eq!(page[1].offset, o2);
assert_eq!(page[1].key, Some("k".to_string()));
assert_eq!(page[2].offset, o3);
assert_eq!(page[2].key, None);
for (i, event) in page.iter().enumerate() {
let decoded: Value = event.payload_as().unwrap();
assert_eq!(decoded, json!({"n": i as i64 + 1}));
}
// Cursor reads: offset > cursor, ASC.
let tail = stream.read_since(o1, 10).await.unwrap();
assert_eq!(tail.len(), 2);
assert_eq!(tail[0].offset, o2);
// Consumers absent from the table read 0 (the pre-save state).
let from_absent = stream
.read_from_consumer("fresh-consumer", 10)
.await
.unwrap();
assert_eq!(from_absent.len(), 3, "absent consumer reads from 0");
assert_eq!(
stream.get_offset("fresh-consumer").await.unwrap(),
0,
"get_offset of an absent consumer = 0"
);
// Save + inspect; reads from the consumer's checkpoint.
stream.save_offset("c", o2).await.unwrap();
assert_eq!(stream.get_offset("c").await.unwrap(), o2);
let from_c = stream.read_from_consumer("c", 10).await.unwrap();
assert_eq!(from_c.len(), 1);
assert_eq!(from_c[0].offset, o3);
store.close();
cleanup(&dir);
}
/// Acceptance: the extent guard (ADR-023 §2) — `read_since` /
/// `read_from_consumer` with `limit <= 0` return the empty `Vec` (never
/// the whole stream, which SQLite's `LIMIT -1` dialect would yield),
/// at the trait-impl entry.
#[tokio::test(flavor = "multi_thread")]
async fn extent_guard_reads_limit_zero_and_negative_yield_empty() {
let dir = temp_dir("extent-guard");
let store = open_store(
temp_path("extent-guard").to_str().unwrap(),
Default::default(),
)
.unwrap();
let stream = store.stream("s").await.unwrap();
stream.publish(json!({"n": 1})).await.unwrap();
stream.publish(json!({"n": 2})).await.unwrap();
stream.save_offset("c", 6).await.unwrap();
for limit in [0i64, -1, i64::MIN] {
let page = stream.read_since(0, limit).await.unwrap();
assert!(
page.is_empty(),
"read_since(limit={limit}) must read nothing"
);
let page = stream.read_from_consumer("c", limit).await.unwrap();
assert!(
page.is_empty(),
"read_from_consumer(limit={limit}) must read nothing"
);
}
// The stream was never handed out in bulk: the second event is
// still readable with a positive extent.
assert_eq!(stream.read_since(0, 10).await.unwrap().len(), 2);
store.close();
cleanup(&dir);
}
/// Acceptance: `trim_to` is exact-boundary (`offset <= horizon`),
/// returns the deleted count, a negative horizon deletes nothing
/// (boundary args are total, ADR-023 §2 — no guard, pass through),
/// and surviving offsets are never renumbered (gaps legal).
#[tokio::test(flavor = "multi_thread")]
async fn trim_to_trims_the_exact_boundary_and_never_renumbers() {
let dir = temp_dir("trim");
let store = open_store(temp_path("trim").to_str().unwrap(), Default::default()).unwrap();
let stream = store.stream("s").await.unwrap();
let mut offsets = Vec::new();
for i in 0..5 {
offsets.push(stream.publish(json!({"n": i})).await.unwrap());
}
// Negative horizon: deletes nothing, idempotently (no guard).
assert_eq!(
stream.trim_to(-7).await.unwrap(),
0,
"a negative horizon deletes nothing"
);
assert_eq!(stream.read_since(0, 10).await.unwrap().len(), 5);
// Zero horizon is a no-op on this store (offsets start above 0 on
// the autoincrement counter — also nothing to delete).
assert_eq!(stream.trim_to(0).await.unwrap(), 0);
// Exact boundary: horizon = offsets[1] deletes exactly rows
// offsets[0..=1] (two rows), keeping offsets[2..].
let deleted = stream.trim_to(offsets[1]).await.unwrap();
assert_eq!(deleted, 2, "trim deletes offset <= horizon exactly");
let survivors = stream.read_since(0, 10).await.unwrap();
assert_eq!(survivors.len(), 3);
let survivor_offsets: Vec<i64> = survivors.iter().map(|e| e.offset).collect();
assert_eq!(
survivor_offsets,
offsets[2..],
"surviving offsets are never renumbered — gaps are legal"
);
// Reads from a trimmed-away region resume at the horizon's first
// remaining row (fewer/no rows, never an error).
let resumed = stream.read_since(offsets[0], 10).await.unwrap();
assert_eq!(resumed.len(), 3);
// Trim is idempotent below the surviving set.
assert_eq!(stream.trim_to(offsets[1]).await.unwrap(), 0);
store.close();
cleanup(&dir);
}
/// Acceptance: monotone saves (ADR-019 §6) — a direct save below the
/// stored checkpoint is a silent no-op (`Ok(())`, offset unchanged);
/// direct and receiver forms compose without regression.
#[tokio::test(flavor = "multi_thread")]
async fn offset_saves_are_monotone_across_forms() {
let dir = temp_dir("monotone");
let store = open_store(temp_path("monotone").to_str().unwrap(), Default::default()).unwrap();
let stream = store.stream("s").await.unwrap();
let o1 = stream.publish(json!({"n": 1})).await.unwrap();
let o2 = stream.publish(json!({"n": 2})).await.unwrap();
assert!(o1 < o2);
stream.save_offset("c", o2).await.unwrap();
// Regression: silent no-op, not an error (the act-differently rule).
stream.save_offset("c", o1).await.unwrap();
assert_eq!(
stream.get_offset("c").await.unwrap(),
o2,
"a save below the stored checkpoint does not rewind the cursor"
);
stream.save_offset("c", o2).await.unwrap();
assert_eq!(stream.get_offset("c").await.unwrap(), o2);
// The receiver form drives the same op: subscribing a fresh
// consumer, reading to tail, saving through the receiver — then a
// direct regression save cannot undo it.
let mut rx = store
.stream("s")
.await
.unwrap()
.subscribe("r")
.await
.unwrap();
let last = must_recv_event(&mut *rx, "monotone-receiver").await;
rx.save_offset().unwrap();
assert_eq!(rx.offset(), last.offset);
let stored = store
.stream("s")
.await
.unwrap()
.get_offset("r")
.await
.unwrap();
assert_eq!(stored, last.offset, "the receiver save checkpointed");
stream.save_offset("r", o1).await.unwrap();
let after = store
.stream("s")
.await
.unwrap()
.get_offset("r")
.await
.unwrap();
assert_eq!(
after, last.offset,
"the direct form cannot regress the receiver's"
);
// Receiver save before any yield is a no-op (the stored checkpoint
// untouched — ADR-021 §5). A fresh subscription over a consumer
// already at a stored offset: nothing yielded yet, save, unchanged.
let fresh_stream = store.stream("s").await.unwrap();
let mut fresh = fresh_stream.subscribe("pre-yield").await.unwrap();
fresh.save_offset().unwrap();
assert_eq!(fresh.offset(), 0);
assert_eq!(fresh_stream.get_offset("pre-yield").await.unwrap(), 0);
store.close();
cleanup(&dir);
}
/// Acceptance: `subscribe` replays from the *stored* offset (a
/// consumer with a prior checkpoint resumes there, not from 0), then
/// delivers post-attach events wake-driven — events published after
/// attach arrive through the receiver.
#[tokio::test(flavor = "multi_thread")]
async fn subscribe_replays_from_stored_offset_then_delivers_wake_driven() {
let dir = temp_dir("subscribe");
let path = temp_path("subscribe");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
let pre = store.stream("s").await.unwrap();
let o1 = pre.publish(json!({"n": 1})).await.unwrap();
let o2 = pre.publish(json!({"n": 2})).await.unwrap();
let o3 = pre.publish(json!({"n": 3})).await.unwrap();
drop(pre);
// A stored checkpoint below the tail: replay starts there (o2 is
// the first event after the checkpoint).
let handle = store.stream("s").await.unwrap();
handle.save_offset("c", o1).await.unwrap();
let mut rx = handle.subscribe("c").await.unwrap();
drop(handle);
let first = must_recv_event(&mut *rx, "subscribe-replay").await;
assert_eq!(
first.offset, o2,
"replay resumes past the stored checkpoint"
);
let second = must_recv_event(&mut *rx, "subscribe-replay").await;
assert_eq!(second.offset, o3);
assert_eq!(second.stream, "s");
// Post-attach publishes arrive wake-driven — no polling needed.
let late = store
.stream("s")
.await
.unwrap()
.publish_with_key(Some("late".to_string()), json!({"n": 4}))
.await
.unwrap();
let wake_event = must_recv_event(&mut *rx, "subscribe-wake").await;
assert_eq!(wake_event.offset, late);
assert_eq!(wake_event.key.as_deref(), Some("late"));
// Receiver checkpoint composes with the substrate's monotone upsert.
rx.save_offset().unwrap();
let stored = store
.stream("s")
.await
.unwrap()
.get_offset("c")
.await
.unwrap();
assert_eq!(stored, wake_event.offset);
store.close();
cleanup(&dir);
}
/// Acceptance: key round-trips exactly — `None` vs `Some` preserved on
/// every read form (direct page and receiver delivery), and the
/// `stream` field carries the contract name (the `topic` column is
/// never surfaced).
#[tokio::test(flavor = "multi_thread")]
async fn keys_round_trip_exactly_and_stream_field_carries_the_name() {
let dir = temp_dir("keys");
let store = open_store(temp_path("keys").to_str().unwrap(), Default::default()).unwrap();
let stream = store.stream("the-name").await.unwrap();
let plain = stream.publish(json!({})).await.unwrap();
let keyed = stream
.publish_with_key(Some("k".to_string()), json!({}))
.await
.unwrap();
let page = stream.read_since(0, 10).await.unwrap();
fn find(events: &[alkstore::StreamEvent], offset: i64) -> Option<&alkstore::StreamEvent> {
events.iter().find(|e| e.offset == offset)
}
assert_eq!(find(&page, plain).unwrap().key, None);
assert_eq!(find(&page, keyed).unwrap().key, Some("k".to_string()));
assert!(page.iter().all(|e| e.stream == "the-name"));
let mut rx = store
.stream("the-name")
.await
.unwrap()
.subscribe("c")
.await
.unwrap();
let e1 = must_recv_event(&mut *rx, "keys-receiver").await;
let e2 = must_recv_event(&mut *rx, "keys-receiver").await;
assert_eq!(e1.key, None);
assert_eq!(e2.key, Some("k".to_string()));
assert_eq!(e1.stream, "the-name");
store.close();
cleanup(&dir);
}
/// Acceptance: watcher death closes the subscription terminally —
/// `recv() -> None` (never reopens); `try_recv` reports `Err(Closed)`.
/// Driven through store close (the death-guard clears every
/// subscriber feed — the bridge exits, the consumer sees the close).
#[tokio::test(flavor = "multi_thread")]
async fn watcher_death_closes_the_subscription_terminally() {
let dir = temp_dir("death");
let store = open_store(temp_path("death").to_str().unwrap(), Default::default()).unwrap();
let stream = store.stream("s").await.unwrap();
stream.publish(json!({"n": 1})).await.unwrap();
let mut rx = stream.subscribe("c").await.unwrap();
let event = must_recv_event(&mut *rx, "death").await;
assert_eq!(event.stream, "s");
store.close();
let none = rx.recv().await;
assert!(
none.is_none(),
"watcher death closes the receiver: got {none:?}"
);
assert!(
matches!(rx.try_recv(), Err(Error::Closed)),
"after close, try_recv is Err(Closed) — not idle"
);
assert!(rx.recv().await.is_none(), "a closed receiver never reopens");
cleanup(&dir);
}
/// Acceptance: a subscribe attached while the stream is at tail idles
/// (`Ok(None)`), then delivers when the next publish lands.
#[tokio::test(flavor = "multi_thread")]
async fn subscribe_at_tail_idles_then_delivers() {
let dir = temp_dir("tail-idle");
let store = open_store(temp_path("tail-idle").to_str().unwrap(), Default::default()).unwrap();
let stream = store.stream("s").await.unwrap();
let mut rx = stream.subscribe("c").await.unwrap();
tokio::time::sleep(Duration::from_millis(150)).await;
assert_eq!(rx.try_recv().unwrap(), None, "at tail, the receiver idles");
let o = stream.publish(json!({"n": 1})).await.unwrap();
let event = must_recv_event(&mut *rx, "tail-idle").await;
assert_eq!(event.offset, o);
store.close();
cleanup(&dir);
}
/// Acceptance: two consumers subscribe independently — each drains a
/// copy of the stream from its own position; one's receiver save never
/// gates the other's delivery (saved checkpoints never gate reads).
#[tokio::test(flavor = "multi_thread")]
async fn subscribers_are_independent_copies() {
let dir = temp_dir("two-consumers");
let store = open_store(
temp_path("two-consumers").to_str().unwrap(),
Default::default(),
)
.unwrap();
let stream = store.stream("s").await.unwrap();
let o1 = stream.publish(json!({"n": 1})).await.unwrap();
let o2 = stream.publish(json!({"n": 2})).await.unwrap();
let handle = store.stream("s").await.unwrap();
let mut a = handle.subscribe("a").await.unwrap();
let mut b = handle.subscribe("b").await.unwrap();
drop(handle);
// a drains and checkpoints both; b must still receive both (a's
// position is a's, not shared).
let _ = must_recv_event(&mut *a, "two-a-1").await;
let _ = must_recv_event(&mut *a, "two-a-2").await;
a.save_offset().unwrap();
let b1 = must_recv_event(&mut *b, "two-b-1").await;
let b2 = must_recv_event(&mut *b, "two-b-2").await;
assert_eq!(b1.offset, o1);
assert_eq!(b2.offset, o2);
store.close();
cleanup(&dir);
}
/// Acceptance: `recv()`'s `Err` arm carries `Database` only
/// (ADR-021 §5) and close is terminal. Driving the transient-failure
/// arm against a healthy engine is not directly reproducible (a
/// healthy SELECT on a WAL file does not fail); the arm's shape is
/// pinned here by construction — the receiver yields whatever the
/// bridge read produced, and this test pins the observable close
/// semantics across drop (dropping the receiver's channel source).
/// The compile-time shape (`Option<Result<StreamEvent>>`) plus the
/// `Database`-only doc rule are the pinned surface; this test pins the
/// runtime close arm over the drop path.
#[tokio::test(flavor = "multi_thread")]
async fn receiver_drop_closes_its_own_subscription() {
let dir = temp_dir("receiver-drop");
let store = open_store(
temp_path("receiver-drop").to_str().unwrap(),
Default::default(),
)
.unwrap();
let stream = store.stream("s").await.unwrap();
stream.publish(json!({})).await.unwrap();
let baseline = store.watcher.subscriber_count();
assert_eq!(baseline, 0, "no subscriptions before any subscribe");
{
let handle = store.stream("s").await.unwrap();
let rx = handle.subscribe("c").await.unwrap();
drop(rx);
}
// The receiver's Drop unsubscribed eagerly.
until("receiver-drop", || {
store.watcher.subscriber_count() == baseline
})
.await;
assert_eq!(
store.watcher.subscriber_count(),
baseline,
"receiver drop must not leak subscriptions"
);
store.close();
cleanup(&dir);
}
/// Acceptance: entry-point validation — stream names validate as
/// shared-namespace kinds at the constructor; consumers as local names
/// (non-empty only; reserved prefixes legal for consumers). Validation
/// fires before any round trip: rejected constructors yield no handle,
/// and rejected consumer names never touch the tables.
#[tokio::test(flavor = "multi_thread")]
async fn entry_points_validate_stream_names_and_consumers() {
let dir = temp_dir("validation");
let path = temp_path("validation");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
// Stream constructors: empty/whitespace → InvalidName; reserved
// prefix → ReservedName.
for name in ["", " ", "__alkstore_x"] {
let handle = match store.stream(name).await {
Err(e) => e,
Ok(_) => panic!("stream({name:?}) must reject by validation"),
};
match handle {
Error::InvalidName { name: n } => assert_eq!(name, n),
Error::ReservedName { name: n } => assert_eq!(name, n),
other => panic!("stream({name:?}) must reject by validation, got {other}"),
}
}
let stream = store.stream("s").await.unwrap();
// Empty-`Some` keys are InvalidName (the tx path's rule).
let err = stream
.publish_with_key(Some(String::new()), json!({}))
.await
.unwrap_err();
match err {
Error::InvalidName { name } => assert_eq!(name, ""),
other => panic!("empty-Some key must be InvalidName, got {other}"),
}
let err = stream
.publish_with_key(Some(" ".to_string()), json!({}))
.await
.unwrap_err();
assert!(matches!(err, Error::InvalidName { .. }));
// Consumers: non-empty only — reserved prefixes are legal.
for consumer in ["", " "] {
let err = stream.read_from_consumer(consumer, 5).await.unwrap_err();
assert!(matches!(err, Error::InvalidName { .. }));
let err = stream.save_offset(consumer, 1).await.unwrap_err();
assert!(matches!(err, Error::InvalidName { .. }));
let err = stream.get_offset(consumer).await.unwrap_err();
assert!(matches!(err, Error::InvalidName { .. }));
assert!(stream.subscribe(consumer).await.is_err());
}
// A legal-but-reserved consumer name works (non-empty only).
let o = stream.publish(json!({"n": 9})).await.unwrap();
stream.save_offset("__alkstore_local", o).await.unwrap();
assert_eq!(stream.get_offset("__alkstore_local").await.unwrap(), o);
// The rejected consumer saves never touched the table.
assert_eq!(stream.get_offset("never-saved").await.unwrap(), 0);
store.close();
cleanup(&dir);
}
/// Acceptance: `StreamHandle` ops on a closed store fail closed with
/// `Error::Database` (the engine-wide closed-store posture); a
/// subscribe on a closed store fails too.
#[tokio::test(flavor = "multi_thread")]
async fn closed_store_fails_stream_ops() {
let dir = temp_dir("closed-store");
let store = open_store(
temp_path("closed-store").to_str().unwrap(),
Default::default(),
)
.unwrap();
let stream = store.stream("s").await.unwrap();
stream.publish(json!({"n": 1})).await.unwrap();
let stored_handle = store.stream("s").await.unwrap();
store.close();
let err = stream.publish(json!({})).await.unwrap_err();
assert!(matches!(err, Error::Database(_)), "publish: {err}");
let err = stream.read_since(0, 5).await.unwrap_err();
assert!(matches!(err, Error::Database(_)), "read_since: {err}");
let err = stream.read_from_consumer("c", 5).await.unwrap_err();
assert!(
matches!(err, Error::Database(_)),
"read_from_consumer: {err}"
);
let err = stream.save_offset("c", 1).await.unwrap_err();
assert!(matches!(err, Error::Database(_)), "save_offset: {err}");
let err = stream.get_offset("c").await.unwrap_err();
assert!(matches!(err, Error::Database(_)), "get_offset: {err}");
let err = stream.trim_to(10).await.unwrap_err();
assert!(matches!(err, Error::Database(_)), "trim_to: {err}");
assert!(stored_handle.subscribe("c").await.is_err());
// New constructors on a closed store fail too.
let err = match store.stream("s").await {
Err(e) => e,
Ok(_) => panic!("stream on a closed store must fail closed"),
};
assert!(matches!(err, Error::Database(_)));
cleanup(&dir);
}
/// Acceptance: the durable subscription survives store restart — a
/// fresh open over the same file resumes from the saved checkpoint
/// (durable consumption; offsets are positions in the durable log).
#[tokio::test(flavor = "multi_thread")]
async fn subscription_state_survives_store_restart() {
let dir = temp_dir("restart");
let path = temp_path("restart");
let saved;
{
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
let stream = store.stream("s").await.unwrap();
let o1 = stream.publish(json!({"n": 1})).await.unwrap();
let _ = stream.publish(json!({"n": 2})).await.unwrap();
stream.save_offset("c", o1).await.unwrap();
saved = stream.get_offset("c").await.unwrap();
store.close();
}
{
let store = open(path.to_str().unwrap(), Default::default()).unwrap();
let stream = store.stream("s").await.unwrap();
// Checkpoint durable across the restart.
assert_eq!(stream.get_offset("c").await.unwrap(), saved);
// Replay resumes from it — the second event is the first read.
let page = stream.read_from_consumer("c", 10).await.unwrap();
assert_eq!(page.len(), 1, "the resumed read starts past the checkpoint");
let saved_offset = saved;
// A fresh subscription replays from the saved offset.
let mut rx = stream.subscribe("c").await.unwrap();
tokio::time::sleep(Duration::from_millis(200)).await;
match rx.try_recv().unwrap() {
Some(event) => assert_eq!(event.offset, saved_offset + 1),
None => panic!("the resumed subscription must replay past the checkpoint"),
}
drop(rx);
drop(stream);
drop(store);
}
cleanup(&dir);
}
/// Acceptance: the tx seam and the handle compose — a `publish_tx`
/// inside a committed business transaction is visible to the handle's
/// reads and wakes a live subscriber; a rolled-back `publish_tx` ghosts
/// (no read, no wake).
#[tokio::test(flavor = "multi_thread")]
async fn tx_publishes_compose_with_the_handle() {
let dir = temp_dir("tx-compose");
let path = temp_path("tx-compose");
let store = open_store(path.to_str().unwrap(), Default::default()).unwrap();
let mut rx = store
.stream("s")
.await
.unwrap()
.subscribe("c")
.await
.unwrap();
let mut tx = store.begin_tx().await.unwrap();
let tx_offset = tx.publish_tx("s", json!({"via": "tx"})).await.unwrap();
tx.commit().await.unwrap();
let event = must_recv_event(&mut *rx, "tx-compose").await;
assert_eq!(event.offset, tx_offset);
assert_eq!(event.payload_as::<Value>().unwrap(), json!({"via": "tx"}));
// Rollback: the publish ghosts.
let mut tx = store.begin_tx().await.unwrap();
tx.publish_tx("s", json!({"via": "ghost"})).await.unwrap();
drop(tx);
tokio::time::sleep(Duration::from_millis(200)).await;
while rx.try_recv().unwrap().is_some() {}
tokio::time::sleep(Duration::from_millis(200)).await;
assert_eq!(
rx.try_recv().unwrap(),
None,
"a rolled-back publish must not deliver"
);
// The handle read confirms the table state.
let stream = store.stream("s").await.unwrap();
let page = stream.read_since(0, 10).await.unwrap();
assert_eq!(page.len(), 1, "the ghost never landed");
assert_eq!(page[0].offset, tx_offset);
store.close();
cleanup(&dir);
}
/// Acceptance: the receiver's `read_since` is offset-anchored and
/// independent of its position (saved checkpoints do not gate reads —
/// replay is a read concern); extent-guarded like the direct form.
#[tokio::test(flavor = "multi_thread")]
async fn receiver_read_since_is_offset_anchored_and_extent_guarded() {
let dir = temp_dir("receiver-read");
let store = open_store(
temp_path("receiver-read").to_str().unwrap(),
Default::default(),
)
.unwrap();
let stream = store.stream("s").await.unwrap();
let o1 = stream.publish(json!({"n": 1})).await.unwrap();
let o2 = stream.publish(json!({"n": 2})).await.unwrap();
let mut rx = stream.subscribe("c").await.unwrap();
let _ = must_recv_event(&mut *rx, "receiver-read").await;
let _ = must_recv_event(&mut *rx, "receiver-read").await;
rx.save_offset().unwrap();
assert_eq!(rx.offset(), o2);
// The receiver replays from any offset regardless of its position.
let replay = rx.read_since(o1 - 1, 10).await.unwrap();
assert_eq!(replay.len(), 2);
assert_eq!(replay[0].offset, o1);
// Extent guard on the receiver form too.
assert_eq!(rx.read_since(0, 0).await.unwrap().len(), 0);
assert_eq!(rx.read_since(0, -3).await.unwrap().len(), 0);
store.close();
cleanup(&dir);
}
+645
View File
@@ -0,0 +1,645 @@
//! The stream mechanism's SQLite arm (ADR-015; ADR-019 §1/§6; ADR-021
//! §5; ADR-023 §2):
//!
//! - `stream(name)` — the validated constructor. Construction is free
//! (validation + a closed-store check only — the table exists from
//! bootstrap); the handle carries the connection machinery's arcs.
//! - Auto-commit publishes — an append to the durable log through the
//! writer slot, returning the engine-assigned offset; keys are
//! carried metadata (ADR-015 §1) with the non-empty-key rule (an
//! empty `Some` key is `InvalidName`, matching the tx path).
//! - Cursor reads — reader-pool reads, `offset ASC` (global FIFO,
//! ADR-015 §4), decoded into [`StreamEvent`] via the core
//! constructor's `from_row` (the `topic` column carries the
//! contract's `stream` field — the one name delta). The extent guard
//! (ADR-023 §2) short-circuits `limit <= 0` to the empty `Vec` at
//! the trait-impl entry, before the substrate — the `LIMIT -1`
//! dialect artifact dies here.
//! - Checkpoints — explicit, monotone saves (a save below the stored
//! checkpoint is a silent no-op, ADR-019 §6); an absent consumer
//! reads 0.
//! - `trim_to` — the `DELETE … WHERE offset <= ?` bounded-growth op
//! inside the writer-slot lease (ADR-015 §5); the horizon is a
//! boundary argument, total per ADR-023 §2 (a negative horizon
//! deletes nothing, idempotently — no guard, pass through).
//! Surviving events keep their offsets.
//! - `subscribe` — the durable receiver: attach, replay from the
//! stored offset to the current tail, then wake-driven re-reads off
//! the watcher fanout (the same bridge shape as `listen`, yielding
//! `Result<StreamEvent>` per `EventReceiver`). Watcher death closes
//! it terminally (`recv() -> None`); recovery is a fresh
//! `subscribe(consumer)` resuming from the saved checkpoint
//! (ADR-021 §5).
//!
//! Subscribe specifics:
//!
//! - The watcher subscription is taken *before* any read, so no commit
//! can fall between the attach read and the wake wait (a tick for a
//! commit already re-read sits buffered; a drain that sees it costs
//! an empty re-poll, never a missed event).
//! - The bridge drains pages ([`SUBSCRIBE_PAGE`]) to the tail, then
//! parks on the wake feed. Wakes are coalesced hints (overtrigger on
//! purpose, ADR-006): a wake may arrive for unrelated writes — the
//! re-read from the in-memory cursor sees nothing new and parks
//! again.
//! - A read failure while the store is open surfaces one
//! `Err(Database)` on the receiver and waits for the next wake to
//! re-read (transient arm); a failure after store close exits the
//! loop (the consumer sees the close). An attach-read failure
//! surfaces one `Err(Database)` and closes the receiver — recovery
//! is a fresh `subscribe(consumer)`. Non-`Database` read errors
//! (the decode-side `Codec` posture guard) are remapped into
//! `Database` with the source chain preserved — `recv()`'s `Err` arm
//! carries `Database` only (ADR-021 §5).
//! - The bridge pushes drained events into a bounded tokio channel
//! and `blocking_send`s when full — a stalled consumer parks the
//! bridge thread, and delivery resumes exactly where it stopped. The
//! receiver's position is the last-*yielded* event's offset, so
//! `save_offset()` never checkpoints an event the consumer did not
//! observe.
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use parking_lot::Mutex;
use serde_json::Value;
use alkstore::{
BoxedFuture, Error, EventReceiver, Result, StreamEvent, StreamHandle, validate_local_name,
validate_shared_name,
};
use crate::seam::{blocking, sqlite_error, with_writer};
use crate::substrate::{
Readers, SchemaError, SharedUpdateWatcher, Writer, stream_get_offset, stream_publish,
stream_read_since, stream_save_offset,
};
/// How many events a wake-driven re-read fetches per page while
/// draining to the tail.
const SUBSCRIBE_PAGE: i64 = 256;
/// The bridge's tokio-channel capacity — the smoothing buffer between
/// the bridge thread's re-read loop and the consumer (the substrate's
/// 1-slot sync feed stays the coalescing point).
const EVENT_CHANNEL_CAPACITY: usize = 16;
/// The auto-commit publish paths (see the module docs). Stream-name
/// validation fires at the entry point, before any round trip and
/// before the payload is touched (ADR-008 §4); a present key must be
/// non-empty.
fn publish_with_key(
writer: Arc<Writer>,
stream: &str,
key: Option<String>,
payload: Value,
) -> BoxedFuture<'static, Result<i64>> {
if let Err(e) = validate_shared_name(stream) {
return Box::pin(async move { Err(e) });
}
let key = match key {
Some(k) if k.is_empty() || k.trim().is_empty() => {
return Box::pin(async move { Err(Error::InvalidName { name: k }) });
}
other => other,
};
let stream = stream.to_string();
Box::pin(async move {
with_writer(writer, move |conn| {
let bytes = alkstore::encode_payload(&payload)?;
stream_publish(
conn,
&stream,
key.as_deref(),
std::str::from_utf8(&bytes)
.map_err(|e| Error::Codec(format!("row payload must be utf-8: {e}")))?,
)
.map_err(sqlite_error)
})
.await
})
}
/// A reader-pool read through the `spawn_blocking` seam: acquire a
/// pooled reader, run the body, release on both arms (a failed SELECT
/// leaves the pooled connection healthy). A closed pool (store closed)
/// fails closed with `Database`.
async fn with_reader<T, F>(readers: Arc<Readers>, f: F) -> Result<T>
where
T: Send + 'static,
F: FnOnce(&rusqlite::Connection) -> Result<T> + Send + 'static,
{
blocking(move || {
let conn = acquire_reader(&readers)?;
let out = f(&conn);
readers.release(conn);
out
})
.await
}
fn acquire_reader(readers: &Arc<Readers>) -> Result<rusqlite::Connection> {
match readers.acquire() {
Ok(conn) => Ok(conn),
Err(SchemaError::Sqlite(e)) => {
if is_closed_err(&e) {
Err(closed_store_error())
} else {
Err(sqlite_error(e))
}
}
}
}
fn is_closed_err(e: &rusqlite::Error) -> bool {
matches!(
e,
rusqlite::Error::SqliteFailure(_, Some(msg)) if msg.contains("Database is closed")
)
}
fn closed_store_error() -> Error {
Error::database(std::io::Error::other("the store is closed"))
}
/// `recv()`'s `Err` arm carries `Database` only (ADR-021 §5): a
/// read error that mapped elsewhere (the decode-side `Codec` posture
/// guard) is remapped into the opaque fallback with the original error
/// preserved as the source.
fn database_only(e: Error) -> Error {
match e {
already @ Error::Database(_) => already,
other => Error::database(other),
}
}
/// One decoded page from the reader pool: `offset > from`, ASC, up to
/// `limit` (the extent guard is the caller's — this helper is
/// contract-blind to it).
fn read_page(
readers: Arc<Readers>,
stream: String,
from: i64,
limit: i64,
) -> BoxedFuture<'static, Result<Vec<StreamEvent>>> {
let read = move |conn: &rusqlite::Connection| {
let page = stream_read_since(conn, &stream, from, limit).map_err(sqlite_error)?;
stream_events_from_json_page(&page, stream.clone())
};
Box::pin(with_reader(readers, read))
}
/// Decode the substrate's stream-page JSON (the fork's inherited
/// return shape) into the contract's [`StreamEvent`] values via the
/// core constructor's `from_row`. Shared by the auto-commit reads
/// (this module) and the tx reads (`tx.rs`) — one decode owner. The
/// `topic` column carries the contract's `stream` field (the one name
/// delta, mapped at this layer); malformed rows map to `Error::Codec`
/// (decode-side failures is exactly what `Codec` is for — ADR-021 §2's
/// posture guard; the fork's writer always wrote well-formed rows).
pub(crate) fn stream_events_from_json_page(page: &str, stream: String) -> Result<Vec<StreamEvent>> {
let array: serde_json::Value = serde_json::from_str(page)
.map_err(|e| Error::Codec(format!("stream page must be json: {e}")))?;
let rows = array
.as_array()
.ok_or_else(|| Error::Codec("stream page must be a json array".to_string()))?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
out.push(StreamEvent::from_row(
serde_json::from_value(row["offset"].clone())
.map_err(|e| Error::Codec(format!("stream row offset must be i64: {e}")))?,
stream.clone(),
if row["key"].is_null() {
None
} else {
serde_json::from_value(row["key"].clone())
.map_err(|e| Error::Codec(format!("stream row key must be a string: {e}")))?
},
row["payload"]
.as_str()
.ok_or_else(|| Error::Codec("stream row payload must be a string".to_string()))?
.as_bytes()
.to_vec(),
serde_json::from_value(row["created_at"].clone())
.map_err(|e| Error::Codec(format!("stream row created_at must be i64: {e}")))?,
));
}
Ok(out)
}
/// The stream-scoped handle (see the module docs).
pub struct SqliteStreamHandle {
name: String,
writer: Arc<Writer>,
readers: Arc<Readers>,
watcher: Arc<SharedUpdateWatcher>,
closed: Arc<AtomicBool>,
}
impl std::fmt::Debug for SqliteStreamHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SqliteStreamHandle")
.field("name", &self.name)
.finish_non_exhaustive()
}
}
impl SqliteStreamHandle {
pub(crate) fn new(
name: String,
writer: Arc<Writer>,
readers: Arc<Readers>,
watcher: Arc<SharedUpdateWatcher>,
closed: Arc<AtomicBool>,
) -> Self {
Self {
name,
writer,
readers,
watcher,
closed,
}
}
}
impl StreamHandle for SqliteStreamHandle {
fn name(&self) -> &str {
&self.name
}
fn publish<'a>(&'a self, payload: Value) -> BoxedFuture<'a, Result<i64>> {
publish_with_key(self.writer.clone(), &self.name, None, payload)
}
fn publish_with_key<'a>(
&'a self,
key: Option<String>,
payload: Value,
) -> BoxedFuture<'a, Result<i64>> {
publish_with_key(self.writer.clone(), &self.name, key, payload)
}
fn read_since<'a>(
&'a self,
offset: i64,
limit: i64,
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
if limit <= 0 {
return Box::pin(async { Ok(Vec::new()) });
}
read_page(self.readers.clone(), self.name.clone(), offset, limit)
}
fn read_from_consumer<'a>(
&'a self,
consumer: &str,
limit: i64,
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
if let Err(e) = validate_local_name(consumer) {
return Box::pin(async move { Err(e) });
}
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
if limit <= 0 {
return Box::pin(async { Ok(Vec::new()) });
}
let readers = self.readers.clone();
let stream = self.name.clone();
let consumer = consumer.to_string();
let read = move |conn: &rusqlite::Connection| {
let from = stream_get_offset(conn, &consumer, &stream).map_err(sqlite_error)?;
let page = stream_read_since(conn, &stream, from, limit).map_err(sqlite_error)?;
stream_events_from_json_page(&page, stream.clone())
};
Box::pin(with_reader(readers, read))
}
fn save_offset<'a>(&'a self, consumer: &str, offset: i64) -> BoxedFuture<'a, Result<()>> {
if let Err(e) = validate_local_name(consumer) {
return Box::pin(async move { Err(e) });
}
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
let writer = self.writer.clone();
let stream = self.name.clone();
let consumer = consumer.to_string();
Box::pin(async move {
with_writer(writer, move |conn| {
stream_save_offset(conn, &consumer, &stream, offset)
.map(|_| ())
.map_err(sqlite_error)
})
.await
})
}
fn get_offset<'a>(&'a self, consumer: &str) -> BoxedFuture<'a, Result<i64>> {
if let Err(e) = validate_local_name(consumer) {
return Box::pin(async move { Err(e) });
}
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
let readers = self.readers.clone();
let stream = self.name.clone();
let consumer = consumer.to_string();
let read = move |conn: &rusqlite::Connection| {
stream_get_offset(conn, &consumer, &stream).map_err(sqlite_error)
};
Box::pin(with_reader(readers, read))
}
fn trim_to<'a>(&'a self, horizon: i64) -> BoxedFuture<'a, Result<i64>> {
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
let writer = self.writer.clone();
let stream = self.name.clone();
Box::pin(async move {
with_writer(writer, move |conn| {
let deleted = conn
.execute(
"DELETE FROM __alkstore_stream WHERE topic = ?1 AND offset <= ?2",
rusqlite::params![stream, horizon],
)
.map_err(sqlite_error)?;
Ok(deleted as i64)
})
.await
})
}
fn subscribe<'a>(&'a self, consumer: &str) -> BoxedFuture<'a, Result<Box<dyn EventReceiver>>> {
if let Err(e) = validate_local_name(consumer) {
return Box::pin(async move { Err(e) });
}
if self.closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
let name = self.name.clone();
let consumer = consumer.to_string();
let writer = self.writer.clone();
let readers = self.readers.clone();
let watcher = self.watcher.clone();
let closed = self.closed.clone();
Box::pin(async move {
spawn_bridge(name, consumer, writer, readers, watcher, closed)
.map(|rx| Box::new(rx) as Box<dyn EventReceiver>)
})
}
}
/// The `subscribe` bridge construction: the watcher subscription is
/// taken before any read, the bridge thread then replays and parks on
/// the wake feed. A thread-spawn failure unsubscribes and fails with
/// `Database` (no panics cross the seam, the W-2 posture).
fn spawn_bridge(
name: String,
consumer: String,
writer: Arc<Writer>,
readers: Arc<Readers>,
watcher: Arc<SharedUpdateWatcher>,
store_closed: Arc<AtomicBool>,
) -> Result<SqliteEventReceiver> {
let (id, wake_rx) = watcher.subscribe();
let wake_rx = Arc::new(Mutex::new(wake_rx));
let (event_tx, event_rx) =
tokio::sync::mpsc::channel::<Result<StreamEvent>>(EVENT_CHANNEL_CAPACITY);
let thread_name = name.clone();
let thread_consumer = consumer.clone();
let thread_readers = readers.clone();
let thread_closed = store_closed.clone();
let thread_tx = event_tx.clone();
let spawned = std::thread::Builder::new()
.name(format!("alkstore-stream-bridge-{thread_name}"))
.spawn(move || {
run_subscribe_loop(
thread_name,
thread_consumer,
thread_readers,
thread_closed,
wake_rx,
thread_tx,
);
});
if let Err(e) = spawned {
watcher.unsubscribe(id);
return Err(Error::database(std::io::Error::other(format!(
"failed to spawn stream bridge thread: {e}"
))));
}
Ok(SqliteEventReceiver {
name,
consumer,
rx: event_rx,
watcher,
id,
position: 0,
writer,
readers,
store_closed,
})
}
/// The bridge loop body (runs on its own thread — this file's honest
/// sync shape, mirroring `listen`'s bridge thread):
///
/// 1. read the consumer's stored offset (the attach read),
/// 2. drain pages until the tail (`offset ASC`),
/// 3. park on the wake feed (any commit ticks it), re-drain on wake,
/// 4. a read failure while the store is open surfaces one
/// `Err(Database)` and waits for the next wake; a read failure
/// after store close exits (the death path — watcher death via the
/// cleared feed, or store close, disconnects the wake feed and ends
/// the loop; the receiver sees the channel close).
fn run_subscribe_loop(
name: String,
consumer: String,
readers: Arc<Readers>,
store_closed: Arc<AtomicBool>,
wake_rx: Arc<Mutex<std::sync::mpsc::Receiver<()>>>,
event_tx: tokio::sync::mpsc::Sender<Result<StreamEvent>>,
) {
let mut cursor = match blocking_read(&readers, &store_closed, |conn| {
stream_get_offset(conn, &consumer, &name).map_err(sqlite_error)
}) {
Ok(offset) => offset,
Err(LoopExit::Closed) => return,
Err(LoopExit::Failed(e)) => {
let _ = event_tx.blocking_send(Err(database_only(e)));
return;
}
};
loop {
loop {
match blocking_read(&readers, &store_closed, |conn| {
let page =
stream_read_since(conn, &name, cursor, SUBSCRIBE_PAGE).map_err(sqlite_error)?;
stream_events_from_json_page(&page, name.clone())
}) {
Err(LoopExit::Closed) => return,
Err(LoopExit::Failed(e)) => {
if event_tx.blocking_send(Err(database_only(e))).is_err() {
return;
}
break;
}
Ok(events) if events.is_empty() => break,
Ok(events) => {
for event in events {
cursor = event.offset;
if event_tx.blocking_send(Ok(event)).is_err() {
return;
}
}
}
}
}
match wake_rx.lock().recv() {
Ok(()) => {}
// Watcher death: the death-guard cleared the feed. Close.
Err(_) => return,
}
}
}
enum LoopExit {
Closed,
Failed(Error),
}
/// A synchronous pooled read used by the bridge thread. The closed
/// arms are distinguished: a closed store exits the loop, an open
/// store's failure is transient (`Failed` carries the mapped error).
fn blocking_read<T>(
readers: &Arc<Readers>,
store_closed: &Arc<AtomicBool>,
f: impl FnOnce(&rusqlite::Connection) -> Result<T>,
) -> std::result::Result<T, LoopExit> {
if store_closed.load(Ordering::Acquire) {
return Err(LoopExit::Closed);
}
let conn = match readers.acquire() {
Ok(conn) => conn,
Err(e) => match e {
SchemaError::Sqlite(inner) if is_closed_err(&inner) => return Err(LoopExit::Closed),
SchemaError::Sqlite(inner) => {
return Err(LoopExit::Failed(sqlite_error(inner)));
}
},
};
let out = f(&conn);
readers.release(conn);
match out {
Ok(value) => Ok(value),
Err(_e) if store_closed.load(Ordering::Acquire) => Err(LoopExit::Closed),
Err(e) => Err(LoopExit::Failed(e)),
}
}
/// The durable subscription receiver: the tokio side of the bridge,
/// owning the consumer's position state. Its [`Drop`] unsubscribes the
/// substrate feed eagerly (before the watcher's next fanout would) —
/// the bridge thread exits on the failed send or the disconnected wake
/// feed, and the consumer sees the channel close (`None`, terminal).
///
/// `position` is the last-yielded event's offset (0 before any yield —
/// `save_offset()` is then a no-op, the stored checkpoint untouched,
/// ADR-021 §5; ADR-019 §6).
struct SqliteEventReceiver {
name: String,
consumer: String,
rx: tokio::sync::mpsc::Receiver<Result<StreamEvent>>,
watcher: Arc<SharedUpdateWatcher>,
id: u64,
position: i64,
writer: Arc<Writer>,
readers: Arc<Readers>,
store_closed: Arc<AtomicBool>,
}
impl Drop for SqliteEventReceiver {
fn drop(&mut self) {
self.watcher.unsubscribe(self.id);
}
}
impl EventReceiver for SqliteEventReceiver {
fn recv<'a>(&'a mut self) -> BoxedFuture<'a, Option<Result<StreamEvent>>> {
Box::pin(async move {
match self.rx.recv().await {
Some(Ok(event)) => {
self.position = event.offset;
Some(Ok(event))
}
// A transient read failure — the bridge keeps running
// and re-reads on the next wake.
Some(Err(e)) => Some(Err(e)),
// The bridge exited: watcher death or store close —
// `None` = closed, terminal (it never reopens).
None => None,
}
})
}
fn try_recv(&mut self) -> Result<Option<StreamEvent>> {
match self.rx.try_recv() {
Ok(Ok(event)) => {
self.position = event.offset;
Ok(Some(event))
}
Ok(Err(e)) => Err(e),
Err(tokio::sync::mpsc::error::TryRecvError::Empty) => Ok(None),
Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => Err(Error::Closed),
}
}
fn read_since<'a>(
&'a mut self,
offset: i64,
limit: i64,
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
if self.store_closed.load(Ordering::Acquire) {
return Box::pin(async { Err(closed_store_error()) });
}
if limit <= 0 {
return Box::pin(async { Ok(Vec::new()) });
}
read_page(self.readers.clone(), self.name.clone(), offset, limit)
}
fn save_offset(&mut self) -> Result<()> {
let position = self.position;
if position == 0 {
return Ok(());
}
// The save runs through the same monotone op the direct form
// uses — the forms interleave and cannot regress (ADR-019 §6).
// The trait's signature is sync, so the storage round trip
// takes the writer slot on this thread (the slot is nearly
// always free between ops; a long open transaction parks this
// call — the honest lease posture surfaced, ADR-007).
if self.store_closed.load(Ordering::Acquire) {
return Err(closed_store_error());
}
let conn = match self.writer.acquire() {
Some(conn) => conn,
None => return Err(closed_store_error()),
};
let out = stream_save_offset(&conn, &self.consumer, &self.name, position)
.map(|_| ())
.map_err(sqlite_error);
self.writer.release(conn);
out
}
fn offset(&self) -> i64 {
self.position
}
}
+1 -33
View File
@@ -32,6 +32,7 @@ use crate::resolution::{
resolve_enqueue_opts,
};
use crate::seam::{blocking, sqlite_error};
use crate::stream::stream_events_from_json_page;
use crate::substrate::{
Stamps, Writer, enqueue, get_job, stream_get_offset, stream_publish, stream_read_since,
stream_save_offset,
@@ -491,36 +492,3 @@ fn job_from_json(value: &serde_json::Value) -> Result<Option<Job>> {
opt_field("died_at")?,
)))
}
/// Decode the substrate's stream-page JSON into the contract's
/// [`StreamEvent`] values (`from_row` — `topic` carries the contract's
/// `stream`).
fn stream_events_from_json_page(page: &str, stream: String) -> Result<Vec<StreamEvent>> {
let array: serde_json::Value = serde_json::from_str(page)
.map_err(|e| Error::Codec(format!("stream page must be json: {e}")))?;
let rows = array
.as_array()
.ok_or_else(|| Error::Codec("stream page must be a json array".to_string()))?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
out.push(StreamEvent::from_row(
serde_json::from_value(row["offset"].clone())
.map_err(|e| Error::Codec(format!("stream row offset must be i64: {e}")))?,
stream.clone(),
if row["key"].is_null() {
None
} else {
serde_json::from_value(row["key"].clone())
.map_err(|e| Error::Codec(format!("stream row key must be a string: {e}")))?
},
row["payload"]
.as_str()
.ok_or_else(|| Error::Codec("stream row payload must be a string".to_string()))?
.as_bytes()
.to_vec(),
serde_json::from_value(row["created_at"].clone())
.map_err(|e| Error::Codec(format!("stream row created_at must be i64: {e}")))?,
));
}
Ok(out)
}
+101 -10
View File
@@ -1,7 +1,7 @@
---
id: sqlite-engine-streams
name: SQLite engine — streams (`StreamHandle`, reads, subscribe, trim)
status: pending
status: completed
depends_on: [sqlite-engine-seam-tx]
scope: moderate
risk: medium
@@ -49,21 +49,21 @@ delta, mapped at this layer):
## Acceptance Criteria
- [ ] Full `StreamHandle` impl; publish/read/save/get/trim/subscribe
- [x] Full `StreamHandle` impl; publish/read/save/get/trim/subscribe
all work against a temp-file store
- [ ] Extent guard: `read_since`/`read_from_consumer` with
- [x] Extent guard: `read_since`/`read_from_consumer` with
`limit <= 0` return the empty `Vec` (never the whole stream) —
pinned per ADR-023 §2's backlog row
- [ ] `trim_to` exact-boundary (`<=`), negative horizon deletes
- [x] `trim_to` exact-boundary (`<=`), negative horizon deletes
nothing, surviving offsets never renumbered
- [ ] Monotone saves: regression saves are silent no-ops; direct and
- [x] Monotone saves: regression saves are silent no-ops; direct and
receiver forms compose without regression
- [ ] `subscribe` replays from the stored offset, then delivers
- [x] `subscribe` replays from the stored offset, then delivers
post-attach events wake-driven; watcher death closes it
(`None`, terminal); `recv()`'s `Err` carries `Database` only
- [ ] Key round-trips exactly (`None` vs `Some`); `stream` field
- [x] Key round-trips exactly (`None` vs `Some`); `stream` field
carries the contract name (not `topic`)
- [ ] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
- [x] `cargo test -p alkstore-sqlite`, clippy `-D warnings`, fmt clean
## References
@@ -76,8 +76,99 @@ delta, mapped at this layer):
## Notes
> To be filled by implementation agent
- **Stream module (stream.rs)**: `SqliteStreamHandle { name, writer,
readers, watcher, closed }` — the handle carries the store's
machinery arcs (writer slot, reader pool, watcher, closed flag).
Publishes (both forms) ride `with_writer` (the seam's short-lived
slot lease); reads ride a new `with_reader` bridge (pooled reader,
`spawn_blocking`, release on both arms); `trim_to`'s DELETE rides
`with_writer` inside the writer-slot lease, straight SQL over the
substrate (the substrate has no trim op — engine-side statement,
`deleted as i64` count).
- **Decode shared**: `stream_events_from_json_page` moved from
`tx.rs` into `stream.rs` as the one `pub(crate)` decode owner (both
the tx reads and the auto-commit reads call it; the `topic` →
`stream` name delta and the `Codec` posture guard live there once).
- **Extent guard placement**: `read_since`/`read_from_consumer` (and
the receiver's `read_since`) check `limit <= 0` at trait-impl entry,
before validation of *any* round trip into the substrate — the
empty-`Vec` is the whole result, `LIMIT -1` can never fire. The
closed-store check precedes the guard, so a closed store still fails
closed even for `limit <= 0` (the guard is about what an *open*
store's limit means).
- **Subscribe bridge**: a dedicated std thread (one per subscription,
the `listen` bridge shape) — (1) `watcher.subscribe()` is taken
*before* any read so a commit cannot fall between the attach read
and the wake wait; (2) initial cursor read of the consumer's *stored*
offset; (3) page-drain to the tail (`SUBSCRIBE_PAGE = 256`,
`offset ASC`); (4) park on the wake feed, re-drain on wake
(overtrigger-coalesced hints); exit on wake-feed disconnect (watcher
death / store close) or failed send. Read failures while open
surface one `Err(Database)` item and wait for the next wake (the
transient arm); attach-read failure yields one `Err` and closes.
A thread-spawn failure unsubscribes and fails `Database` (W-2
posture, no panics cross the seam).
- **Receiver position semantics**: `position` = the last-*yielded*
event's offset (0 before any yield ⇒ `save_offset()` before any
yield is a no-op, ADR-021 §5); advances only on delivered events —
the bridge's internal cursor (which advances on
*drained-from-storage* events) is separate, so a stalled consumer's
`save_offset()` never checkpoints unobserved events. Backpressure:
bounded tokio channel + `blocking_send` parks the bridge; delivery
resumes exactly where it stopped.
- **`EventReceiver::save_offset` is sync in the trait** — the
receiver's save takes the writer slot on the calling thread (the
slot is free between ops; a long open transaction parks this call —
the honest lease posture, surfaced in the doc comment). `offset()`
reports the in-memory position, not a db read.
- **`recv()`'s `Err` carries `Database` only (ADR-021 §5)**: a
`database_only` helper remaps any non-`Database` read error (e.g.
the decode-side `Codec` posture guard) into `Database` with the
original error preserved as the source. `try_recv` arms mirror the
notify receiver: idle `Ok(None)`, `Err(Closed)` after the channel
disconnected.
- **Closed-store postures**: every handle op (and the receiver's
`read_since`/`save_offset`) checks the store's closed flag first and
fails closed with `Database`; the reader pool's closed error
(`Database is closed` MISUSE shape) is recognized and remapped to
the same closed-store error. New `stream(name)` constructors on a
closed store fail `Database` (like `listen`). The bridge exits on
the closed reads — receiver sees `None`/`Err(Closed)` terminal.
- **Validation**: stream names via `validate_shared_name` at `store.
stream()` and (defensively) the handle's reads; consumers via
`validate_local_name` (non-empty only — reserved prefixes legal for
consumers, tested); empty-`Some` keys rejected `InvalidName` per the
tx path's rule.
- **Test infra**: `store/stream_tests.rs`, 15 tests over temp-dir
stores (two-store-open cross-connection probes where wake arrival
needs a second writer; raw-reader assertions only where the task
demands table-level honesty — trim count/durability read through
the handle, not raw SQL — the substrate's ops.rs already probes the
raw rows). Stability: suite ran 4× green.
## Summary
> To be filled on completion
Implemented the stream mechanism's SQLite arm: `SqliteStreamHandle`
(full `StreamHandle` trait — `publish`/`publish_with_key` auto-commit
through the writer slot with `encode_payload` at the seam and
`Codec` propagation, `read_since`/`read_from_consumer` reader-pool
reads with the ADR-023 §2 extent guard and ASC ordering, monotone
`save_offset`/`get_offset`, `trim_to`'s writer-lease DELETE with the
total boundary argument) and the durable `subscribe` receiver (a
std-thread bridge: attach-before-read, replay from the stored offset,
wake-driven re-reads off the watcher fanout, terminal close on
watcher death, position = last-yielded offset, sync receiver-form
`save_offset` through the same monotone op; `recv()`'s `Err` remapped
to `Database`-only). Shared the stream-page decoder between the tx
and auto-commit paths (`stream.rs` owns it, `tx.rs` imports). Wired
`Store::stream` in `store.rs` (entry-point validation + closed-store
fail-closed). 15 new tests (round-trip + stream-field + key, extent
guard, trim boundary/negative/no-renumbering, monotone saves across
forms, replay + wake-driven delivery, terminal close, tail-idle,
independent subscribers, receiver-drop unsubscribe, entry-point
validation incl. reserved-prefix-local consumers, closed-store fail
closures, restart durability, tx-compose with ghost verification,
receiver `read_since` anchoring + extent guard). Verified: `cargo
test` (workspace: 25 + 3 suite + 161 sqlite),
`cargo clippy --all-targets -- -D warnings`, `cargo fmt --check` all
clean; suite run 4× green.