From 68bc5e1fe0eea1f2fcd33dccd4f2006870d05ca9 Mon Sep 17 00:00:00 2001 From: phil Date: Wed, 20 May 2026 19:40:42 -0400 Subject: [PATCH] host status summaries --- mini/src/host_summary.rs | 287 ++++++++++++++++++++++++++++ mini/src/main.rs | 63 +++++- mini/src/storage/repo_stats.rs | 32 ++++ stats-backfill/sql/host_summary.sql | 78 ++++++++ 4 files changed, 458 insertions(+), 2 deletions(-) create mode 100644 mini/src/host_summary.rs create mode 100644 stats-backfill/sql/host_summary.sql diff --git a/mini/src/host_summary.rs b/mini/src/host_summary.rs new file mode 100644 index 0000000..5ff4d09 --- /dev/null +++ b/mini/src/host_summary.rs @@ -0,0 +1,287 @@ +//! `host-summary` subcommand: drill into one PDS host and report the +//! disposition of every DID we've seen there, plus byte totals over the +//! captured set. +//! +//! Sibling to `count-repos` (per-host capture counts only) and +//! `summarize-stats` (global totals). The cross-tool equivalent lives in +//! `stats-backfill/sql/host_summary.sql` — see the plan file +//! `local/mini-listrepos-fetch-split-proposal.md` for the parallel. +//! +//! Conceptual note: mini's R/ is keyed `(pds, did)`, so a multi-homed +//! DID appears once per host here. stats-backfill's `repos.host` is +//! the upsert-chosen host (one row per DID), so its host-summary view +//! collapses the multi-homing direction. The two are complementary; +//! they're not expected to agree exactly on the same population. + +use std::collections::BTreeMap; + +use hubble_pds::Hostname; + +use crate::storage::{Db, DbError, commit, repo_state, repo_stats}; + +#[derive(Debug, thiserror::Error)] +pub enum Error { + #[error("storage: {0}")] + Storage(#[from] DbError), +} + +#[derive(Debug, Default)] +pub struct Summary { + /// All buckets are mutually exclusive — every R/ row contributes + /// to exactly one. Order of classification: tombstoned → + /// inactive → captured → failed → seen_uncaptured. + pub captured: u64, + pub seen_uncaptured: u64, + pub failed: u64, + pub inactive: u64, + pub tombstoned: u64, + + /// Sub-breakdown of `inactive` by `listrepos_status` (server- + /// reported string). `None` keyed as `"(none)"` for readability. + pub inactive_by_status: BTreeMap, + + /// Sub-breakdown of `failed` by `last_error`. Sorted at output + /// time; not capped here. + pub failed_by_error: BTreeMap, + + /// `expects_big = true` across the full R/ set (cuts across the + /// other buckets — a tombstoned big repo still counts here). + pub expects_big: u64, + + /// Sums over the captured set (one `s/` lookup per captured R/). + pub wire_bytes: u64, + pub car_bytes: u64, + pub record_count: u64, + pub blob_ref_count: u64, + pub blob_total_size: u64, +} + +impl Summary { + /// Sum of all five buckets. Equal to the count of R/ rows for the + /// host (every row landed in exactly one bucket). + pub fn total_seen(&self) -> u64 { + self.captured + + self.seen_uncaptured + + self.failed + + self.inactive + + self.tombstoned + } +} + +pub fn run(db: &Db, pds: &Hostname) -> Result { + let mut s = Summary::default(); + for (did, r) in repo_state::dids_for_pds(db, pds)? { + if r.expects_big { + s.expects_big += 1; + } + if r.tombstoned { + s.tombstoned += 1; + continue; + } + if !r.listrepos_active { + s.inactive += 1; + let key = r + .listrepos_status + .clone() + .unwrap_or_else(|| "(none)".to_string()); + *s.inactive_by_status.entry(key).or_insert(0) += 1; + continue; + } + // Active. Captured iff a C/ row exists. + if commit::get(db, pds, &did)?.is_some() { + s.captured += 1; + if let Some(stats) = repo_stats::get(db, pds, &did)? { + s.wire_bytes += stats.wire_bytes; + s.car_bytes += stats.car_bytes; + s.record_count += stats.record_count; + s.blob_ref_count += stats.blob_ref_count; + s.blob_total_size += stats.blob_total_size; + } + continue; + } + if let Some(err) = r.last_error.clone() { + s.failed += 1; + *s.failed_by_error.entry(err).or_insert(0) += 1; + } else { + s.seen_uncaptured += 1; + } + } + Ok(s) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::storage::open_temporary; + use crate::storage::repo_state::RepoState; + use crate::storage::repo_stats::RepoStats; + use rocksdb::WriteBatch; + + fn h(s: &str) -> Hostname { + Hostname::new(s).unwrap() + } + + fn r( + active: bool, + tombstoned: bool, + last_error: Option<&str>, + listrepos_status: Option<&str>, + expects_big: bool, + ) -> RepoState { + RepoState { + last_seen_rev: "3krev".into(), + last_fetched_at_unix: 0, + listrepos_active: active, + listrepos_status: listrepos_status.map(|s| s.into()), + tombstoned, + last_error: last_error.map(|s| s.into()), + expects_big, + } + } + + fn sample_commit() -> commit::StoredCommit { + commit::StoredCommit { + did: "did:plc:x".into(), + version: 3, + data: "bafyreidkjxqphmt3yhvdgmwx5p64q7s7tcmw3rj7ynvtu7tnp7nb6yurxe".into(), + rev: "3kabc".into(), + prev: None, + sig: vec![0xde, 0xad], + } + } + + fn sample_stats(wire: u64, car: u64, records: u64) -> RepoStats { + RepoStats { + wire_bytes: wire, + car_bytes: car, + record_count: records, + blob_ref_count: 1, + blob_total_size: 100, + ..RepoStats::default() + } + } + + #[test] + fn empty_host_is_all_zero() { + let (db, _g) = open_temporary(); + let s = run(&db, &h("example.com")).unwrap(); + assert_eq!(s.total_seen(), 0); + assert_eq!(s.captured, 0); + assert_eq!(s.wire_bytes, 0); + assert_eq!(s.expects_big, 0); + assert!(s.inactive_by_status.is_empty()); + assert!(s.failed_by_error.is_empty()); + } + + #[test] + fn classifies_each_row_into_exactly_one_bucket() { + let (db, _g) = open_temporary(); + let pds = h("a.example"); + + // captured: active, no error, has C/ + s/. + repo_state::put_fetched(&db, &pds, "did:plc:cap", &r(true, false, None, None, false)) + .unwrap(); + let mut batch = WriteBatch::default(); + commit::put_into(&mut batch, &db, &pds, "did:plc:cap", &sample_commit()).unwrap(); + repo_stats::put_into(&mut batch, &db, &pds, "did:plc:cap", &sample_stats(100, 400, 5)) + .unwrap(); + db.inner.write(batch).unwrap(); + + // inactive(deactivated): listrepos_active=false. + repo_state::put_fetched( + &db, + &pds, + "did:plc:inact", + &r(false, false, None, Some("deactivated"), false), + ) + .unwrap(); + + // tombstoned: takes precedence over inactive/active flags. + repo_state::put_fetched(&db, &pds, "did:plc:tomb", &r(true, true, None, None, false)) + .unwrap(); + + // failed: active, has last_error, no C/. + repo_state::put_fetched( + &db, + &pds, + "did:plc:err", + &r(true, false, Some("car too large"), None, true), + ) + .unwrap(); + + // seen_uncaptured: active, no error, no C/. + repo_state::put_fetched(&db, &pds, "did:plc:sun", &r(true, false, None, None, false)) + .unwrap(); + + let s = run(&db, &pds).unwrap(); + assert_eq!(s.captured, 1); + assert_eq!(s.inactive, 1); + assert_eq!(s.tombstoned, 1); + assert_eq!(s.failed, 1); + assert_eq!(s.seen_uncaptured, 1); + assert_eq!(s.total_seen(), 5); + assert_eq!(s.inactive_by_status.get("deactivated"), Some(&1)); + assert_eq!(s.failed_by_error.get("car too large"), Some(&1)); + // expects_big counted across all R/ — only the failed row had it. + assert_eq!(s.expects_big, 1); + } + + #[test] + fn byte_totals_sum_over_captured_set() { + let (db, _g) = open_temporary(); + let pds = h("a.example"); + + for (did, w, c, recs) in + [("did:plc:1", 100, 400, 5), ("did:plc:2", 50, 250, 3)].iter().copied() + { + repo_state::put_fetched(&db, &pds, did, &r(true, false, None, None, false)).unwrap(); + let mut batch = WriteBatch::default(); + commit::put_into(&mut batch, &db, &pds, did, &sample_commit()).unwrap(); + repo_stats::put_into(&mut batch, &db, &pds, did, &sample_stats(w, c, recs)).unwrap(); + db.inner.write(batch).unwrap(); + } + + let s = run(&db, &pds).unwrap(); + assert_eq!(s.captured, 2); + assert_eq!(s.wire_bytes, 150); + assert_eq!(s.car_bytes, 650); + assert_eq!(s.record_count, 8); + } + + #[test] + fn tombstoned_with_expects_big_still_counted_for_expects_big() { + let (db, _g) = open_temporary(); + let pds = h("a.example"); + // Tombstoned row that was previously known-big — expects_big + // should still count it even though its bucket is tombstoned. + repo_state::put_fetched( + &db, + &pds, + "did:plc:tomb_big", + &r(true, true, None, None, true), + ) + .unwrap(); + + let s = run(&db, &pds).unwrap(); + assert_eq!(s.tombstoned, 1); + assert_eq!(s.expects_big, 1); + } + + #[test] + fn scoped_to_one_host() { + let (db, _g) = open_temporary(); + let a = h("a.example"); + let b = h("b.example"); + + repo_state::put_fetched(&db, &a, "did:plc:1", &r(true, false, None, None, false)).unwrap(); + repo_state::put_fetched(&db, &b, "did:plc:1", &r(true, true, None, None, false)).unwrap(); + + let s_a = run(&db, &a).unwrap(); + assert_eq!(s_a.total_seen(), 1); + assert_eq!(s_a.tombstoned, 0); + + let s_b = run(&db, &b).unwrap(); + assert_eq!(s_b.total_seen(), 1); + assert_eq!(s_b.tombstoned, 1); + } +} diff --git a/mini/src/main.rs b/mini/src/main.rs index 5f7dc92..1b69516 100644 --- a/mini/src/main.rs +++ b/mini/src/main.rs @@ -8,6 +8,7 @@ mod compare_dump; mod count_repos; mod export_repo; mod find_pdses; +mod host_summary; mod pds; mod repo; mod snapshot_repos; @@ -54,6 +55,9 @@ pub enum Error { #[error("compare_dump: {0}")] CompareDump(#[from] compare_dump::Error), + #[error("host_summary: {0}")] + HostSummary(#[from] host_summary::Error), + #[error("io: {0}")] Io(#[from] std::io::Error), } @@ -154,6 +158,16 @@ enum Command { /// flushed-to-disk state). CompareDump(CompareDumpArgs), + /// Drill into one PDS host: print the disposition (captured / + /// seen-uncaptured / failed / inactive / tombstoned) of every DID + /// mini has under it, plus byte totals over the captured set, an + /// `expects_big` count, and top failure reasons. Sibling + /// `stats-backfill/sql/host_summary.sql` does the same shape on + /// the stats-backfill side. Pass `--read-only` to run alongside + /// an active `snapshot-repos` writer (reflects only flushed-to- + /// disk state). + HostSummary(HostSummaryArgs), + /// `find-pdses` then `snapshot-repos` in one run. Auto { #[command(flatten)] @@ -216,6 +230,13 @@ struct ExportRepoArgs { output: Option, } +#[derive(Args, Debug)] +struct HostSummaryArgs { + /// PDS hostname to summarise. + #[arg(long)] + pds: Hostname, +} + #[derive(Args, Debug)] struct CompareDumpArgs { /// Output path for the NDJSON dump (one JSON object per DID). @@ -270,6 +291,7 @@ async fn main() -> Result<(), Error> { Command::CountRepos | Command::SummarizeStats | Command::CompareDump(_) + | Command::HostSummary(_) | Command::CheckRepo(_) | Command::CheckDid(_) | Command::ExportRepo(_) @@ -277,8 +299,8 @@ async fn main() -> Result<(), Error> { if cli.read_only && !read_only_supported { eprintln!( "error: --read-only is only valid with the read commands \ - (count-repos, summarize-stats, compare-dump, check-repo, \ - check-did, export-repo)" + (count-repos, summarize-stats, compare-dump, host-summary, \ + check-repo, check-did, export-repo)" ); std::process::exit(2); } @@ -322,6 +344,9 @@ async fn main() -> Result<(), Error> { Command::CompareDump(args) => { run_compare_dump(&db, &args)?; } + Command::HostSummary(args) => { + run_host_summary(&db, &args)?; + } Command::Auto { find_pdses: fp, snapshot_repos: sr, @@ -500,6 +525,40 @@ fn run_summarize_stats(db: &storage::Db) -> Result<(), Error> { Ok(()) } +fn run_host_summary(db: &storage::Db, args: &HostSummaryArgs) -> Result<(), Error> { + let s = host_summary::run(db, &args.pds)?; + // Mutually exclusive bucket counts. Emit even at zero so the + // output shape is stable across hosts. + tracing::info!(host = %args.pds, bucket = "captured", dids = s.captured, "host-summary: bucket"); + tracing::info!(host = %args.pds, bucket = "seen_uncaptured", dids = s.seen_uncaptured, "host-summary: bucket"); + tracing::info!(host = %args.pds, bucket = "failed", dids = s.failed, "host-summary: bucket"); + tracing::info!(host = %args.pds, bucket = "inactive", dids = s.inactive, "host-summary: bucket"); + tracing::info!(host = %args.pds, bucket = "tombstoned", dids = s.tombstoned, "host-summary: bucket"); + for (status, count) in &s.inactive_by_status { + tracing::info!(host = %args.pds, status, dids = count, "host-summary: inactive_by_status"); + } + // Failure reasons: top 25, sorted desc by count. Mostly noise to + // log every long-tail variant. + let mut errs: Vec<(&String, &u64)> = s.failed_by_error.iter().collect(); + errs.sort_by(|a, b| b.1.cmp(a.1)); + for (err, count) in errs.iter().take(25) { + tracing::info!(host = %args.pds, error = %err, dids = count, "host-summary: failed_by_error"); + } + tracing::info!( + host = %args.pds, + total_seen = s.total_seen(), + captured = s.captured, + expects_big = s.expects_big, + wire_bytes = s.wire_bytes, + car_bytes = s.car_bytes, + record_count = s.record_count, + blob_ref_count = s.blob_ref_count, + blob_total_size = s.blob_total_size, + "host-summary: totals" + ); + Ok(()) +} + fn run_compare_dump(db: &storage::Db, args: &CompareDumpArgs) -> Result<(), Error> { let s = compare_dump::run(db, &args.dump_out)?; tracing::info!( diff --git a/mini/src/storage/repo_stats.rs b/mini/src/storage/repo_stats.rs index 9500639..ba62078 100644 --- a/mini/src/storage/repo_stats.rs +++ b/mini/src/storage/repo_stats.rs @@ -108,6 +108,20 @@ pub fn put_into( Ok(()) } +/// Point lookup of a single `s\0` row. Parallel to +/// [`super::commit::get`]. Used by `host-summary` to accumulate +/// captured-set byte totals without paying for a full `s/` scan. +pub fn get(db: &Db, pds: &Hostname, did: &str) -> Result> { + let cf = db.cf_default()?; + match db.inner.get_cf(cf, key(pds, did))? { + None => Ok(None), + Some(bytes) => Ok(Some(decode_raw( + &format!("s{}\\0{}", pds.as_ref(), did), + &bytes, + )?)), + } +} + /// Every `s\0` row, key parsed and value decoded. One row per /// repo that has been successfully fetched at least once (the row is /// overwritten on each successful re-fetch, so totals reflect the most @@ -240,4 +254,22 @@ mod tests { let (db, _g) = open_temporary(); assert!(scan_all(&db).unwrap().is_empty()); } + + #[test] + fn get_point_lookup() { + let (db, _g) = open_temporary(); + let pds = h("example.com"); + let did = "did:plc:abc"; + assert!(get(&db, &pds, did).unwrap().is_none()); + + let stats = sample(); + let mut batch = WriteBatch::default(); + put_into(&mut batch, &db, &pds, did, &stats).unwrap(); + db.inner.write(batch).unwrap(); + + let got = get(&db, &pds, did).unwrap().expect("row"); + assert_eq!(got, stats); + // Wrong did stays None. + assert!(get(&db, &pds, "did:plc:nope").unwrap().is_none()); + } } diff --git a/stats-backfill/sql/host_summary.sql b/stats-backfill/sql/host_summary.sql new file mode 100644 index 0000000..3a87d32 --- /dev/null +++ b/stats-backfill/sql/host_summary.sql @@ -0,0 +1,78 @@ +-- Per-host summary for stats-backfill. Cross-tool sibling of +-- `mini host-summary --pds `. Same shape buckets where they +-- apply, so the two outputs can be eyeballed side-by-side. +-- +-- Schema reference: stats-backfill/src/db.rs:38-104. +-- +-- Run with sqlite3, parameter set ahead of `.read`: +-- +-- sqlite3 path/to/stats-backfill.db -bail \ +-- -cmd ".parameter set :host 'berlin-user.eurosky.social'" \ +-- ".read stats-backfill/sql/host_summary.sql" +-- +-- Conceptual caveat: `repos.host` is the upsert-chosen host for a DID +-- (one row per DID, max-rev tie-break). mini's R/ is keyed (pds, did), +-- so a multi-homed DID appears under every host there. For migrated +-- accounts (started on bsky, now on indie), this view may put the DID +-- under the legacy host while mini has it under the indie host. Same +-- skew as compare-dump's resolver. + +.mode column +.headers on + +-- 1. Disposition by fetch_state (captured = 'done', failed_terminal, +-- still pending). Sum across rows = total DIDs sb chose this host for. +SELECT '--- fetch_state ---' AS section; +SELECT fetch_state, COUNT(*) AS dids +FROM repos +WHERE host = :host +GROUP BY fetch_state +ORDER BY dids DESC; + +-- 2. listRepos disposition. Mirrors mini's inactive bucket sub- +-- breakdown by status. +SELECT '--- listrepos disposition ---' AS section; +SELECT listrepos_active, listrepos_status, COUNT(*) AS dids +FROM repos +WHERE host = :host +GROUP BY listrepos_active, listrepos_status +ORDER BY dids DESC; + +-- 3. Top failure reasons among failed_terminal. Mirrors mini's +-- `failed_by_error` histogram. +SELECT '--- last_error (failed_terminal, top 25) ---' AS section; +SELECT last_error, COUNT(*) AS dids +FROM repos +WHERE host = :host AND fetch_state = 'failed_terminal' +GROUP BY last_error +ORDER BY dids DESC +LIMIT 25; + +-- 4. Measured totals over captured DIDs. Mirrors mini's totals line. +-- `star_zstd_bytes` is stats-backfill-only (mini doesn't compute it). +SELECT '--- measured totals (captured set) ---' AS section; +SELECT + COUNT(*) AS measured_dids, + SUM(m.car_bytes) AS total_car_bytes, + SUM(m.car_wire_bytes) AS total_car_wire_bytes, + SUM(m.star_zstd_bytes) AS total_star_zstd_bytes, + SUM(m.record_count) AS total_records, + SUM(m.blob_ref_count) AS total_blob_refs, + SUM(m.blob_total_size) AS total_blob_bytes +FROM repos r +JOIN repo_measurements m ON m.did = r.did +WHERE r.host = :host; + +-- 5. Recent fetch_attempt_log activity. Surfaces clustering of errors +-- in time. mini has no equivalent (it doesn't keep an attempt log). +SELECT '--- recent attempts (last 20) ---' AS section; +SELECT + datetime(attempted_at, 'unixepoch') AS attempted_at_iso, + substr(did, 1, 40) AS did, + http_status, + classification, + substr(error_summary, 1, 80) AS error_summary +FROM fetch_attempt_log +WHERE host = :host +ORDER BY attempted_at DESC +LIMIT 20; -- 2.51.2