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.
Rust
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296//! per-mode keyspace groups. this is the single place where mode-specific//! database state is declared: each feature contributes one group struct,//! opened and initialized here, and `Db` embeds each group behind one cfg.//!//! keyspace names, tuning, and debug rendering live in `db::schema`; the//! by-name and enumeration tables are generated in `db::registry`.
use fjall::{CompressionType, Database, Keyspace, KeyspaceCreateOptions};use miette::{IntoDiagnostic, Result};use std::cell::RefCell;use std::sync::Arc;
use crate::config::Config;
#[cfg(any( feature = "indexer", feature = "indexer_stream", feature = "jetstream", feature = "relay"))]use super::schema::{self, Ks};
#[cfg(any(feature = "indexer_stream", feature = "jetstream", feature = "relay"))]use miette::Context;#[cfg(any(feature = "indexer_stream", feature = "jetstream", feature = "relay"))]use std::sync::atomic::AtomicU64;
/// everything a keyspace needs to open itself.////// nameable so `Schema::options` can reference it in its signature, but all/// fields and methods are `pub(super)`: it cannot be constructed or used/// outside `db`.pub struct OpenCx<'a> { pub(super) db: &'a Arc<Database>, pub(super) cfg: &'a Config, pub(super) compression: &'a dyn Fn(&str, i32) -> CompressionType, /// names opened so far, checked against the registry at the end of open. pub(super) opened: RefCell<Vec<&'static str>>,}
impl OpenCx<'_> { pub(super) fn open_ks(&self, name: &str, opts: KeyspaceCreateOptions) -> Result<Keyspace> { self.db.keyspace(name, move || opts).into_diagnostic() }
pub(super) fn record_opened(&self, name: &'static str) { self.opened.borrow_mut().push(name); }}
#[cfg(feature = "indexer")]#[derive(Clone)]pub struct IndexerDb { /// maps `{DID} 00 {COL} 2f {RKey}` -> DAG-CBOR body or redacted CID marker pub(super) records: Ks<schema::Heads>, /// superseded bodies/markers: `{record key} 00 {death_rev}` -> body/CID/empty pub(super) history: Ks<schema::History>, /// operator-redacted versions: `{record key} 00 {cid}` -> empty pub(super) redactions: Ks<schema::Redactions>, /// backfill queue of `{ID}` -> empty pub pending: Ks<schema::Pending>, /// per-repo resync/retry state pub resync: Ks<schema::Resync>, /// live events buffered during backfill pub resync_buffer: Ks<schema::ResyncBuffer>, /// serializes lifecycle count rebuilds pub(crate) lifecycle_count_lock: Arc<std::sync::Mutex<()>>,}
#[cfg(feature = "indexer")]impl IndexerDb { pub(super) fn open(cx: &OpenCx) -> Result<Self> { Ok(Self { records: Ks::open(cx)?, history: Ks::open(cx)?, redactions: Ks::open(cx)?, pending: Ks::open(cx)?, resync: Ks::open(cx)?, resync_buffer: Ks::open(cx)?, lifecycle_count_lock: Arc::new(std::sync::Mutex::new(())), }) } pub(crate) fn record<K: AsRef<[u8]>>(&self, key: K) -> fjall::Result<Option<fjall::Slice>> { self.records.get(key) }
pub(crate) fn record_prefix<K: AsRef<[u8]>>(&self, prefix: K) -> fjall::Iter { self.records.prefix(prefix) }
pub(crate) fn record_range<K, R>(&self, range: R) -> fjall::Iter where K: AsRef<[u8]>, R: std::ops::RangeBounds<K>, { self.records.range(range) }
// read-side only: stream inflation is today's sole history reader #[cfg(feature = "indexer_stream")] pub(crate) fn history_range<K, R>(&self, range: R) -> fjall::Iter where K: AsRef<[u8]>, R: std::ops::RangeBounds<K>, { self.history.range(range) }
#[cfg(any(test, feature = "indexer_stream"))] pub(crate) fn stage_record<K, V>(&self, batch: &mut fjall::OwnedWriteBatch, key: K, value: V) where K: Into<fjall::UserKey>, V: Into<fjall::UserValue>, { batch.insert(&self.records, key, value); }
#[cfg(test)] #[allow(dead_code)] pub(crate) fn stage_history<K, V>(&self, batch: &mut fjall::OwnedWriteBatch, key: K, value: V) where K: Into<fjall::UserKey>, V: Into<fjall::UserValue>, { batch.insert(&self.history, key, value); }}
#[cfg(feature = "indexer_stream")]#[derive(Clone)]pub struct StreamDb { /// maps `{ID}` (u64 BE) -> `StoredEvent`, the source for the json stream api pub(super) events: Ks<schema::Events>, /// bodies needed only by legacy events whose record is no longer the /// current head: `{record key} 00 {cid}` -> raw DAG-CBOR body pub(crate) event_bodies: Ks<schema::EventBodies>, pub(crate) event_tx: tokio::sync::broadcast::Sender<crate::types::BroadcastEvent>, pub next_event_id: Arc<AtomicU64>,}
#[cfg(feature = "indexer_stream")]impl StreamDb { pub(super) fn open(cx: &OpenCx) -> Result<Self> { let (event_tx, _) = tokio::sync::broadcast::channel(512);
Ok(Self { events: Ks::open(cx)?, event_bodies: Ks::open(cx)?, event_tx, next_event_id: Arc::new(AtomicU64::new(0)), }) }
#[cfg(feature = "jetstream")] pub(crate) fn event<K: AsRef<[u8]>>(&self, key: K) -> fjall::Result<Option<fjall::Slice>> { self.events.get(key) }
pub(crate) fn event_range<K, R>(&self, range: R) -> fjall::Iter where K: AsRef<[u8]>, R: std::ops::RangeBounds<K>, { self.events.range(range) }
pub(crate) fn stage_event<K, V>(&self, batch: &mut fjall::OwnedWriteBatch, key: K, value: V) where K: Into<fjall::UserKey>, V: Into<fjall::UserValue>, { batch.insert(&self.events, key, value); }
pub(crate) fn approximate_event_count(&self) -> usize { self.events.approximate_len() }
pub(crate) fn event_head(&self) -> Result<Option<u64>> { let Some(guard) = self.events.iter().next_back() else { return Ok(None); }; let key = guard.key().into_diagnostic()?; let id = u64::from_be_bytes( key.as_ref() .try_into() .into_diagnostic() .wrap_err("expected stream event ID to be 8 bytes")?, ); Ok(Some(id)) }
/// resume event ids after the last stored event. pub(super) fn init(&self) -> Result<()> { let mut last_id = 0; if let Some(guard) = self.events.iter().next_back() { let k = guard.key().into_diagnostic()?; last_id = u64::from_be_bytes( k.as_ref() .try_into() .into_diagnostic() .wrap_err("expected to be id (8 bytes)")?, ); } // relaxed is fine since we are just initializing the db self.next_event_id .store(last_id + 1, std::sync::atomic::Ordering::Relaxed); Ok(()) }}
#[cfg(feature = "jetstream")]#[derive(Clone)]pub(crate) struct JetstreamDb { /// maps `{time_us}|{ID}` (16 bytes) -> jetstream event data pub(crate) events: Ks<schema::JetstreamEvents>, pub(crate) tx: tokio::sync::broadcast::Sender<crate::types::JetstreamBroadcast>, pub(crate) next_id: Arc<AtomicU64>, pub(crate) last_time_us: Arc<std::sync::atomic::AtomicI64>, /// serializes jetstream staging with batch commit in relay mode #[cfg(feature = "relay")] pub(crate) lock: Arc<parking_lot::Mutex<()>>,}
#[cfg(feature = "jetstream")]impl JetstreamDb { pub(super) fn open(cx: &OpenCx) -> Result<Self> { let (tx, _) = tokio::sync::broadcast::channel(512);
Ok(Self { events: Ks::open(cx)?, tx, next_id: Arc::new(AtomicU64::new(0)), last_time_us: Arc::new(std::sync::atomic::AtomicI64::new(0)), #[cfg(feature = "relay")] lock: Arc::new(parking_lot::Mutex::new(())), }) }
/// resume jetstream ids and time watermark after the last stored event. pub(super) fn init(&self) -> Result<()> { let mut last_id = 0; let mut last_time_us = 0; if let Some(guard) = self.events.iter().next_back() { let k = guard.key().into_diagnostic()?; let (time_us, id) = super::keys::parse_jetstream_event_key(&k)?; last_id = id; last_time_us = time_us; } self.next_id .store(last_id + 1, std::sync::atomic::Ordering::Relaxed); self.last_time_us .store(last_time_us as i64, std::sync::atomic::Ordering::Relaxed); Ok(()) }}
#[cfg(feature = "relay")]#[derive(Clone)]pub(crate) struct RelayDb { /// maps `{SEQ}` (u64 BE) -> re-encoded relay frame pub(crate) events: Ks<schema::RelayEvents>, pub(crate) next_seq: Arc<AtomicU64>, pub(crate) broadcast_tx: tokio::sync::broadcast::Sender<crate::types::RelayBroadcast>,}
#[cfg(feature = "relay")]impl RelayDb { pub(super) fn open(cx: &OpenCx) -> Result<Self> { let (broadcast_tx, _) = tokio::sync::broadcast::channel(512);
Ok(Self { events: Ks::open(cx)?, next_seq: Arc::new(AtomicU64::new(0)), broadcast_tx, }) }
/// resume relay sequence numbers after the last stored frame. pub(super) fn init(&self) -> Result<()> { let mut last_relay_seq = 0u64; if let Some(guard) = self.events.iter().next_back() { let k = guard.key().into_diagnostic()?; last_relay_seq = u64::from_be_bytes( k.as_ref() .try_into() .into_diagnostic() .wrap_err("relay_events: invalid key length")?, ); } self.next_seq .store(last_relay_seq + 1, std::sync::atomic::Ordering::Relaxed); Ok(()) }}