very fast at protocol indexer with flexible filtering, xrpc queries, cursor-backed event stream, and more, built on fjall
rust fjall at-protocol atproto indexer
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396//! 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<AppState>, 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<usize>) { 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<T> { rx: mpsc::Receiver<T>, did_of: fn(T) -> String, dids: Vec<String>,}
impl<T> Seen<T> { 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<AppState>, opts: StreamOptions, cursor: Option<u64>) -> Seen<StreamItem> { let (tx, rx) = mpsc::channel(4096); let state = state.clone(); std::thread::spawn(move || { super::event_stream_thread::<super::Decoded>(state, tx, cursor, opts) }); Seen { rx, did_of: |item| item.unwrap().record.unwrap().did.as_str().to_owned(), dids: Vec::new(), }}
type StreamItem = Result<crate::control::Event, crate::control::StreamError>;
/// 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<bytes::Bytes, crate::control::JetstreamStreamError>;
fn subscribe( state: &Arc<AppState>, opts: StreamOptions, cursor: Option<i64>, wanted_event_types: &[&str], ) -> Seen<Item> { subscribe_with_capacity(state, opts, 4096, cursor, wanted_event_types) }
fn subscribe_with_capacity( state: &Arc<AppState>, opts: StreamOptions, capacity: usize, cursor: Option<i64>, wanted_event_types: &[&str], ) -> Seen<Item> { 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<serde_json::Value> = (0..3) .map(|_| { let item = rx.blocking_recv().unwrap().unwrap(); serde_json::from_slice::<serde_json::Value>(&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)); }}