//! audit hubble-sync's derived repo-count indexes against the repos themselves //! //! hubble-sync keeps two counter indexes over `ri|` (repo info): //! //! - `rc|` repos by sync state (storage/repo/info_idx_state_count.rs) //! - `ra|` repos by effective account status (info_idx_status_count.rs) //! //! both are maintained incrementally (+1/-1 per transition, in the same atomic //! batch as the `ri|` write) and are never rebuilt. this walks every `ri|` //! record, tallies what each counter *should* be, and diffs. //! //! it also prints every `rc|`/`ra|` key actually present, including ones the //! current bucket scheme neither reads nor writes -- a leftover key holding a //! stale positive is the signature of a bucket-scheme change under a live db. //! //! read-only: opens as a rocksdb secondary by default, so it is safe to run //! against a live instance. use std::collections::{BTreeMap, HashSet}; use std::error::Error; use std::hash::{Hash, Hasher}; use std::path::PathBuf; use clap::Parser; use dasl::drisl; use rocksdb::{ ColumnFamilyDescriptor, DB, IteratorMode, MergeOperands, Options, ReadOptions, }; use serde::Deserialize; /// column families, per hubble-sync-rocksdb's defaults const POINT_CF: &str = "default"; const QUEUE_CF: &str = "queue"; const COUNTER_CF: &str = "counter"; /// hubble-sync key prefixes (storage/mod.rs), under the engine-level prefix const REPO_INFO: &[u8] = b"ri|"; const STATE_COUNT: &[u8] = b"rc|"; const STATUS_COUNT: &[u8] = b"ra|"; const RESYNC_QUEUE: &[u8] = b"rq|"; #[derive(Parser, Debug)] #[command(about = "audit hubble-sync repo-count indexes against ri| records")] struct Args { /// path to the rocksdb directory (the `db` dir, not the data dir) #[arg(long)] db: PathBuf, /// scratch dir for the secondary instance (must be writable) #[arg(long)] secondary: Option, /// open read-only instead of as a secondary /// /// only for a db with no live writer: read-only does not replay the WAL, so /// recent writes from a running instance would be invisible. #[arg(long, conflicts_with = "secondary")] read_only: bool, /// hubble-sync's storage prefix, as hex (SyncConfig::hubble_sync_storage_prefix) #[arg(long, default_value = "00")] prefix_hex: String, /// also audit the resync queue index (rq|) against ri| /// /// same reconcile call sites as the counters, but plain puts/deletes in the /// queue CF -- so it separates "the reconcile inputs were wrong" from /// "something happened to the counters". costs ~24 bytes of memory per /// queued repo. #[arg(long)] audit_queue: bool, /// log progress every N repos scanned #[arg(long, default_value_t = 1_000_000)] progress_every: u64, } fn main() -> Result<(), Box> { let args = Args::parse(); let prefix = decode_hex(&args.prefix_hex)?; let db = open(&args)?; println!("db: {}", args.db.display()); println!("prefix: 0x{}\n", args.prefix_hex); let state_counters = read_counters(&db, &prefix, STATE_COUNT)?; let status_counters = read_counters(&db, &prefix, STATUS_COUNT)?; let tally = scan_repos(&db, &prefix, args.progress_every, args.audit_queue)?; println!( "scanned {} repo info records ({} failed to decode)\n", tally.repos, tally.undecodable ); let state_bad = report( "by sync state (rc|)", &state_counters, &tally.by_state, &expected_state_buckets(), ); let status_bad = report( "by account status (ra|)", &status_counters, &tally.by_status, &expected_status_buckets(), ); let queue_bad = if args.audit_queue { report_queue(&db, &prefix, &tally)? } else { 0 }; if tally.undecodable > 0 { println!("\n{} record(s) could not be decoded -- tallies are incomplete", tally.undecodable); } if state_bad == 0 && status_bad == 0 && queue_bad == 0 { println!("\nOK: every index matches the repos on disk"); return Ok(()); } println!( "\nMISMATCH: {state_bad} state bucket(s), {status_bad} status bucket(s), \ {queue_bad} queue discrepancy kind(s)" ); std::process::exit(1); } fn open(args: &Args) -> Result> { let mut opts = Options::default(); opts.create_if_missing(false); let mut counter_opts = Options::default(); // required to read the counter CF: without the operator, un-compacted merge // operands are not applied on read. counter_opts.set_merge_operator_associative("hubble_sync.counter.i64", counter_merge); let cfs = vec![ ColumnFamilyDescriptor::new(POINT_CF, Options::default()), ColumnFamilyDescriptor::new(QUEUE_CF, Options::default()), ColumnFamilyDescriptor::new(COUNTER_CF, counter_opts), ]; if args.read_only { return Ok(DB::open_cf_descriptors_read_only(&opts, &args.db, cfs, false)?); } let secondary = args .secondary .clone() .unwrap_or_else(|| std::env::temp_dir().join("hubble-count-audit")); std::fs::create_dir_all(&secondary)?; let db = DB::open_cf_descriptors_as_secondary(&opts, &args.db, &secondary, cfs)?; // pull in whatever the primary has written since it last flushed db.try_catch_up_with_primary()?; Ok(db) } /// same arithmetic as hubble-sync-rocksdb's `counter_merge` fn counter_merge(_k: &[u8], existing: Option<&[u8]>, operands: &MergeOperands) -> Option> { let mut current = match existing.map(i64_be) { Some(Some(v)) => v, Some(None) => return existing.map(<[u8]>::to_vec), // corrupt: pass through None => 0, }; for op in operands { let Some(delta) = i64_be(op) else { return Some(op.to_vec()); // corrupt: pass through }; current = current.saturating_add(delta); } Some(current.to_be_bytes().to_vec()) } fn i64_be(bytes: &[u8]) -> Option { bytes.try_into().ok().map(i64::from_be_bytes) } /// every `` counter key present, keyed by its bucket suffix fn read_counters( db: &DB, prefix: &[u8], kind: &[u8], ) -> Result, Box> { let scan_prefix = [prefix, kind].concat(); let mut out = BTreeMap::new(); for item in scan(db, COUNTER_CF, &scan_prefix) { let (key, value) = item?; let bucket = String::from_utf8_lossy(&key).into_owned(); let v = match i64_be(&value) { Some(v) => CounterValue::Count(v), None => CounterValue::Corrupt(value.len()), }; out.insert(bucket, v); } Ok(out) } #[derive(Debug, Clone, Copy)] enum CounterValue { Count(i64), Corrupt(usize), } #[derive(Default)] struct Tally { repos: u64, undecodable: u64, by_state: BTreeMap, by_status: BTreeMap, /// fingerprints of the rq| keys the repos say should exist expected_queue_keys: HashSet, expected_queue: u64, } fn scan_repos( db: &DB, prefix: &[u8], progress_every: u64, audit_queue: bool, ) -> Result> { let scan_prefix = [prefix, REPO_INFO].concat(); let mut tally = Tally::default(); for item in scan(db, POINT_CF, &scan_prefix) { let (did_key, value) = item?; // app-owned slots live at `ri|\0u`; DIDs never contain NUL, so a // NUL in the suffix means this is a slot, not a repo info record. if did_key.contains(&0x00) { continue; } tally.repos += 1; if tally.repos % progress_every == 0 { eprintln!(" ...{} repos", tally.repos); } let info: RepoInfo = match drisl::from_slice(&value) { Ok(info) => info, Err(err) => { tally.undecodable += 1; if tally.undecodable <= 10 { eprintln!( " undecodable ri| record for {}: {err}", String::from_utf8_lossy(&did_key) ); } continue; } }; *tally.by_state.entry(info.sync_status.bucket()).or_default() += 1; *tally .by_status .entry(info.account_status_kind().to_string()) .or_default() += 1; if audit_queue && let Some(key) = info.expected_queue_key(&did_key) { tally.expected_queue += 1; tally.expected_queue_keys.insert(fingerprint(&key)); } } Ok(tally) } /// prefix-bounded scan of one CF, yielding prefix-stripped keys fn scan<'a>( db: &'a DB, cf_name: &str, scan_prefix: &[u8], ) -> impl Iterator, Vec), Box>> + 'a { let cf = db.cf_handle(cf_name).expect("column family present"); let mut opts = ReadOptions::default(); opts.set_iterate_lower_bound(scan_prefix.to_vec()); if let Some(end) = prefix_end_exclusive(scan_prefix) { opts.set_iterate_upper_bound(end); } opts.fill_cache(false); // don't evict a live instance's working set let strip = scan_prefix.len(); db.iterator_cf_opt(&cf, opts, IteratorMode::Start) .map(move |r| match r { Ok((k, v)) => Ok((k[strip..].to_vec(), v.into_vec())), Err(e) => Err(Box::new(e) as Box), }) } fn prefix_end_exclusive(prefix: &[u8]) -> Option> { let mut end = prefix.to_vec(); while let Some(last) = end.last_mut() { if *last < 0xFF { *last += 1; return Some(end); } end.pop(); } None } /// print counter-vs-actual for one index; returns the number of bad rows fn report( title: &str, counters: &BTreeMap, actual: &BTreeMap, expected_buckets: &[String], ) -> usize { println!("{title}"); println!( " {:<44}{:>12}{:>12}{:>12}", "bucket", "counter", "actual", "delta" ); let mut bad = 0; let mut buckets: Vec<&String> = expected_buckets.iter().collect(); // any key present on disk that the current scheme does not know about let unknown: Vec<&String> = counters .keys() .filter(|k| !expected_buckets.contains(k)) .collect(); buckets.extend(unknown.iter().copied()); for bucket in buckets { let known = expected_buckets.contains(bucket); let counter = counters.get(bucket).copied(); let actual_n = actual.get(bucket).copied().unwrap_or(0); let (counter_str, delta_str, row_bad) = match counter { Some(CounterValue::Count(c)) => { let delta = c - actual_n as i64; (c.to_string(), delta.to_string(), delta != 0) } Some(CounterValue::Corrupt(len)) => { (format!("<{len}B>"), "?".to_string(), true) } None => ("-".to_string(), (-(actual_n as i64)).to_string(), actual_n != 0), }; if row_bad { bad += 1; } let flag = if !known { " <- key not in the current bucket scheme" } else if row_bad { " <-" } else { "" }; let actual_str = if known { actual_n.to_string() } else { "-".into() }; println!(" {bucket:<44}{counter_str:>12}{actual_str:>12}{delta_str:>12}{flag}"); } println!(); bad } /// audit `rq|` against what the repos say should be in it /// /// the index should hold exactly one entry per (upstream-active, desynchronized) /// repo, keyed by its pds host and desync `due_at` -- see `Index::should` and /// `encode_key` in storage/repo/info_idx_resync.rs. fn report_queue(db: &DB, prefix: &[u8], tally: &Tally) -> Result> { let scan_prefix = [prefix, RESYNC_QUEUE].concat(); let mut present = 0u64; let mut matched = 0u64; for item in scan(db, QUEUE_CF, &scan_prefix) { let (key, _) = item?; present += 1; if tally.expected_queue_keys.contains(&fingerprint(&key)) { matched += 1; } } let missing = tally.expected_queue.saturating_sub(matched); let orphaned = present.saturating_sub(matched); println!("resync queue index (rq|)"); println!(" {:<44}{:>12}", "repos that should be queued", tally.expected_queue); println!(" {:<44}{:>12}", "entries present", present); println!(" {:<44}{:>12}", "present and correctly keyed", matched); println!( " {:<44}{:>12}{}", "missing (repo desynced, not queued)", missing, if missing > 0 { " <-" } else { "" } ); println!( " {:<44}{:>12}{}", "orphaned (queued, stale key or no repo)", orphaned, if orphaned > 0 { " <-" } else { "" } ); println!(); Ok(usize::from(missing > 0) + usize::from(orphaned > 0)) } /// 64-bit fingerprint of a queue key, so the expected set stays small /// /// collisions would under-report a discrepancy; at these cardinalities the odds /// are negligible and the audit is a diagnostic, not a proof. fn fingerprint(key: &[u8]) -> u64 { let mut h = std::collections::hash_map::DefaultHasher::new(); key.hash(&mut h); h.finish() } // --- the stored shapes, trimmed to what the buckets are computed from --- // // field and variant names must match hubble-sync's `DbRepoInfo` (storage/repo/ // info.rs). unknown fields are ignored; an unknown *variant* is a decode error, // which gets reported rather than silently mis-bucketed. #[derive(Debug, Deserialize)] struct RepoInfo { upstream_status: AccountStatus, #[serde(default)] moderation: Option, sync_status: SyncStatus, identity: Identity, } #[derive(Debug, Deserialize)] struct Identity { /// already normalized when written (it comes from the host interner) pds: String, } impl RepoInfo { /// the `rq|` key this repo should be indexed under, if any /// /// mirrors `Index::should` -- which keys off *upstream* status, ignoring /// local moderation -- and `encode_key`, minus the scan prefix. fn expected_queue_key(&self, did: &[u8]) -> Option> { let (AccountStatus::Active, SyncStatus::Desynchronized(d)) = (&self.upstream_status, &self.sync_status) else { return None; }; let mut k = Vec::with_capacity(self.identity.pds.len() + 9 + did.len()); k.extend_from_slice(self.identity.pds.as_bytes()); k.push(0x00); k.extend_from_slice(&d.due_at.to_be_bytes()); k.extend_from_slice(did); Some(k) } /// mirrors `RepoInfo::account_status()` + `AccountStatusKind::name()` fn account_status_kind(&self) -> &'static str { match (&self.upstream_status, &self.moderation) { (AccountStatus::Deleted, _) => "deleted", (_, Some(m)) => match m.action { ModAction::Suspend => "suspended", ModAction::Takedown => "takendown", }, (upstream, None) => upstream.kind_name(), } } } #[derive(Debug, Deserialize)] enum AccountStatus { Active, Deleted, Deactivated, Suspended, Takendown, #[expect(dead_code, reason = "payload must decode, but only the kind is bucketed")] Inactive(String), } impl AccountStatus { fn kind_name(&self) -> &'static str { match self { Self::Active => "active", Self::Deleted => "deleted", Self::Deactivated => "deactivated", Self::Suspended => "suspended", Self::Takendown => "takendown", Self::Inactive(_) => "other", } } } #[derive(Debug, Deserialize)] struct Moderation { action: ModAction, } #[derive(Debug, Deserialize)] enum ModAction { #[serde(rename = "suspend")] Suspend, #[serde(rename = "takedown")] Takedown, } #[derive(Debug, Deserialize)] enum SyncStatus { Synchronized, Desynchronized(Desynchronized), OutOfScope {}, NonActive {}, Gone {}, } impl SyncStatus { /// mirrors `Bucket::encode_key` (info_idx_state_count.rs), minus the prefix fn bucket(&self) -> String { match self { Self::Synchronized => "synchronized".into(), Self::OutOfScope {} => "outOfScope".into(), Self::NonActive {} => "nonActive".into(), Self::Gone {} => "gone".into(), Self::Desynchronized(d) => format!("desynchronized|{}", d.reason.name()), } } } #[derive(Debug, Deserialize)] struct Desynchronized { reason: DesyncReason, /// unix ms, same units the rq| key encodes due_at: u64, } #[derive(Debug, Deserialize)] enum DesyncReason { FirstSeen, CameIntoScope, UnresolvableIdentity, FirehoseSync {}, FirehoseFail {}, FirehoseAccountDesynchronized, FirehoseAccountThrottled, FutureRev {}, Throttled {}, Sync11Lax, AppRequested {}, } impl DesyncReason { /// mirrors `DesyncReasonKind::name()` fn name(&self) -> &'static str { match self { Self::FirstSeen => "first_seen", Self::CameIntoScope => "came_into_scope", Self::UnresolvableIdentity => "unresolvable_identity", Self::FirehoseSync {} => "firehose_sync", Self::FirehoseFail {} => "firehose_fail", Self::FirehoseAccountDesynchronized => "firehose_account_desync", Self::FirehoseAccountThrottled => "firehose_account_throttled", Self::FutureRev {} => "future_rev", Self::Throttled {} => "throttled", Self::Sync11Lax => "sync11_lax", Self::AppRequested {} => "app_requested", } } } const DESYNC_REASONS: &[&str] = &[ "first_seen", "came_into_scope", "unresolvable_identity", "firehose_sync", "firehose_fail", "firehose_account_desync", "firehose_account_throttled", "future_rev", "throttled", "sync11_lax", "app_requested", ]; fn expected_state_buckets() -> Vec { let mut v: Vec = ["synchronized", "outOfScope", "nonActive", "gone"] .iter() .map(|s| s.to_string()) .collect(); v.extend(DESYNC_REASONS.iter().map(|r| format!("desynchronized|{r}"))); v } fn expected_status_buckets() -> Vec { [ "active", "deactivated", "suspended", "takendown", "deleted", "other", ] .iter() .map(|s| s.to_string()) .collect() } fn decode_hex(s: &str) -> Result, Box> { if s.len() % 2 != 0 { return Err("prefix hex must have an even number of digits".into()); } (0..s.len()) .step_by(2) .map(|i| u8::from_str_radix(&s[i..i + 2], 16).map_err(|e| e.to_string().into())) .collect() }