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.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434//! stream positions, assigned while committing.//!//! every event stream position (`/stream` ids, `/subscribe` keys, relay//! seqs) is handed out inside [`Sequencer::commit`], which holds one lock across assigning the//! positions, committing the batch that writes their rows, and broadcasting//! the events. so by construction://!//! - position order is commit order. every position at or below a sequence's//! head is committed, so a missing row below the head is a permanent hole,//! never a write still in flight.//! - a broadcast is only sent once its row is readable.//! - broadcasts arrive in position order, across every writer.//! - positions never repeat, even across restarts: the same batch writes//! down where each counter stopped, for positions that leave no row or//! whose rows get pruned.
use std::sync::atomic::{AtomicU64, Ordering};use std::sync::{Mutex, PoisonError};
use fjall::OwnedWriteBatch;use miette::{IntoDiagnostic, Result, WrapErr};
use super::schema::{Cursors, Ks};
/// the last position a counter hands out. every stream position fits in an/// i64, the width relay seqs and jetstream `time_us` have on the wire.const LAST_POSITION: u64 = i64::MAX as u64;
/// the one place stream positions are handed out.pub(crate) struct Sequencer { lock: Mutex<()>, /// where each counter's next position is written down marks: Ks<Cursors>,}
impl Sequencer { pub(crate) fn new(marks: Ks<Cursors>) -> Self { Self { lock: Mutex::new(()), marks, } }
/// assign positions and stage their rows, commit, then publish, all under /// one lock. positions only advance if the batch commits. pub(crate) fn commit<'a, T, R>( &self, mut batch: OwnedWriteBatch, stage: impl FnOnce(&mut Assign<'a>, &mut OwnedWriteBatch) -> Result<T>, publish: impl FnOnce(T) -> R, ) -> Result<R> { let _guard = self.lock.lock().unwrap_or_else(PoisonError::into_inner); let mut assign = Assign { pending: Vec::new(), }; let staged = stage(&mut assign, &mut batch)?; assign.mark(&self.marks, &mut batch); batch.commit().into_diagnostic()?; assign.advance(); Ok(publish(staged)) }}
/// positions taken by one sequenced commit. only [`Sequencer::commit`] can/// make one, so positions can't be taken anywhere else.pub(crate) struct Assign<'a> { /// the value each touched counter moves to once the batch commits pending: Vec<(&'a Counter, u64)>,}
impl<'a> Assign<'a> { /// take `counter`'s next value. the first take in a commit begins at /// `start(committed next)` and later ones count up from there. fn take(&mut self, counter: &'a Counter, start: impl FnOnce(u64) -> u64) -> Result<u64> { let slot = match self .pending .iter() .position(|(pending, _)| std::ptr::eq(*pending, counter)) { Some(slot) => slot, None => { let next = start(counter.next.load(Ordering::Acquire)); self.pending.push((counter, next)); self.pending.len() - 1 } }; let taken = self.pending[slot].1; miette::ensure!(taken <= LAST_POSITION, "stream positions are exhausted"); self.pending[slot].1 = taken + 1; Ok(taken) }
/// write down where each touched counter moves to, in the same batch fn mark(&self, marks: &Ks<Cursors>, batch: &mut OwnedWriteBatch) { for (counter, next) in &self.pending { batch.insert(marks, counter.mark, next.to_be_bytes()); } }
fn advance(self) { for (counter, next) in self.pending { counter.next.store(next, Ordering::Release); } }}
/// a position counter. every commit that moves it also writes its next/// position under `mark`, so reopening the database resumes past positions/// that left no row, or whose rows were pruned.struct Counter { /// the first position not committed yet next: AtomicU64, mark: &'static [u8],}
impl Counter { const fn new(mark: &'static [u8]) -> Self { Self { next: AtomicU64::new(0), mark, } }
/// continue at `next`, or at the mark if a commit got further. fn resume_at(&self, marks: &Ks<Cursors>, next: u64) -> Result<()> { let marked = marks .get(self.mark) .into_diagnostic()? .map(|mark| { <[u8; 8]>::try_from(mark.as_ref()) .map(u64::from_be_bytes) .into_diagnostic() .wrap_err("a stream position mark must be 8 bytes") }) .transpose()?; let next = next.max(marked.unwrap_or(0)); // only a corrupted or hand-edited database gets here miette::ensure!( next <= LAST_POSITION, "the database's next stream position {next} leaves no room for more" ); self.next.store(next, Ordering::Release); Ok(()) }}
/// a counter whose values are only taken inside a sequenced commit.pub(crate) struct Sequence(Counter);
impl Sequence { pub(crate) const fn new(mark: &'static [u8]) -> Self { Self(Counter::new(mark)) }
/// the last committed position, if any. pub(crate) fn head(&self) -> Option<u64> { self.next().checked_sub(1) }
/// the position the next commit will take first. pub(crate) fn next(&self) -> u64 { self.0.next.load(Ordering::Acquire) }
/// continue numbering after `last`, the last stored position, or at 1 /// for none, when opening the database. pub(in crate::db) fn resume_after(&self, marks: &Ks<Cursors>, last: Option<u64>) -> Result<()> { self.0.resume_at(marks, next_after(last).max(1)) }
#[cfg(test)] pub(in crate::db) fn resume_at(&self, next: u64) { self.0.next.store(next, Ordering::Release); }
pub(crate) fn take<'a>(&'a self, assign: &mut Assign<'a>) -> Result<u64> { assign.take(&self.0, |next| next) }}
/// microsecond timestamps that strictly increase in commit order, even when/// the wall clock stalls or steps back. a commit's timestamps are consecutive,/// starting at the wall time or just after the previous commit's last,/// whichever is later, so the wall clock is read once per commit.#[cfg(feature = "jetstream")]pub(crate) struct Clock(Counter);
#[cfg(feature = "jetstream")]impl Clock { pub(crate) const fn new(mark: &'static [u8]) -> Self { Self(Counter::new(mark)) }
/// continue after `last`, the last stored timestamp, when opening the /// database. pub(in crate::db) fn resume_after(&self, marks: &Ks<Cursors>, last: Option<u64>) -> Result<()> { self.0.resume_at(marks, next_after(last)) }
/// start the next commit `by` ahead of the wall clock, as if the wall /// clock had stepped back after it. #[cfg(all(test, feature = "indexer_stream"))] pub(crate) fn run_ahead(&self, by: std::time::Duration) { let now = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap(); let ahead = now + u64::try_from(by.as_micros()).unwrap(); self.0.next.store(ahead, Ordering::Release); }
pub(crate) fn tick<'a>(&'a self, assign: &mut Assign<'a>) -> Result<u64> { assign.take(&self.0, |next| { let now = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap_or(0); next.max(now) }) }}
/// the position after `last`, saturating so a stored `u64::MAX` fails/// resuming instead of wrapping.fn next_after(last: Option<u64>) -> u64 { last.map_or(0, |last| last.saturating_add(1))}
/// send `events` in position order. a run of rows that aren't broadcast is/// announced by one `marker` for its last position, sent before the next/// broadcast, so a subscriber reads them from the db before moving past them.pub(crate) fn publish_in_order<B, P>( tx: &tokio::sync::broadcast::Sender<B>, events: impl IntoIterator<Item = (P, Option<B>)>, marker: impl Fn(P) -> B,) { let mut unannounced = None; for (position, event) in events { let Some(event) = event else { unannounced = Some(position); continue; }; if let Some(position) = unannounced.take() { let _ = tx.send(marker(position)); } let _ = tx.send(event); } if let Some(position) = unannounced { let _ = tx.send(marker(position)); }}
#[cfg(test)]mod tests { use super::*; use crate::config::Config;
fn db() -> (tempfile::TempDir, crate::db::Db) { let tmp = tempfile::tempdir().unwrap(); let db = crate::db::Db::open(&Config { database_path: tmp.path().to_path_buf(), ..Default::default() }) .unwrap(); (tmp, db) }
#[test] fn positions_advance_only_when_the_batch_commits() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let seq = Sequence::new(b"test_seq"); seq.resume_at(5);
let failed = sequencer.commit( db.inner.batch(), |assign, _| { assert_eq!((seq.take(assign)?, seq.take(assign)?), (5, 6)); Err::<(), _>(miette::miette!("stage failed")) }, |_| {}, ); assert!(failed.is_err()); assert_eq!(seq.next(), 5);
let taken = sequencer .commit( db.inner.batch(), |assign, _| { let taken = [seq.take(assign)?, seq.take(assign)?]; // readers never see positions that haven't committed assert_eq!(seq.head(), Some(4)); Ok(taken) }, |taken| taken, ) .unwrap(); assert_eq!(taken, [5, 6]); assert_eq!(seq.head(), Some(6)); }
#[test] fn positions_stop_at_the_last_one() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let seq = Sequence::new(b"test_seq"); seq.resume_at(LAST_POSITION);
let failed = sequencer.commit( db.inner.batch(), |assign, _| { assert_eq!(seq.take(assign)?, LAST_POSITION); seq.take(assign) }, |_| {}, ); assert!(failed.is_err()); assert_eq!(seq.next(), LAST_POSITION); }
#[test] fn a_stored_position_needs_room_after_it() { let (_tmp, db) = db(); let seq = Sequence::new(b"test_seq"); seq.resume_after(&db.cursors, None).unwrap(); assert_eq!(seq.next(), 1); seq.resume_after(&db.cursors, Some(LAST_POSITION - 1)) .unwrap(); assert_eq!(seq.next(), LAST_POSITION); assert!(seq.resume_after(&db.cursors, Some(LAST_POSITION)).is_err()); assert!(seq.resume_after(&db.cursors, Some(u64::MAX)).is_err()); }
#[test] fn a_counter_resumes_past_what_its_commits_took() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let seq = Sequence::new(b"test_seq"); seq.resume_after(&db.cursors, None).unwrap(); sequencer .commit( db.inner.batch(), |assign, _| (0..3).try_for_each(|_| seq.take(assign).map(drop)), |_| {}, ) .unwrap();
// none of those positions has a row to resume after let reopened = Sequence::new(b"test_seq"); reopened.resume_after(&db.cursors, None).unwrap(); assert_eq!(reopened.next(), 4); // rows further along still win reopened.resume_after(&db.cursors, Some(9)).unwrap(); assert_eq!(reopened.next(), 10); }
#[cfg(feature = "jetstream")] #[test] fn clock_increases_even_when_ahead_of_the_wall_clock() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let clock = Clock::new(b"test_clock"); let ahead = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap() + 60_000_000; clock.resume_after(&db.cursors, Some(ahead)).unwrap();
// a failed commit takes nothing let failed = sequencer.commit( db.inner.batch(), |assign, _| { clock.tick(assign)?; Err::<(), _>(miette::miette!("stage failed")) }, |_| {}, ); assert!(failed.is_err());
let ticks = sequencer .commit( db.inner.batch(), |assign, _| Ok([clock.tick(assign)?, clock.tick(assign)?]), |ticks| ticks, ) .unwrap(); assert_eq!(ticks, [ahead + 1, ahead + 2]); }
#[cfg(feature = "jetstream")] #[test] fn a_commit_reads_the_wall_clock_once() { let (_tmp, db) = db(); let sequencer = Sequencer::new(db.cursors.clone()); let clock = Clock::new(b"test_clock"); let wall = || u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap();
let before = wall(); let ticks = sequencer .commit( db.inner.batch(), |assign, _| { let first = clock.tick(assign)?; // slow staging doesn't spread the commit's timestamps out std::thread::sleep(std::time::Duration::from_millis(5)); Ok([first, clock.tick(assign)?, clock.tick(assign)?]) }, |ticks| ticks, ) .unwrap(); assert!(ticks[0] >= before); assert_eq!(ticks, [ticks[0], ticks[0] + 1, ticks[0] + 2]);
// the next commit starts at the wall time, which has moved on let next = sequencer .commit(db.inner.batch(), |assign, _| clock.tick(assign), |t| t) .unwrap(); assert!(next >= ticks[0] + 5_000); }
#[test] fn unbroadcast_runs_are_announced_before_the_next_broadcast() { let (tx, mut rx) = tokio::sync::broadcast::channel(16); publish_in_order( &tx, [ (1, None), (2, None), (3, Some("live 3")), (4, Some("live 4")), (5, None), ], |position| match position { 2 => "marker 2", 5 => "marker 5", _ => panic!("unexpected marker for {position}"), }, ); let received: Vec<_> = std::iter::from_fn(|| rx.try_recv().ok()).collect(); assert_eq!(received, ["marker 2", "live 3", "live 4", "marker 5"]); }}