From 2d8dbd1cac4acdaff4da5bec06c7a35a06da98fa Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Tue, 6 Oct 2026 00:23:59 +0300 Subject: [PATCH] [tests] pin the stream engine's starts and replay chunks the representation audit lists every crate-visible function in src/control/stream as a boundary, and the shared engine landed with none of them inventoried. Start::live and Start::replay had no test of their own: nothing checked that a live start skips rows through its head, that a replay resumes after its cursor, or that a cursor already at the head reads nothing. ReplayChunk::read, which turns stored keys into positions for every stream, had no direct test of its scan bound, its skipping of keys that aren't positions, or how last_seen follows filtered rows. the jetstream outbox, JetstreamEphemeral::to_json and stored_event_to_bytes also had nothing checking that the json a live subscriber gets is the json a replay renders for the same commit. live_json_is_what_a_replay_renders compares the two and pins the commit shape, record body and cid. it needs the jetstream feature, like the rest of that module. --- src/control/stream/engine.rs | 47 ++++++++++++++++++++++ src/control/stream/interleaving.rs | 52 ++++++++++++++++++++++++ src/control/stream/types.rs | 64 ++++++++++++++++++++++++++++++ 3 files changed, 163 insertions(+) diff --git a/src/control/stream/engine.rs b/src/control/stream/engine.rs index 3960241..c5d64fd 100644 --- a/src/control/stream/engine.rs +++ b/src/control/stream/engine.rs @@ -606,6 +606,53 @@ mod tests { finish(&log, handle, out); } + #[test] + fn a_live_start_skips_everything_through_its_head() { + let (log, rx) = Log::new(16); + // already broadcast, so they wait in the subscriber's channel + log.backfill(3); + log.live(); + let (handle, mut out) = spawn(&log, rx, Start::live(log.head()), 16, opts(), |_, _| { + panic!("rows through the head need no db read"); + }); + + let live = log.live(); + assert_eq!(take(&mut out, 1), [live]); + finish(&log, handle, out); + } + + #[test] + fn a_replay_start_reads_only_past_its_cursor() { + let (log, rx) = Log::new(16); + let rows = log.backfill(4); + let (handle, mut out) = spawn( + &log, + rx, + Start::replay(Some(rows[1]), log.head()), + 16, + opts(), + nothing_on_read, + ); + assert_eq!(take(&mut out, 2), rows[2..]); + let live = log.live(); + assert_eq!(take(&mut out, 1), [live]); + finish(&log, handle, out); + + let (log, rx) = Log::new(16); + log.backfill(2); + let (handle, mut out) = spawn( + &log, + rx, + Start::replay(log.head(), log.head()), + 16, + opts(), + |_, _| panic!("a cursor at the head has nothing to replay"), + ); + let live = log.live(); + assert_eq!(take(&mut out, 1), [live]); + finish(&log, handle, out); + } + #[test] fn a_marker_reads_its_rows_from_the_db() { let (log, rx) = Log::new(16); diff --git a/src/control/stream/interleaving.rs b/src/control/stream/interleaving.rs index d625f40..1259b71 100644 --- a/src/control/stream/interleaving.rs +++ b/src/control/stream/interleaving.rs @@ -325,6 +325,58 @@ mod jetstream { assert!(commits[2].get("cid").is_some()); } + #[test] + fn live_json_is_what_a_replay_renders() { + let (_tmp, state, opts) = state(); + // live bodies are only kept inline while /stream has a subscriber too + let _stream = stream(&state, opts, None); + wait_subscribed(1, || state.db.stream.event_tx.receiver_count()); + let mut broadcasts = state.db.jetstream.tx.subscribe(); + let cursor = a_second_ago(); + commit_live(&state); + + let crate::types::JetstreamBroadcast::Live(live) = broadcasts.blocking_recv().unwrap() + else { + panic!("a live commit is broadcast live"); + }; + let json = live.json.expect("built while a subscriber listened"); + let live_json: serde_json::Value = serde_json::from_slice(&json).unwrap(); + + let (tx, mut rx) = mpsc::channel(4); + let filter = JetstreamFilter::new(JetstreamSubscriberOptions::default()); + let thread_state = state.clone(); + std::thread::spawn(move || { + super::super::jetstream_stream_thread(thread_state, tx, Some(cursor), filter, opts) + }); + let replayed = rx.blocking_recv().unwrap().unwrap(); + let replayed: serde_json::Value = serde_json::from_slice(&replayed).unwrap(); + assert_eq!(live_json, replayed); + + let body = serde_json::json!({ "$type": COLLECTION, "n": REPO_SIZE }); + let cid = crate::car::CarBlock::from_body(serde_ipld_dagcbor::to_vec(&body).unwrap()) + .cid() + .to_string(); + let rev = replayed["commit"]["rev"].as_str().unwrap(); + assert!(Tid::new(rev).is_ok(), "{rev}"); + assert_eq!( + replayed, + serde_json::json!({ + "did": LIVE, + "time_us": live.position.time_us, + "kind": "commit", + "commit": { + "rev": rev, + "operation": "create", + "collection": COLLECTION, + "rkey": format!("r{REPO_SIZE}"), + "record": body, + "cid": cid, + "live": true, + }, + }) + ); + } + #[test] fn replay_sends_backfills_by_default() { let (_tmp, state, opts) = state(); diff --git a/src/control/stream/types.rs b/src/control/stream/types.rs index 0f2703b..11bb33b 100644 --- a/src/control/stream/types.rs +++ b/src/control/stream/types.rs @@ -177,3 +177,67 @@ impl ReplayChunk { Ok(chunk) } } + +#[cfg(test)] +mod tests { + use super::*; + + /// keys 1 through 10 as big-endian positions, values naming them, plus one + /// key that isn't a position between 3 and 4 + fn rows() -> (tempfile::TempDir, fjall::Database, fjall::Keyspace) { + let tmp = tempfile::tempdir().unwrap(); + let db = fjall::Database::builder(tmp.path()).open().unwrap(); + let ks = db + .keyspace("rows", fjall::KeyspaceCreateOptions::default) + .unwrap(); + for n in 1..=10_u64 { + ks.insert(n.to_be_bytes(), format!("v{n}")).unwrap(); + } + ks.insert([0, 0, 0, 0, 0, 0, 0, 3, 0], "junk").unwrap(); + (tmp, db, ks) + } + + fn position(key: &[u8]) -> Option { + <[u8; 8]>::try_from(key).ok().map(u64::from_be_bytes) + } + + fn evens(position: u64, value: &[u8]) -> Option<(u64, String)> { + position + .is_multiple_of(2) + .then(|| (position, String::from_utf8(value.to_vec()).unwrap())) + } + + #[test] + fn a_chunk_stops_at_its_limit_and_remembers_filtered_rows() { + let (_tmp, _db, ks) = rows(); + + let chunk = ReplayChunk::read(ks.iter(), None, 2, position, evens).unwrap(); + assert_eq!(chunk.events, [(2, "v2".into()), (4, "v4".into())]); + assert_eq!(chunk.last_seen, Some(4)); + assert!(!chunk.exhausted); + + // four rows scanned for one wanted output, the junk key among them + let chunk = ReplayChunk::read(ks.iter(), None, 1, position, |_, _| None::<()>).unwrap(); + assert!(chunk.events.is_empty()); + assert_eq!(chunk.last_seen, Some(3)); + assert!(!chunk.exhausted); + + let chunk = ReplayChunk::read(ks.range(9_u64.to_be_bytes()..), Some(8), 4, position, evens) + .unwrap(); + assert_eq!(chunk.events, [(10, "v10".into())]); + assert_eq!(chunk.last_seen, Some(10)); + assert!(chunk.exhausted); + + let chunk = ReplayChunk::read( + ks.range(11_u64.to_be_bytes()..), + Some(10), + 4, + position, + evens, + ) + .unwrap(); + assert!(chunk.events.is_empty()); + assert_eq!(chunk.last_seen, Some(10)); + assert!(chunk.exhausted); + } +} -- 2.51.2