//! stream subscribers against the real write path, with a live commit landing //! while a backfill batch is still open: the interleaving that once made //! both streams skip the backfill (hydrant-abq). use std::sync::Arc; use std::time::{Duration, Instant}; use jacquard_common::types::did::Did; use jacquard_common::types::string::Tid; use tokio::sync::mpsc; use crate::config::Config; use crate::db::Txn; use crate::db::types::{DbAction, DbRkey, DbTid}; use crate::state::AppState; use super::StreamOptions; const BACKFILLED: &str = "did:plc:aaaaaaaaaaaaaaaaaaaaaaaa"; const LIVE: &str = "did:plc:zzzzzzzzzzzzzzzzzzzzzzzz"; const COLLECTION: &str = "app.bsky.feed.post"; const REPO_SIZE: usize = 100; fn state() -> (tempfile::TempDir, Arc, StreamOptions) { let tmp = tempfile::tempdir().unwrap(); let config = Config { database_path: tmp.path().to_path_buf(), ..Default::default() }; let state = Arc::new(AppState::new(&config).unwrap()); (tmp, state, StreamOptions::from_config(&config)) } fn put(records: &mut crate::db::RecordTxn<'_, '_, '_>, n: usize) { let block = crate::car::CarBlock::from_body( serde_ipld_dagcbor::to_vec(&serde_json::json!({ "$type": COLLECTION, "n": n })).unwrap(), ); let rkey = DbRkey::Str(smol_str::format_smolstr!("r{n}")); records .put_record(COLLECTION, &rkey, &block, DbAction::Create) .unwrap(); } fn rev() -> DbTid { DbTid::from(&Tid::now_0()) } /// the backfill worker's write: a whole repo in one uncommitted batch. fn stage_backfill(state: &AppState) -> Txn<'_> { let did = Did::new_static(BACKFILLED).unwrap(); let mut txn = Txn::new(&state.db); let mut records = txn.backfill_records(state, &rev(), &did).unwrap(); (0..REPO_SIZE).for_each(|n| put(&mut records, n)); records.finish().unwrap(); txn } /// one live commit, the way the indexer shard writes it. fn commit_live(state: &AppState) { commit_live_records(state, REPO_SIZE..REPO_SIZE + 1); } /// one live commit writing a record for each of `records`. fn commit_live_records(state: &AppState, records: std::ops::Range) { let did = Did::new_static(LIVE).unwrap(); let mut txn = Txn::new(&state.db); let mut staged = txn.records(state, &rev(), &did).unwrap(); records.for_each(|n| put(&mut staged, n)); staged.finish().unwrap(); txn.commit().unwrap(); } /// dids of what a subscriber sent, collected until `until` holds or a /// deadline passes. struct Seen { rx: mpsc::Receiver, did_of: fn(T) -> String, dids: Vec, } impl Seen { fn until(&mut self, until: impl Fn(&[String]) -> bool) -> &mut Self { let deadline = Instant::now() + Duration::from_secs(5); while Instant::now() < deadline && !until(&self.dids) { match self.rx.try_recv() { Ok(item) => self.dids.push((self.did_of)(item)), Err(_) => std::thread::sleep(Duration::from_millis(5)), } } self } fn live(&mut self) -> &mut Self { self.until(|dids| dids.iter().any(|did| did == LIVE)) } /// wait a moment longer, for anything sent that shouldn't have been fn settle(&mut self) -> &mut Self { std::thread::sleep(Duration::from_millis(100)); while let Ok(item) = self.rx.try_recv() { self.dids.push((self.did_of)(item)); } self } fn counts(&self) -> (usize, usize) { let count = |did| self.dids.iter().filter(|seen| *seen == did).count(); (count(BACKFILLED), count(LIVE)) } } fn stream(state: &Arc, opts: StreamOptions, cursor: Option) -> Seen { let (tx, rx) = mpsc::channel(4096); let state = state.clone(); std::thread::spawn(move || { super::event_stream_thread::(state, tx, cursor, opts) }); Seen { rx, did_of: |item| item.unwrap().record.unwrap().did.as_str().to_owned(), dids: Vec::new(), } } type StreamItem = Result; /// wait until `n` live-only subscribers have subscribed, so they see what /// follows. fn wait_subscribed(n: usize, receivers: impl Fn() -> usize) { let deadline = Instant::now() + Duration::from_secs(5); while receivers() < n { assert!(Instant::now() < deadline, "subscriber never subscribed"); std::thread::sleep(Duration::from_millis(5)); } } #[test] fn stream_tail_gets_a_backfill_committed_after_a_live_commit() { let (_tmp, state, opts) = state(); let mut seen = stream(&state, opts, None); wait_subscribed(1, || state.db.stream.event_tx.receiver_count()); let backfill = stage_backfill(&state); commit_live(&state); seen.live(); backfill.commit().unwrap(); seen.until(|dids| dids.len() > REPO_SIZE).settle(); assert_eq!(seen.counts(), (REPO_SIZE, 1)); } #[test] fn stream_replay_gets_a_backfill_committed_after_a_live_commit() { let (_tmp, state, opts) = state(); let backfill = stage_backfill(&state); commit_live(&state); let mut seen = stream(&state, opts, Some(0)); seen.live(); backfill.commit().unwrap(); seen.until(|dids| dids.len() > REPO_SIZE).settle(); assert_eq!(seen.counts(), (REPO_SIZE, 1)); } #[cfg(feature = "jetstream")] mod jetstream { use super::*; use crate::control::{JetstreamFilter, JetstreamSubscriberOptions}; type Item = Result; fn subscribe( state: &Arc, opts: StreamOptions, cursor: Option, wanted_event_types: &[&str], ) -> Seen { subscribe_with_capacity(state, opts, 4096, cursor, wanted_event_types) } fn subscribe_with_capacity( state: &Arc, opts: StreamOptions, capacity: usize, cursor: Option, wanted_event_types: &[&str], ) -> Seen { let (tx, rx) = mpsc::channel(capacity); let state = state.clone(); let wanted: Vec<_> = wanted_event_types.iter().map(|t| t.to_string()).collect(); let filter = JetstreamFilter::new(JetstreamSubscriberOptions::parse(&[], &[], 0, &wanted).unwrap()); std::thread::spawn(move || { super::super::jetstream_stream_thread(state, tx, cursor, filter, opts) }); Seen { rx, did_of: |item| { let json: serde_json::Value = serde_json::from_slice(&item.unwrap()).unwrap(); json["did"].as_str().unwrap().to_owned() }, dids: Vec::new(), } } fn a_second_ago() -> i64 { chrono::Utc::now().timestamp_micros() - 1_000_000 } #[test] fn replay_gets_a_backfill_committed_after_a_live_commit() { let (_tmp, state, opts) = state(); let cursor = a_second_ago(); let backfill = stage_backfill(&state); commit_live(&state); let mut seen = subscribe(&state, opts, Some(cursor), &["live", "historical"]); seen.live(); backfill.commit().unwrap(); seen.until(|dids| dids.len() > REPO_SIZE).settle(); assert_eq!(seen.counts(), (REPO_SIZE, 1)); } #[test] fn tail_gets_backfills_only_when_asked() { let (_tmp, state, opts) = state(); let mut historical = subscribe(&state, opts, None, &["live", "historical"]); let mut default = subscribe(&state, opts, None, &[]); wait_subscribed(2, || state.db.jetstream.tx.receiver_count()); let backfill = stage_backfill(&state); commit_live(&state); backfill.commit().unwrap(); historical.until(|dids| dids.len() > REPO_SIZE).settle(); assert_eq!(historical.counts(), (REPO_SIZE, 1)); default.live().settle(); assert_eq!(default.counts(), (0, 1)); } #[test] fn a_lagging_tail_gets_backfills_only_when_asked() { // more live events than the broadcast channel holds, per commit const LAG: usize = 600; let (_tmp, state, opts) = state(); let opts = StreamOptions { pending_event_limit: 1, ..opts }; let mut historical = subscribe_with_capacity(&state, opts, 64, None, &["live", "historical"]); let mut default = subscribe_with_capacity(&state, opts, 64, None, &[]); wait_subscribed(2, || state.db.jetstream.tx.receiver_count()); // nobody reads, so both subscribers lag and catch up from the db commit_live_records(&state, REPO_SIZE..REPO_SIZE + LAG); stage_backfill(&state).commit().unwrap(); commit_live_records(&state, REPO_SIZE + LAG..REPO_SIZE + 2 * LAG); historical .until(|dids| dids.len() >= REPO_SIZE + 2 * LAG) .settle(); assert_eq!(historical.counts(), (REPO_SIZE, 2 * LAG)); default.until(|dids| dids.len() >= 2 * LAG).settle(); assert_eq!(default.counts(), (0, 2 * LAG)); } #[test] fn a_cursor_ahead_of_the_wall_clock_replays() { let (_tmp, state, opts) = state(); // committed time_us stays ahead of the wall clock after it steps back state.db.jetstream.clock.run_ahead(Duration::from_secs(60)); commit_live(&state); let seen_through = state.db.jetstream.head().unwrap().unwrap().time_us; commit_live_records(&state, REPO_SIZE + 1..REPO_SIZE + 2); let cursor = i64::try_from(seen_through + 1).unwrap(); let mut seen = subscribe(&state, opts, Some(cursor), &[]); seen.live().settle(); assert_eq!(seen.counts(), (0, 1)); } #[test] fn a_redaction_doesnt_reach_events_already_queued() { let (_tmp, state, opts) = state(); let (tx, mut rx) = mpsc::channel(2); let filter = JetstreamFilter::new(JetstreamSubscriberOptions::default()); let thread_state = state.clone(); std::thread::spawn(move || { super::super::jetstream_stream_thread(thread_state, tx, None, filter, opts) }); // live bodies are only kept inline while /stream has a subscriber too let _stream = stream(&state, opts, None); wait_subscribed(1, || state.db.jetstream.tx.receiver_count()); wait_subscribed(1, || state.db.stream.event_tx.receiver_count()); let taken = || wait_subscribed(1, || 1 - state.db.jetstream.tx.len().min(1)); // fill the output and leave the subscriber waiting to send the next, // so the one after sits unrendered in its backlog commit_live_records(&state, REPO_SIZE..REPO_SIZE + 2); taken(); commit_live_records(&state, REPO_SIZE + 2..REPO_SIZE + 3); taken(); let did = Did::new_static(LIVE).unwrap(); let rkey = DbRkey::Str(smol_str::format_smolstr!("r{}", REPO_SIZE + 2)); crate::db::redact_record_bodies( &state.db, &did, COLLECTION, &rkey, crate::types::DeleteBodyTarget::Head, ) .unwrap(); let commits: Vec = (0..3) .map(|_| { let item = rx.blocking_recv().unwrap().unwrap(); serde_json::from_slice::(&item).unwrap()["commit"].take() }) .collect(); // the third was queued before the redaction completed, so it goes out as committed assert!( commits.iter().all(|commit| commit.get("record").is_some()), "{commits:?}" ); } #[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(); let cursor = a_second_ago(); stage_backfill(&state).commit().unwrap(); commit_live(&state); let mut seen = subscribe(&state, opts, Some(cursor), &[]); seen.until(|dids| dids.len() > REPO_SIZE).settle(); assert_eq!(seen.counts(), (REPO_SIZE, 1)); let mut live_only = subscribe(&state, opts, Some(cursor), &["live"]); live_only.live().settle(); assert_eq!(live_only.counts(), (0, 1)); } }