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.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509//! events a transaction stages, written and broadcast when it commits.//!//! stream rows are staged here without a position. [`Outbox::commit`] gives//! them one inside the [`Sequencer`](super::sequencer::Sequencer)'s critical//! section, which is the only place stream positions come from.
use fjall::OwnedWriteBatch;use miette::{IntoDiagnostic, Result};
use super::Db;
#[cfg(feature = "jetstream")]mod jetstream;#[cfg(feature = "relay")]mod relay;#[cfg(feature = "indexer_stream")]mod stream;
#[cfg(feature = "jetstream")]pub(crate) use jetstream::JetstreamOutbox;#[cfg(feature = "relay")]pub(crate) use relay::{RelayFrameRef, RelayOutbox};#[cfg(feature = "indexer_stream")]pub(crate) use stream::{StreamEventRef, StreamOutbox};
#[derive(Default)]pub(crate) struct Outbox { #[cfg(feature = "indexer_stream")] pub(crate) stream: StreamOutbox, #[cfg(feature = "relay")] pub(crate) relay: RelayOutbox, #[cfg(feature = "jetstream")] pub(crate) jetstream: JetstreamOutbox,}
/// the positions a commit gave its events.#[derive(Default)]pub(crate) struct Assigned { #[cfg(feature = "indexer_stream")] first_stream_id: Option<u64>,}
#[cfg(feature = "indexer_stream")]impl Assigned { pub(crate) fn stream_id(&self, event: StreamEventRef) -> Option<u64> { self.first_stream_id .map(|first| first + event.index() as u64) }}
#[cfg(not(any(feature = "indexer_stream", feature = "relay")))]impl Outbox { pub(super) fn commit( self, _db: &Db, batch: OwnedWriteBatch, committed: impl FnOnce(), ) -> Result<Assigned> { batch.commit().into_diagnostic()?; committed(); Ok(Assigned::default()) }}
#[cfg(any(feature = "indexer_stream", feature = "relay"))]impl Outbox { /// commit `batch` together with this outbox's rows, run `committed`, then /// broadcast them. pub(super) fn commit( self, db: &Db, batch: OwnedWriteBatch, committed: impl FnOnce(), ) -> Result<Assigned> { if self.is_empty() { batch.commit().into_diagnostic()?; committed(); return Ok(Assigned::default()); } db.sequencer.commit( batch, |assign, batch| self.stage(db, assign, batch), |staged| { committed(); staged.publish(db) }, ) }
fn is_empty(&self) -> bool { #[cfg(feature = "indexer_stream")] if !self.stream.is_empty() { return false; } #[cfg(feature = "relay")] if !self.relay.is_empty() { return false; } #[cfg(feature = "jetstream")] if !self.jetstream.is_empty() { return false; } true }
fn stage<'a>( self, db: &'a Db, assign: &mut super::sequencer::Assign<'a>, batch: &mut OwnedWriteBatch, ) -> Result<Staged> { #[cfg(feature = "indexer_stream")] let stream = self.stream.stage(&db.stream, assign, batch)?; #[cfg(feature = "relay")] let relay = self.relay.stage(&db.relay, assign, batch)?; // jetstream rows refer to the upstream rows staged just above #[cfg(all(feature = "jetstream", feature = "indexer_stream"))] let jetstream = self .jetstream .stage(&db.jetstream, assign, batch, &stream)?; #[cfg(all(feature = "jetstream", feature = "relay"))] let jetstream = self.jetstream.stage(&db.jetstream, assign, batch, &relay)?; Ok(Staged { #[cfg(feature = "indexer_stream")] stream, #[cfg(feature = "relay")] relay, #[cfg(feature = "jetstream")] jetstream, }) }}
#[cfg(any(feature = "indexer_stream", feature = "relay"))]struct Staged { #[cfg(feature = "indexer_stream")] stream: stream::StagedStream, #[cfg(feature = "relay")] relay: relay::StagedRelay, #[cfg(feature = "jetstream")] jetstream: jetstream::StagedJetstream,}
#[cfg(any(feature = "indexer_stream", feature = "relay"))]impl Staged { fn publish(self, db: &Db) -> Assigned { #[cfg(feature = "relay")] self.relay.publish(&db.relay); #[cfg(feature = "jetstream")] self.jetstream.publish(&db.jetstream); Assigned { #[cfg(feature = "indexer_stream")] first_stream_id: self.stream.publish(&db.stream), } }}
#[cfg(all(test, feature = "indexer_stream"))]mod tests { use std::sync::Arc;
use jacquard_common::types::did::Did; use jacquard_common::types::string::Tid;
use crate::config::Config; use crate::db::Txn; use crate::db::types::{DbAction, DbRkey, DbTid}; use crate::state::AppState; use crate::types::BroadcastEvent;
const BACKFILLED: &str = "did:plc:aaaaaaaaaaaaaaaaaaaaaaaa"; const LIVE: &str = "did:plc:zzzzzzzzzzzzzzzzzzzzzzzz"; const COLLECTION: &str = "app.bsky.feed.post";
fn state() -> (tempfile::TempDir, Arc<AppState>) { let tmp = tempfile::tempdir().unwrap(); let state = AppState::new(&Config { database_path: tmp.path().to_path_buf(), ..Default::default() }) .unwrap(); (tmp, Arc::new(state)) }
fn put(records: &mut crate::db::txn::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(); }
/// a backfill batch is staged, a live commit lands, then the backfill /// commits: the order ingest produces when both run at once. fn interleave(state: &AppState, backfilled: usize) { let rev = DbTid::from(&Tid::now_0()); let did = Did::new_static(BACKFILLED).unwrap(); let mut backfill = Txn::new(&state.db); let mut records = backfill.backfill_records(state, &rev, &did).unwrap(); (0..backfilled).for_each(|n| put(&mut records, n)); records.finish().unwrap();
let did = Did::new_static(LIVE).unwrap(); let mut live = Txn::new(&state.db); let mut records = live.records(state, &rev, &did).unwrap(); put(&mut records, backfilled); records.finish().unwrap(); live.commit().unwrap();
backfill.commit().unwrap(); }
fn event_did(state: &AppState, id: u64) -> String { let row = state .db .stream .event_range(crate::db::keys::event_key(id)..=crate::db::keys::event_key(id)) .next() .unwrap() .value() .unwrap(); let event: crate::types::StoredEvent<'_> = rmp_serde::from_slice(&row).unwrap(); event.did.to_did().as_str().to_owned() }
#[test] fn stream_ids_follow_commit_order() { let (_tmp, state) = state(); let mut rx = state.db.stream.event_tx.subscribe(); interleave(&state, 3);
// the live commit took the first id even though it staged last let first = state.db.stream.ids.head().unwrap() - 3; assert_eq!(event_did(&state, first), LIVE); for id in first + 1..=first + 3 { assert_eq!(event_did(&state, id), BACKFILLED); }
let Ok(BroadcastEvent::LiveRecord(live)) = rx.try_recv() else { panic!("expected the live record first"); }; assert_eq!(live.id, first); let Ok(BroadcastEvent::Persisted(head)) = rx.try_recv() else { panic!("expected the backfill's persisted marker"); }; assert_eq!(head, first + 3); assert!(rx.try_recv().is_err()); }
#[test] fn the_committed_hook_runs_before_the_broadcast() { let (_tmp, state) = state(); let rx = state.db.stream.event_tx.subscribe(); let mut outbox = super::Outbox::default(); outbox.stream.push_identity(crate::types::IdentityEvt { did: Did::new_static(LIVE).unwrap(), handle: None, });
let mut ran = false; outbox .commit(&state.db, state.db.inner.batch(), || { assert!(rx.is_empty(), "broadcast before the hook ran"); ran = true; }) .unwrap(); assert!(ran); assert!(!rx.is_empty()); }
#[test] fn stream_ids_are_not_reused_after_a_restart() { let tmp = tempfile::tempdir().unwrap(); let config = Config { database_path: tmp.path().to_path_buf(), ..Default::default() }; let head = { let db = crate::db::Db::open(&config).unwrap(); // identity events take ids but leave no rows let mut txn = Txn::new(&db); txn.outbox.stream.push_identity(crate::types::IdentityEvt { did: Did::new_static(LIVE).unwrap(), handle: None, }); txn.commit().unwrap(); db.persist().unwrap(); db.stream.ids.head() }; let db = crate::db::Db::open(&config).unwrap(); let head = head.unwrap(); // past every id taken, and no further than one reserved block let next = db.stream.ids.next(); assert!(head < next && next <= head + 1 + crate::db::sequencer::SEQUENCE_BLOCK); }
#[cfg(feature = "jetstream")] #[test] fn jetstream_keys_follow_commit_order() { use crate::types::StoredJetstreamEvent;
let (_tmp, state) = state(); interleave(&state, 3);
let rows: Vec<_> = state .db .jetstream .events .iter() .map(|guard| { let (key, value) = guard.into_inner().unwrap(); let position = crate::db::keys::parse_jetstream_event_key(&key).unwrap(); let event: StoredJetstreamEvent<'_> = rmp_serde::from_slice(&value).unwrap(); let StoredJetstreamEvent::Commit { event_id, live, .. } = event else { panic!("expected commit rows"); }; (position, event_id, live) }) .collect();
assert_eq!(rows.len(), 4); assert!(rows.windows(2).all(|pair| { let ((time_a, id_a), _, _) = pair[0]; let ((time_b, id_b), _, _) = pair[1]; time_a < time_b && id_a + 1 == id_b })); let (_, live_event, live) = rows[0]; assert!(live); assert_eq!(event_did(&state, live_event), LIVE); for &(_, event_id, live) in &rows[1..] { assert!(!live); assert_eq!(event_did(&state, event_id), BACKFILLED); } }}
#[cfg(all(test, feature = "relay"))]mod relay_tests { use jacquard_common::types::string::Did;
use crate::config::Config; use crate::db::{Db, Txn}; use crate::ingest::stream::{Datetime, Identity, SubscribeReposMessage, decode_frame}; use crate::types::RelayBroadcast;
const FIRST: &str = "did:plc:aaaaaaaaaaaaaaaaaaaaaaaa"; const SECOND: &str = "did:plc:zzzzzzzzzzzzzzzzzzzzzzzz";
fn identity(did: &'static str) -> Identity { Identity { did: Did::new_static(did).unwrap(), handle: None, seq: 0, time: Datetime::try_from("2026-08-10T12:34:56Z".to_owned()).unwrap(), } }
fn stage<'a>(db: &'a Db, did: &'static str) -> Txn<'a> { let mut txn = Txn::new(db); let _frame = txn.outbox.relay.push_identity(identity(did)); #[cfg(feature = "jetstream")] txn.outbox.jetstream.push( crate::types::StoredJetstreamEvent::RelayIdentity { did: crate::db::types::TrimmedDid::from(&Did::new_static(did).unwrap()) .into_static(), relay_seq: _frame, }, None, ); txn }
#[test] fn relay_seqs_continue_after_their_rows_are_pruned() { let tmp = tempfile::tempdir().unwrap(); let config = Config { database_path: tmp.path().to_path_buf(), ..Default::default() }; let head = { let db = Db::open(&config).unwrap(); stage(&db, FIRST).commit().unwrap(); stage(&db, SECOND).commit().unwrap(); // ttl pruning can remove every row let mut batch = db.inner.batch(); for guard in db.relay.events.iter() { batch.remove(&db.relay.events, guard.key().unwrap()); } batch.commit().unwrap(); db.persist().unwrap(); db.relay.seqs.head() }; let db = Db::open(&config).unwrap(); let head = head.unwrap(); // past every seq taken, and no further than one reserved block let next = db.relay.seqs.next(); assert!(head < next && next <= head + 1 + crate::db::sequencer::SEQUENCE_BLOCK); }
#[test] fn a_frame_that_fails_to_encode_is_dropped_alone() { let tmp = tempfile::tempdir().unwrap(); let db = Db::open(&Config { database_path: tmp.path().to_path_buf(), ..Default::default() }) .unwrap(); let mut rx = db.relay.broadcast_tx.subscribe();
let mut txn = stage(&db, FIRST); txn.outbox.relay.push_unencodable(); let _ = txn.outbox.relay.push_identity(identity(SECOND)); txn.batch.insert( &db.cursors, b"test_cursor".as_slice(), b"written".as_slice(), ); txn.commit().unwrap();
// the rest of the transaction lands, and the lost frame's seq stays a hole assert!(db.cursors.get(b"test_cursor").unwrap().is_some()); let seqs: Vec<u64> = db .relay .events .iter() .map(|guard| u64::from_be_bytes(guard.key().unwrap().as_ref().try_into().unwrap())) .collect(); let head = db.relay.seqs.head().unwrap(); assert_eq!(seqs, [head - 2, head]);
let broadcast: Vec<_> = std::iter::from_fn(|| rx.try_recv().ok()).collect(); assert!(matches!( broadcast.as_slice(), [ RelayBroadcast::Ephemeral(first, _), RelayBroadcast::Persisted(hole), RelayBroadcast::Ephemeral(second, _), ] if *first == head - 2 && *hole == head - 1 && *second == head )); }
#[test] fn relay_seqs_follow_commit_order() { let tmp = tempfile::tempdir().unwrap(); let db = Db::open(&Config { database_path: tmp.path().to_path_buf(), ..Default::default() }) .unwrap(); let mut rx = db.relay.broadcast_tx.subscribe();
let first = stage(&db, FIRST); stage(&db, SECOND).commit().unwrap(); first.commit().unwrap();
let frames: Vec<_> = db .relay .events .iter() .map(|guard| { let (key, frame) = guard.into_inner().unwrap(); let seq = u64::from_be_bytes(key.as_ref().try_into().unwrap()); let SubscribeReposMessage::Identity(identity) = decode_frame(&frame).unwrap() else { panic!("expected identity frames"); }; assert_eq!(identity.seq as u64, seq, "a frame carries its own seq"); (seq, identity.did.as_str().to_owned()) }) .collect(); let head = db.relay.seqs.head().unwrap(); assert_eq!( frames, [(head - 1, SECOND.to_owned()), (head, FIRST.to_owned())] );
let broadcast: Vec<_> = std::iter::from_fn(|| rx.try_recv().ok()) .map(|event| match event { RelayBroadcast::Ephemeral(seq, _) => seq, RelayBroadcast::Persisted(seq) => panic!("unexpected marker at {seq}"), }) .collect(); assert_eq!(broadcast, [head - 1, head]);
#[cfg(feature = "jetstream")] for guard in db.jetstream.events.iter() { let (_, value) = guard.into_inner().unwrap(); let event: crate::types::StoredJetstreamEvent<'_> = rmp_serde::from_slice(&value).unwrap(); let crate::types::StoredJetstreamEvent::RelayIdentity { did, relay_seq } = event else { panic!("expected relay identity rows"); }; let frame = db .relay .events .get(crate::db::keys::relay_event_key(relay_seq)) .unwrap() .unwrap(); let SubscribeReposMessage::Identity(identity) = decode_frame(&frame).unwrap() else { panic!("jetstream row points at a non-identity frame"); }; assert_eq!(identity.did.as_str(), did.to_did().as_str()); } }}