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.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142use fjall::{Keyspace, OwnedWriteBatch};use miette::{IntoDiagnostic, Result};
use crate::db::Db;use crate::pds_discovery::canonical_pds_host;use crate::pds_meta::HostStatus;
use super::{ChunkBudget, Pass};
const STATUS_PREFIX: &[u8] = b"pds|status|";
/// v11 splits operator bans from observed status. existing durable `Banned`/// statuses are treated as operator bans, since that is how the pre-v11 value/// was enforced; the old status row is removed so `unban` can clear the ban.////// bans move first, then remaining status aliases canonicalize.pub(super) const PASSES: &[Pass] = &[ Pass { name: "pds_bans", scan: "filter", visit: rewrite_ban, budget: ChunkBudget::DEFAULT, optional_on_absent: false, }, Pass { name: "pds_statuses", scan: "filter", visit: rewrite_status, budget: ChunkBudget::DEFAULT, optional_on_absent: false, },];
fn canonical_host(key: &[u8]) -> Option<(String, bool)> { let raw = key.strip_prefix(STATUS_PREFIX)?; let host = std::str::from_utf8(raw).ok()?; let canonical = canonical_pds_host(host).ok()?; Some(( canonical.as_str().to_string(), canonical.as_str().as_bytes() == raw, ))}
fn status(value: &[u8]) -> Result<HostStatus> { rmp_serde::from_slice(value).into_diagnostic()}
fn rewrite_ban( _db: &Db, batch: &mut OwnedWriteBatch, scanned: &Keyspace, key: &[u8], value: &[u8],) -> Result<usize> { let Some((host, _)) = canonical_host(key) else { return Ok(0); }; if status(value)? != HostStatus::Banned { return Ok(0); }
let ban_key = crate::db::pds_meta::pds_ban_key(&host); batch.insert(scanned, ban_key.clone(), []); batch.remove(scanned, key); Ok(ban_key.len())}
fn rewrite_status( _db: &Db, batch: &mut OwnedWriteBatch, scanned: &Keyspace, key: &[u8], value: &[u8],) -> Result<usize> { let Some((host, already_canonical)) = canonical_host(key) else { return Ok(0); }; if already_canonical || status(value)? == HostStatus::Banned { return Ok(0); }
let new_key = crate::db::pds_meta::pds_status_key(&host); batch.insert(scanned, new_key.clone(), value); batch.remove(scanned, key); Ok(new_key.len() + value.len())}
#[cfg(test)]mod tests { use super::*; use crate::config::Config; use crate::db::migration::rewind_version_for_test; use tempfile::tempdir;
#[test] fn splits_legacy_bans_before_canonicalizing_other_status_aliases() -> Result<()> { let tmp = tempdir().into_diagnostic()?; let cfg = Config { database_path: tmp.path().to_path_buf(), ..Default::default() }; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); let raw = crate::db::pds_meta::pds_status_key("BANNED.EXAMPLE."); let canonical = crate::db::pds_meta::pds_status_key("banned.example"); batch.insert( &db.filter, raw, rmp_serde::to_vec(&HostStatus::Banned).into_diagnostic()?, ); batch.insert( &db.filter, canonical, rmp_serde::to_vec(&HostStatus::Active).into_diagnostic()?, ); batch.commit().into_diagnostic()?; rewind_version_for_test(&db, 9)?; db.persist()?; }
let db = Db::open(&cfg)?; assert!( !db.filter .contains_key(crate::db::pds_meta::pds_status_key("BANNED.EXAMPLE.")) .into_diagnostic()? ); let observed = db .filter .get(crate::db::pds_meta::pds_status_key("banned.example")) .into_diagnostic()? .expect("non-ban status must remain canonical"); assert_eq!(status(observed.as_ref())?, HostStatus::Active); assert!( db.filter .contains_key(crate::db::pds_meta::pds_ban_key("banned.example")) .into_diagnostic()? ); Ok(()) }}