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.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283//! v12 copies the backfill queue into `pending_dids`, so a queue filled before//! the did-order copy existed can be listed by did too. the pass writes a//! different keyspace than it scans and inserting the same key twice changes//! nothing, so a crash that re-visits a chunk is harmless.
#[cfg(feature = "indexer")]use fjall::{Keyspace, OwnedWriteBatch};#[cfg(feature = "indexer")]use miette::Result;
#[cfg(feature = "indexer")]use crate::db::{Db, keys};
#[cfg(feature = "indexer")]use super::ChunkBudget;use super::Pass;
pub(super) const PASSES: &[Pass] = &[ #[cfg(feature = "indexer")] Pass { name: "pending_dids", scan: "pending", visit: copy_pending_entry, budget: ChunkBudget::DEFAULT, optional_on_absent: false, },];
#[cfg(feature = "indexer")]fn copy_pending_entry( db: &Db, batch: &mut OwnedWriteBatch, _scanned: &Keyspace, key: &[u8], value: &[u8],) -> Result<usize> { let did_key = keys::pending_did_key(value, key); let staged = did_key.len(); batch.insert(&db.indexer.pending_dids, did_key, []); Ok(staged)}
#[cfg(all(test, feature = "indexer"))]mod tests { use super::*; use crate::config::Config; use crate::db::migration::rewind_version_for_test; use jacquard_common::types::did::Did; use miette::IntoDiagnostic; use tempfile::tempdir;
#[test] fn copies_a_queue_that_has_no_did_order_copy_yet() -> Result<()> { let tmp = tempdir().into_diagnostic()?; let cfg = Config { database_path: tmp.path().to_path_buf(), ..Default::default() }; let plc = Did::new_static("did:plc:yk4q3id7id6p5z3bypvshc64").into_diagnostic()?; let web = Did::new_static("did:web:example.com").into_diagnostic()?; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); // an older build leaves queue entries with no did-order copy for (id, did) in [(3, &plc), (9, &web), (1, &web)] { batch.insert( &db.indexer.pending, keys::pending_key(id), keys::repo_key(did), ); } batch.commit().into_diagnostic()?; rewind_version_for_test(&db, 9)?; db.persist()?; }
let db = Db::open(&cfg)?; db.indexer.assert_pending_dids_match(); assert_eq!(db.indexer.pending_dids.iter().count(), 3); Ok(()) }}