From 21bb6bb7f59f55074e1c6344388876f548bee21a Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Mon, 5 Oct 2026 16:05:13 +0300 Subject: [PATCH] [backfill] stream filtered full repo cars instead of buffering them when a filtered backfill falls back to a full getRepo, the whole body used to be collected, every block mapped and every leaf walked before the filter dropped what it didn't want, so one 975 MB atcr hold took hydrant to 2.7 GB. the car is now decoded in chunks as it arrives, on a blocking thread fed through a bounded channel, and only the commit, records whose $type the filter wants and mst nodes that could still reach a wanted range are kept. every other node is remembered as a sha-256 digest so the scan over what was kept can tell a node it ruled out from one that's missing. a leaf that went missing is fetched by cid, or from one whole car with the same root if the pds has no getBlocks. the fetch and the insert share one filter snapshot, so a filter change in between can't delete records the fetch never kept. on real cars the atcr hold peaks at 0.95 GB instead of 2.0 GB, and an 83 MB bsky repo keeps 5 MB of it. --- src/backfill/filtered_car.rs | 317 +++++++++++++++++++++++++++++ src/backfill/mod.rs | 1 + src/backfill/sparse.rs | 2 +- src/backfill/worker/process.rs | 305 ++++++++++++++++++++++----- src/backfill/worker/task.rs | 119 ++++++++++- src/bin/backfill_strategy_bench.rs | 1 + src/car.rs | 289 +++++++++++++++++++++++--- src/sparse_mst.rs | 157 +++++++++++++- 8 files changed, 1109 insertions(+), 82 deletions(-) create mode 100644 src/backfill/filtered_car.rs diff --git a/src/backfill/filtered_car.rs b/src/backfill/filtered_car.rs new file mode 100644 index 0000000..e72dc8d --- /dev/null +++ b/src/backfill/filtered_car.rs @@ -0,0 +1,317 @@ +//! a full repo car read block by block, keeping only what a filtered backfill can use. a +//! repo's car can be gigabytes of records in collections the filter doesn't want, and parsing +//! it whole means holding every one of them at once. + +use std::sync::Arc; + +use bytes::Bytes; +use cid::Cid as IpldCid; +use miette::Result; +use serde::Deserialize; + +use crate::car::{Blocks, CarStream, CommitBlocks}; +use crate::filter::FilterConfig; +use crate::mst::decode_node; +use crate::sparse_mst::{KeyRange, RuledOut, node_may_reach, sparse_ranges}; + +pub(crate) struct FilteredCar { + stream: CarStream, + keep: Keep, +} + +struct Keep { + filter: Arc, + ranges: Vec, + blocks: CommitBlocks, + ruled_out: RuledOut, + stats: FilteredStats, +} + +#[derive(Debug, Default, Clone, Copy)] +pub(crate) struct FilteredStats { + pub(crate) blocks: usize, + pub(crate) bytes: usize, + pub(crate) kept_blocks: usize, + pub(crate) kept_bytes: usize, +} + +/// what a filtered car left behind: its commit, the records the filter wants and the mst +/// nodes that could lead to them, plus the nodes it ruled out on the way +#[derive(Debug)] +pub(crate) struct FilteredRepo { + pub(crate) root: IpldCid, + pub(crate) ranges: Vec, + pub(crate) blocks: Blocks, + pub(crate) ruled_out: RuledOut, + pub(crate) stats: FilteredStats, +} + +impl FilteredCar { + /// `None` when `filter` takes every collection, since then every block is wanted anyway + pub(crate) fn new(filter: Arc, verify_cids: bool) -> Option { + let ranges = sparse_ranges(&filter.collections); + if ranges.is_empty() { + return None; + } + Some(Self { + stream: CarStream::new(verify_cids), + keep: Keep { + filter, + ranges, + blocks: CommitBlocks::new(), + ruled_out: RuledOut::default(), + stats: FilteredStats::default(), + }, + }) + } + + pub(crate) fn push(&mut self, chunk: &[u8]) -> Result<()> { + let keep = &mut self.keep; + self.stream.push(chunk, |root, cid, bytes| { + keep.offer(root, cid, bytes); + Ok(()) + }) + } + + pub(crate) fn finish(self) -> Result { + let Keep { + ranges, + blocks, + ruled_out, + stats, + .. + } = self.keep; + let car = self.stream.finish_with(blocks)?; + Ok(FilteredRepo { + root: car.root, + ranges, + blocks: car.blocks, + ruled_out, + stats, + }) + } +} + +impl Keep { + fn offer(&mut self, root: &IpldCid, cid: IpldCid, bytes: &[u8]) { + self.stats.blocks += 1; + self.stats.bytes += bytes.len(); + let wanted = if cid == *root { + true + } else if let Some(collection) = record_type(bytes) { + // a record filed under a collection other than its $type is dropped here, and the + // scan finds it missing later + self.filter.matches_collection(collection) + } else if let Ok(entries) = decode_node(bytes) { + // a node whose cid can't be remembered as ruled out has to be kept instead + node_may_reach(&self.ranges, &entries) || !self.ruled_out.insert(&cid) + } else { + false + }; + if wanted { + self.stats.kept_blocks += 1; + self.stats.kept_bytes += bytes.len(); + // copied, so a kept block doesn't pin the chunk it arrived in + self.blocks.insert(cid, Bytes::copy_from_slice(bytes)); + } + } +} + +/// a record's `$type`. commits, mst nodes and anything that isn't a map with a string `$type` +/// have none +fn record_type(bytes: &[u8]) -> Option<&str> { + #[derive(Deserialize)] + struct Typed<'a> { + #[serde(rename = "$type", borrow)] + ty: Option<&'a str>, + } + serde_ipld_dagcbor::from_slice::(bytes).ok()?.ty +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::filter::FilterMode; + use crate::sparse_mst::SparseScanner; + use jacquard_repo::mst::util::compute_cid; + use jacquard_repo::{BlockStore, MemoryBlockStore, Mst}; + use smol_str::SmolStr; + + fn record(collection: &str, n: usize) -> Vec { + serde_ipld_dagcbor::to_vec(&serde_json::json!({ "$type": collection, "n": n })).unwrap() + } + + /// a whole repo car: the commit, every mst node and every record, keyed `collection/rkey` + async fn repo_car(records: &[(String, Vec)]) -> (Bytes, IpldCid) { + let store = Arc::new(MemoryBlockStore::new()); + let mut mst = Mst::new(store.clone()); + for (key, body) in records { + let cid = compute_cid(body).unwrap(); + mst.add_mut(key, cid).await.unwrap(); + store + .put_many([(cid, Bytes::copy_from_slice(body))]) + .await + .unwrap(); + } + let mst_root = mst.persist().await.unwrap(); + let commit = jacquard_repo::commit::Commit:: { + did: jacquard_common::types::did::Did::new_static("did:plc:aaaaaaaaaaaaaaaaaaaaaaaa") + .unwrap(), + version: 3, + data: mst_root, + rev: "3l2abcdefgh22".parse().unwrap(), + prev: None, + sig: Bytes::new(), + }; + let commit_cid = commit.to_cid().unwrap(); + let mut car = Vec::new(); + let mut writer = + iroh_car::CarWriter::new(iroh_car::CarHeader::new_v1(vec![commit_cid]), &mut car); + writer + .write(commit_cid, commit.to_cbor().unwrap()) + .await + .unwrap(); + mst.write_blocks_to_car(&mut writer).await.unwrap(); + writer.finish().await.unwrap(); + (car.into(), mst_root) + } + + fn filter(collections: &[&str]) -> Arc { + let mut filter = FilterConfig::new(FilterMode::Filter); + filter.collections = collections.iter().map(SmolStr::new).collect(); + Arc::new(filter) + } + + fn stream(car: &[u8], collections: &[&str], chunk: usize) -> FilteredRepo { + let mut filtered = FilteredCar::new(filter(collections), true).unwrap(); + for piece in car.chunks(chunk) { + filtered.push(piece).unwrap(); + } + filtered.finish().unwrap() + } + + fn scan( + repo: FilteredRepo, + mst_root: IpldCid, + ) -> (Vec<(SmolStr, IpldCid)>, Blocks) { + let mut scanner = SparseScanner::new(repo.ranges, repo.blocks); + scanner.rule_out(repo.ruled_out); + let leaves = scanner.scan(mst_root).unwrap().unwrap().leaves; + (leaves, scanner.take_blocks()) + } + + fn mixed_repo() -> Vec<(String, Vec)> { + (0..2000usize) + .map(|i| { + let collection = match i % 400 { + 0 => "sh.tangled.repo", + 200 => "sh.tangled.repo.collaborator", + n if n % 3 == 0 => "app.bsky.feed.post", + _ => "app.bsky.feed.like", + }; + (format!("{collection}/{i:013}"), record(collection, i)) + }) + .collect() + } + + #[test] + fn takes_every_block_when_the_filter_takes_every_collection() { + assert!(FilteredCar::new(filter(&[]), true).is_none()); + } + + #[tokio::test] + async fn keeps_only_the_commit_wanted_records_and_nodes_that_can_reach_them() { + let records = mixed_repo(); + let (car, mst_root) = repo_car(&records).await; + let wanted = ["sh.tangled.repo", "sh.tangled.repo.collaborator"]; + let want = records + .iter() + .filter(|(key, _)| wanted.iter().any(|c| key.split('/').next() == Some(*c))) + .map(|(key, body)| (SmolStr::new(key), compute_cid(body).unwrap())) + .collect::>() + .into_iter() + .collect::>(); + assert_eq!(want.len(), 10); + + for chunk in [1, 7, 4096, car.len()] { + let repo = stream(&car, &wanted, chunk); + let stats = repo.stats; + assert!( + repo.blocks.get(&repo.root).is_some(), + "chunk {chunk}: commit kept" + ); + // the records here are tiny, so nodes are most of the car and this only guards + // against keeping all of it + assert!( + stats.kept_bytes * 3 < stats.bytes, + "chunk {chunk}: kept {} of {} bytes", + stats.kept_bytes, + stats.bytes + ); + + let (leaves, blocks) = scan(repo, mst_root); + + assert_eq!(leaves, want, "chunk {chunk}"); + for (key, cid) in &leaves { + let body = &records.iter().find(|(k, _)| k == key).unwrap().1; + assert_eq!( + blocks.get(cid).unwrap().as_ref(), + body.as_slice(), + "chunk {chunk}" + ); + } + for (key, body) in records + .iter() + .filter(|(key, _)| key.starts_with("app.bsky.")) + { + let cid = compute_cid(body).unwrap(); + assert!(blocks.get(&cid).is_none(), "chunk {chunk}: kept {key}"); + } + } + } + + #[tokio::test] + async fn a_repo_without_wanted_records_keeps_none_of_them() { + let records = (0..2000usize) + .map(|i| { + ( + format!("io.atcr.hold.layer/{i:013}"), + record("io.atcr.hold.layer", i), + ) + }) + .collect::>(); + let (car, mst_root) = repo_car(&records).await; + + let repo = stream(&car, &["sh.tangled.repo"], 64); + let stats = repo.stats; + let (leaves, blocks) = scan(repo, mst_root); + + assert!(leaves.is_empty()); + for (key, body) in &records { + assert!( + blocks.get(&compute_cid(body).unwrap()).is_none(), + "kept {key}" + ); + } + assert!( + stats.kept_bytes * 3 < stats.bytes, + "kept {} of {} bytes", + stats.kept_bytes, + stats.bytes + ); + } + + #[tokio::test] + async fn a_record_filed_under_another_type_is_left_for_the_scan_to_find_missing() { + let mut records = mixed_repo(); + let odd = record("app.bsky.feed.post", 9999); + records.push(("sh.tangled.repo/odd".to_string(), odd.clone())); + let (car, mst_root) = repo_car(&records).await; + + let (leaves, blocks) = scan(stream(&car, &["sh.tangled.repo"], 512), mst_root); + + let odd_cid = compute_cid(&odd).unwrap(); + assert!(leaves.contains(&(SmolStr::new("sh.tangled.repo/odd"), odd_cid))); + assert!(blocks.get(&odd_cid).is_none()); + } +} diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index 0a1b09f..12a02ce 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -1,6 +1,7 @@ mod admission; pub mod client; pub mod error; +pub(crate) mod filtered_car; pub mod manager; pub mod sparse; pub mod worker; diff --git a/src/backfill/sparse.rs b/src/backfill/sparse.rs index 88e18eb..c07a3a9 100644 --- a/src/backfill/sparse.rs +++ b/src/backfill/sparse.rs @@ -393,7 +393,7 @@ pub(crate) fn check_commit_did(commit_did: &str, did: &Did) -> Result<(), Backfi Ok(()) } -async fn fetch_blocks( +pub(crate) async fn fetch_blocks( http: &ThrottledHttpClient, pds: &url::Url, did: &Did, diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index 8d4e54d..58c2468 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -11,20 +11,24 @@ use tracing::{debug, trace, warn}; use jacquard_api::com_atproto::sync::get_repo::GetRepoError; use jacquard_common::IntoStatic; -use jacquard_common::types::cid::Cid as AtCid; +use jacquard_common::types::cid::{Cid as AtCid, IpldCid}; use jacquard_common::types::did::Did; use crate::backfill::admission::BackfillAdmission; use crate::backfill::client::{ThrottledHttpClient, collect_body_bounded}; use crate::backfill::error::BackfillError; -use crate::backfill::sparse::{SparseBackfillResult, check_commit_did, process_did_sparse}; +use crate::backfill::filtered_car::{FilteredCar, FilteredRepo}; +use crate::backfill::sparse::{ + SparseBackfillResult, check_commit_did, fetch_blocks, process_did_sparse, +}; +use crate::car::{Blocks, CarBlock}; use crate::config::{BackfillStrategy, RateTier}; use crate::db::types::{DbAction, DbRkey}; use crate::db::{self, Txn as DbTxn, keys}; use crate::filter::FilterMode; use crate::net::PublicEndpoint; use crate::ops; -use crate::sparse_mst::sparse_probe_collection; +use crate::sparse_mst::{SparseScanner, sparse_probe_collection}; use crate::state::AppState; use crate::types::{Commit, GaugeState, RepoState, RepoStatus, ResyncState}; use crate::util::{parse_retry_after, throttle::ThrottleHandle}; @@ -191,21 +195,25 @@ pub(super) async fn process_did( } } - // 2. fetch repo (car) + // 2. fetch repo (car). one filter snapshot serves both the fetch and the insert below, so a + // filter change in between can't delete stored records the fetch never kept let start = Instant::now(); + let filter = app_state.filter.load_full(); let throttle = app_state.throttler.get_handle(&pds).await; let tier = app_state.resolve_pds_tier(pds.host_str().unwrap_or("")); - let car_bytes = match fetch_full_repo_car( + let car = match fetch_full_repo_car( http, &pds, did, &throttle, &tier, app_state.max_car_body_bytes, + FilteredCar::new(filter.clone(), app_state.verify_cids), ) .await? { - FullRepoOutcome::Car(body) => body, + FullRepoOutcome::Car(body) => FetchedCar::Whole(body), + FullRepoOutcome::Filtered(repo) => FetchedCar::Filtered(repo), FullRepoOutcome::NotFound => { // `RepoNotFound` is not lifecycle evidence: it says the pds has no repo, not that the // account was deleted, and the repo may only be missing on this target. erasure stays @@ -263,31 +271,50 @@ pub(super) async fn process_did( }; drop(admission_permit); - trace!( - bytes = car_bytes.len(), - elapsed = ?start.elapsed(), - "fetched car bytes" - ); - // 3. import repo - let start = Instant::now(); - let verify_cids = app_state.verify_cids; - let parsed = tokio::task::spawn_blocking(move || crate::car::parse_car(car_bytes, verify_cids)) - .await - .into_diagnostic()??; - trace!( - elapsed = %start.elapsed().as_secs_f32(), - blocks = parsed.blocks.map().len(), - verify_cids, - "parsed car" - ); + let car = match car { + FetchedCar::Whole(car_bytes) => { + trace!( + bytes = car_bytes.len(), + elapsed = ?start.elapsed(), + "fetched car bytes" + ); + let start = Instant::now(); + let verify_cids = app_state.verify_cids; + let parsed = + tokio::task::spawn_blocking(move || crate::car::parse_car(car_bytes, verify_cids)) + .await + .into_diagnostic()??; + trace!( + elapsed = %start.elapsed().as_secs_f32(), + blocks = parsed.blocks.map().len(), + verify_cids, + "parsed car" + ); + ParsedRepo::Whole(parsed) + } + FetchedCar::Filtered(repo) => { + let stats = repo.stats; + trace!( + bytes = stats.bytes, + blocks = stats.blocks, + kept_bytes = stats.kept_bytes, + kept_blocks = stats.kept_blocks, + ruled_out = repo.ruled_out.len(), + elapsed = ?start.elapsed(), + "streamed car" + ); + ParsedRepo::Filtered(repo) + } + }; // 4. parse root commit to get mst root - let blocks = parsed.blocks; - let root_bytes = blocks - .get(&parsed.root) - .cloned() - .ok_or_else(|| miette::miette!("root block missing from CAR"))?; + let root_bytes = match &car { + ParsedRepo::Whole(parsed) => parsed.blocks.get(&parsed.root), + ParsedRepo::Filtered(repo) => repo.blocks.get(&repo.root), + } + .cloned() + .ok_or_else(|| miette::miette!("root block missing from CAR"))?; let root_commit = jacquard_repo::commit::Commit::::from_cbor(&root_bytes) @@ -313,16 +340,27 @@ pub(super) async fn process_did( // 5. walk mst and fetch every record block, off the runtime let start = Instant::now(); let mst_root = root_commit.data; - let records = tokio::task::spawn_blocking(move || { - let leaves = crate::mst::leaves(blocks.map(), mst_root)?; - let records: Vec<_> = leaves - .into_iter() - .filter_map(|(key, cid)| blocks.block(&cid).map(|block| (key, block))) - .collect(); - Ok::<_, miette::Report>(records) - }) - .await - .into_diagnostic()??; + let records = match car { + ParsedRepo::Whole(parsed) => { + let blocks = parsed.blocks; + tokio::task::spawn_blocking(move || { + let leaves = crate::mst::leaves(blocks.map(), mst_root)?; + let records: Vec<_> = leaves + .into_iter() + .filter_map(|(key, cid)| blocks.block(&cid).map(|block| (key, block))) + .collect(); + Ok::<_, miette::Report>(records) + }) + .await + .into_diagnostic()?? + } + ParsedRepo::Filtered(repo) => { + filtered_records( + app_state, http, did, &pds, &throttle, &tier, *repo, mst_root, + ) + .await? + } + }; trace!(elapsed = %start.elapsed().as_secs_f32(), "walked mst"); // 6. insert records into db @@ -333,7 +371,6 @@ pub(super) async fn process_did( let rev = root_commit.rev; tokio::task::spawn_blocking(move || { - let filter = app_state.filter.load(); let mut count = 0; let mut collection_counts: HashMap = HashMap::new(); let mut txn = DbTxn::new(&app_state.db); @@ -500,6 +537,171 @@ async fn remove_discarded_repo( Ok(()) } +/// a getRepo body, as fetched +enum FetchedCar { + Whole(Bytes), + Filtered(Box), +} + +/// a getRepo body, ready for the mst walk +enum ParsedRepo { + Whole(crate::car::ParsedCar), + Filtered(Box), +} + +/// how many chunks a slow filter can fall behind the download before the download waits +const FILTER_CHUNKS_IN_FLIGHT: usize = 16; + +/// feeds a getRepo body to `car` on a blocking thread, since hashing and decoding every block +/// of a big repo would otherwise hold up the runtime. `None` once the body passes `max_bytes`, +/// counted on decoded bytes like [`collect_body_bounded`] +async fn stream_filtered( + mut resp: reqwest::Response, + max_bytes: usize, + mut car: FilteredCar, +) -> Result, BackfillError> { + if resp + .content_length() + .is_some_and(|len| len > max_bytes as u64) + { + return Ok(None); + } + let (tx, mut rx) = tokio::sync::mpsc::channel::(FILTER_CHUNKS_IN_FLIGHT); + let filter = tokio::task::spawn_blocking(move || { + while let Some(chunk) = rx.blocking_recv() { + car.push(&chunk)?; + } + car.finish() + }); + let mut received = 0usize; + let read = async { + while let Some(chunk) = resp.chunk().await? { + received += chunk.len(); + if received > max_bytes { + return Ok(false); + } + // a closed channel means the filter stopped on an error, which joining it reports + if tx.send(chunk).await.is_err() { + break; + } + } + Ok::<_, reqwest::Error>(true) + } + .await; + drop(tx); + let filtered = filter.await.into_diagnostic()?; + match read.map_err(|e| BackfillError::Transport(e.to_string().into()))? { + false => Ok(None), + true => Ok(Some(filtered?)), + } +} + +/// the records a filtered car's scan reaches. a record whose $type isn't its collection was +/// dropped while the car streamed, so any leaf the scan finds missing is fetched again: by cid +/// if the pds has getBlocks, else from one whole car. leaves still missing after that weren't +/// in the car at all, and are skipped like a whole car's missing leaves are +#[allow(clippy::too_many_arguments)] +async fn filtered_records( + app_state: &AppState, + http: &ThrottledHttpClient, + did: &Did, + pds: &url::Url, + throttle: &ThrottleHandle, + tier: &RateTier, + repo: FilteredRepo, + mst_root: IpldCid, +) -> Result, BackfillError> { + let FilteredRepo { + root, + ranges, + blocks, + ruled_out, + .. + } = repo; + let (leaves, blocks) = tokio::task::spawn_blocking(move || { + let mut scanner = SparseScanner::new(ranges, blocks); + scanner.rule_out(ruled_out); + let scan = scanner + .scan(mst_root)? + .map_err(|missing| miette::miette!("mst node {} missing from car", missing[0]))?; + Ok::<_, miette::Report>((scan.leaves, scanner.take_blocks())) + }) + .await + .into_diagnostic()??; + + let mut missing = leaves + .iter() + .map(|(_, cid)| *cid) + .filter(|cid| blocks.get(cid).is_none()) + .collect::>(); + missing.sort_unstable(); + missing.dedup(); + let mut fetched = Blocks::default(); + let mut whole = None; + if !missing.is_empty() { + debug!( + missing = missing.len(), + "fetching records the filtered car dropped" + ); + match fetch_blocks( + http, + pds, + did, + &missing, + throttle, + tier, + app_state.verify_cids, + app_state.max_car_body_bytes, + ) + .await + { + Ok(blocks) => fetched = blocks, + Err(err) => debug!(%err, "getBlocks failed, fetching the whole car instead"), + } + if missing.iter().any(|cid| fetched.get(cid).is_none()) { + whole = Some(whole_car_blocks(app_state, http, did, pds, throttle, tier, root).await?); + } + } + + Ok(leaves + .into_iter() + .filter_map(|(key, cid)| { + blocks + .block(&cid) + .or_else(|| fetched.block(&cid)) + .or_else(|| whole.as_ref().and_then(|whole| whole.block(&cid))) + .map(|block| (key, block)) + }) + .collect()) +} + +/// the blocks of a whole car for a repo whose filtered car came from commit `root`. a repo that +/// moved on in between is retried, since its records may no longer match the scan +async fn whole_car_blocks( + app_state: &AppState, + http: &ThrottledHttpClient, + did: &Did, + pds: &url::Url, + throttle: &ThrottleHandle, + tier: &RateTier, + root: IpldCid, +) -> Result, BackfillError> { + let max = app_state.max_car_body_bytes; + let FullRepoOutcome::Car(car) = + fetch_full_repo_car(http, pds, did, throttle, tier, max, None).await? + else { + return Err(miette::miette!("repo changed state while its records were refetched").into()); + }; + let verify_cids = app_state.verify_cids; + let parsed = tokio::task::spawn_blocking(move || crate::car::parse_car(car, verify_cids)) + .await + .into_diagnostic()??; + if parsed.root != root { + return Err(miette::miette!("repo moved on while its records were refetched").into()); + } + Ok(parsed.blocks) +} + /// upper bound on a buffered xrpc error body. error responses are tiny json objects; this only /// guards against a peer streaming an unbounded body on the error path. const ERROR_BODY_MAX_BYTES: usize = 64 * 1024; @@ -510,6 +712,8 @@ const ERROR_BODY_MAX_BYTES: usize = 64 * 1024; enum FullRepoOutcome { /// the streamed CAR body, within the configured size ceiling. Car(Bytes), + /// the parts of the car a filtered backfill can use, kept as the body streamed past. + Filtered(Box), /// the PDS has no repository for this DID (`RepoNotFound`). this is a fetch result, not /// lifecycle evidence: callers must not delete the stored snapshot for it. NotFound, @@ -530,6 +734,7 @@ async fn fetch_full_repo_car( throttle: &ThrottleHandle, tier: &RateTier, max_body_bytes: usize, + filtered: Option, ) -> Result { let pds_endpoint = PublicEndpoint::parse_http(pds).map_err(|error| { BackfillError::Generic(miette::miette!("unsafe public PDS endpoint {pds}: {error}")) @@ -580,15 +785,20 @@ async fn fetch_full_repo_car( let status = resp.status(); if status.is_success() { - return match collect_body_bounded(resp, max_body_bytes) - .await - .map_err(|e| BackfillError::Transport(e.to_string().into()))? - { - Some(body) => Ok(FullRepoOutcome::Car(body)), - None => Err(BackfillError::Generic(miette::miette!( - "getRepo response for {did} exceeded max body size of {max_body_bytes} bytes" - ))), + let outcome = match filtered { + Some(car) => stream_filtered(resp, max_body_bytes, car) + .await? + .map(|repo| FullRepoOutcome::Filtered(Box::new(repo))), + None => collect_body_bounded(resp, max_body_bytes) + .await + .map_err(|e| BackfillError::Transport(e.to_string().into()))? + .map(FullRepoOutcome::Car), }; + return outcome.ok_or_else(|| { + BackfillError::Generic(miette::miette!( + "getRepo response for {did} exceeded max body size of {max_body_bytes} bytes" + )) + }); } // non-success: decode the (tiny) xrpc error body to preserve typed error handling. @@ -678,6 +888,7 @@ mod tests { &throttle, &RateTier::trusted(), max_body_bytes, + None, ) .await } diff --git a/src/backfill/worker/task.rs b/src/backfill/worker/task.rs index e00a908..5951c5c 100644 --- a/src/backfill/worker/task.rs +++ b/src/backfill/worker/task.rs @@ -546,14 +546,25 @@ mod tests { commit_did: &'static str, rev: &str, records: &[(&str, &[u8])], + ) -> miette::Result { + let records = records + .iter() + .map(|&(rkey, body)| (format!("{COLLECTION}/{rkey}"), body)) + .collect::>(); + build_car_keyed(commit_did, rev, &records).await + } + + /// a car whose records sit under full `collection/rkey` keys + async fn build_car_keyed( + commit_did: &'static str, + rev: &str, + records: &[(String, &[u8])], ) -> miette::Result { let store = Arc::new(MemoryBlockStore::new()); let mut mst = Mst::new(store.clone()); - for &(rkey, body) in records { + for (key, body) in records { let cid = jacquard_repo::mst::util::compute_cid(body).into_diagnostic()?; - mst.add_mut(&format!("{COLLECTION}/{rkey}"), cid) - .await - .into_diagnostic()?; + mst.add_mut(key, cid).await.into_diagnostic()?; store .put_many([(cid, Bytes::copy_from_slice(body))]) .await @@ -849,6 +860,106 @@ mod tests { Ok(()) } + const WANTED: &str = "sh.tangled.repo"; + + fn tangled_repo(name: &str) -> Vec { + serde_ipld_dagcbor::to_vec(&serde_json::json!({ "$type": WANTED, "name": name })).unwrap() + } + + fn only_collection(fixture: &Fixture, collection: &str) { + let mut filter = crate::filter::FilterConfig::new(crate::filter::FilterMode::Filter); + filter.collections = vec![collection.into()]; + fixture.state.filter.store(Arc::new(filter)); + } + + fn stored(fixture: &Fixture, collection: &str, rkey: &str) -> miette::Result>> { + let key = keys::record_key(&fixture.did, collection, &DbRkey::new(rkey)); + Ok(fixture + .state + .db + .indexer + .record(key) + .into_diagnostic()? + .map(|body| body.to_vec())) + } + + #[tokio::test] + async fn filtered_full_backfill_stores_what_a_whole_car_would_for_the_wanted_collection() + -> miette::Result<()> { + let (core, site) = (tangled_repo("core"), tangled_repo("site")); + let posts = (0..300) + .map(|i| post(&format!("post {i}"))) + .collect::>(); + let mut records = posts + .iter() + .enumerate() + .map(|(i, body)| (format!("{COLLECTION}/3jzfcijpj{i:04}"), body.as_slice())) + .collect::>(); + records.push((format!("{WANTED}/core"), core.as_slice())); + records.push((format!("{WANTED}/site"), site.as_slice())); + let car = build_car_keyed(DID, "3jzfcijpj2z2a", &records).await?; + + let filtered = Fixture::new().await?; + only_collection(&filtered, WANTED); + filtered.seed(&RepoState::backfilling().into_static(), 9, &[])?; + filtered.pds.serve_car(car.bytes.clone()).await; + let whole = Fixture::new().await?; + whole.seed(&RepoState::backfilling().into_static(), 9, &[])?; + whole.pds.serve_car(car.bytes).await; + + assert_eq!( + filtered.run(keys::pending_key(9)).await?, + TaskDisposition::Finished + ); + assert_eq!( + whole.run(keys::pending_key(9)).await?, + TaskDisposition::Finished + ); + + assert_eq!(filtered.record_rkeys()?, vec!["core", "site"]); + assert_eq!(stored(&filtered, WANTED, "core")?, Some(core.clone())); + for rkey in ["core", "site"] { + assert_eq!( + stored(&filtered, WANTED, rkey)?, + stored(&whole, WANTED, rkey)? + ); + } + assert_eq!( + filtered.repo()?.root.map(|root| root.data), + whole.repo()?.root.map(|root| root.data) + ); + assert_eq!(filtered.pds.requests(), 1); + Ok(()) + } + + #[tokio::test] + async fn record_filed_under_another_type_comes_back_from_a_whole_car() -> miette::Result<()> { + // the stub has no getBlocks, like the atcr hold, so the dropped record can only come + // from a second, whole car + let odd = post("filed under the wrong collection"); + let core = tangled_repo("core"); + let records = [ + (format!("{WANTED}/core"), core.as_slice()), + (format!("{WANTED}/odd"), odd.as_slice()), + (format!("{COLLECTION}/3jzfcijpj0001"), odd.as_slice()), + ]; + let car = build_car_keyed(DID, "3jzfcijpj2z2a", &records).await?; + let fixture = Fixture::new().await?; + only_collection(&fixture, WANTED); + fixture.seed(&RepoState::backfilling().into_static(), 9, &[])?; + fixture.pds.serve_car(car.bytes).await; + + assert_eq!( + fixture.run(keys::pending_key(9)).await?, + TaskDisposition::Finished + ); + + assert_eq!(fixture.record_rkeys()?, vec!["core", "odd"]); + assert_eq!(stored(&fixture, WANTED, "odd")?, Some(odd)); + assert_eq!(fixture.pds.requests(), 2); + Ok(()) + } + #[tokio::test] async fn repo_not_found_keeps_a_rooted_snapshot_and_parks_a_retry() -> miette::Result<()> { let fixture = Fixture::new().await?; diff --git a/src/bin/backfill_strategy_bench.rs b/src/bin/backfill_strategy_bench.rs index c38ed17..27e45a5 100644 --- a/src/bin/backfill_strategy_bench.rs +++ b/src/bin/backfill_strategy_bench.rs @@ -5,6 +5,7 @@ mod car; #[allow(dead_code)] mod mst; #[path = "../sparse_mst.rs"] +#[allow(dead_code)] mod sparse_mst; use std::io::Write; diff --git a/src/car.rs b/src/car.rs index e14da9f..677fe66 100644 --- a/src/car.rs +++ b/src/car.rs @@ -1,7 +1,7 @@ use std::collections::{BTreeMap, HashMap}; use std::io::Cursor; -use bytes::Bytes; +use bytes::{Buf, Bytes, BytesMut}; use cid::Cid as IpldCid; use miette::{IntoDiagnostic, Result}; @@ -207,14 +207,19 @@ fn ordered_blocks( fn car_root(data: &[u8]) -> Result<(IpldCid, usize)> { let (header_start, blocks_start) = car_header_bounds(data)?; - let header = - iroh_car::CarHeader::decode(&data[header_start..blocks_start]).into_diagnostic()?; - let root = header + Ok(( + header_root(&data[header_start..blocks_start])?, + blocks_start, + )) +} + +fn header_root(header: &[u8]) -> Result { + let header = iroh_car::CarHeader::decode(header).into_diagnostic()?; + header .roots() .first() .copied() - .ok_or_else(|| miette::miette!("CAR file has no roots"))?; - Ok((root, blocks_start)) + .ok_or_else(|| miette::miette!("CAR file has no roots")) } fn car_header_bounds(data: &[u8]) -> Result<(usize, usize)> { @@ -253,12 +258,7 @@ fn for_each_block( let section = data.slice(offset..section_end); offset = section_end; - let mut cursor = Cursor::new(section.as_ref()); - let cid = IpldCid::read_bytes(&mut cursor).into_diagnostic()?; - let block_start = cursor.position() as usize; - if block_start >= section.len() { - return Err(miette::miette!("CAR block has no payload for {cid}")); - } + let (cid, block_start) = split_section(§ion)?; let payload = section.slice(block_start..); if verify_cids { validate_block_cid(&cid, &payload)?; @@ -269,6 +269,92 @@ fn for_each_block( Ok(()) } +/// a section's cid, and where its payload starts +fn split_section(section: &[u8]) -> Result<(IpldCid, usize)> { + let mut cursor = Cursor::new(section); + let cid = IpldCid::read_bytes(&mut cursor).into_diagnostic()?; + let block_start = cursor.position() as usize; + if block_start >= section.len() { + return Err(miette::miette!("CAR block has no payload for {cid}")); + } + Ok((cid, block_start)) +} + +/// a car read as its chunks arrive, so a caller can keep a few of its blocks without the whole +/// file ever being in memory. besides what the caller copies out it holds one unfinished +/// section at most. blocks come out in file order, hashed against their cids when +/// `verify_cids` is set, the same as [`parse_car`]. +#[cfg_attr(not(feature = "indexer"), allow(dead_code))] +pub(crate) struct CarStream { + pending: BytesMut, + root: Option, + verify_cids: bool, +} + +#[cfg_attr(not(feature = "indexer"), allow(dead_code))] +impl CarStream { + pub(crate) fn new(verify_cids: bool) -> Self { + Self { + pending: BytesMut::new(), + root: None, + verify_cids, + } + } + + /// adds the next chunk and hands each block it completes to `on_block`, with the car's + /// root, since the header always comes first + pub(crate) fn push( + &mut self, + chunk: &[u8], + mut on_block: impl FnMut(&IpldCid, IpldCid, &[u8]) -> Result<()>, + ) -> Result<()> { + self.pending.extend_from_slice(chunk); + let mut consumed = 0; + while let Some((len, len_bytes)) = try_read_uvarint(&self.pending[consumed..])? { + let start = consumed + len_bytes; + let end = start + .checked_add(len) + .ok_or_else(|| miette::miette!("CAR section length overflow"))?; + if end > self.pending.len() { + break; + } + let section = &self.pending[start..end]; + match self.root { + None => self.root = Some(header_root(section)?), + Some(root) => { + let (cid, block_start) = split_section(section)?; + let payload = §ion[block_start..]; + if self.verify_cids { + validate_block_cid(&cid, payload)?; + } + on_block(&root, cid, payload)?; + } + } + consumed = end; + } + self.pending.advance(consumed); + Ok(()) + } + + /// ends the car, with `kept` holding the blocks the caller copied out of it. a car that + /// stops partway through its header or a block is truncated, not just short + pub(crate) fn finish_with(self, kept: CommitBlocks) -> Result { + let root = match (self.root, self.pending.is_empty()) { + (None, true) => return Err(miette::miette!("empty CAR file")), + (None, false) => return Err(miette::miette!("truncated CAR header")), + (Some(_), false) => return Err(miette::miette!("truncated CAR block")), + (Some(root), true) => root, + }; + Ok(ParsedCar { + root, + blocks: Blocks { + map: kept, + cids_checked: self.verify_cids, + }, + }) + } +} + /// checks that a block's payload hashes to its claimed CID (sha2-256, dag-cbor). fn validate_block_cid(claimed_cid: &IpldCid, bytes: &[u8]) -> Result<()> { let computed_cid = jacquard_repo::mst::util::compute_cid(bytes).into_diagnostic()?; @@ -280,28 +366,31 @@ fn validate_block_cid(claimed_cid: &IpldCid, bytes: &[u8]) -> Result<()> { Ok(()) } -fn read_uvarint(data: &[u8], offset: &mut usize) -> Result> { - if *offset == data.len() { - return Ok(None); - } - +/// a uvarint at the start of `data` and how many bytes it took, or `None` while `data` ends +/// before it does +fn try_read_uvarint(data: &[u8]) -> Result> { let mut value = 0u64; - for shift in (0..64).step_by(7) { - if *offset >= data.len() { - return Err(miette::miette!("truncated uvarint in CAR file")); - } - - let byte = data[*offset]; - *offset += 1; - value |= u64::from(byte & 0x7f) << shift; - + for (idx, byte) in data.iter().take(10).enumerate() { + value |= u64::from(byte & 0x7f) << (7 * idx); if byte & 0x80 == 0 { - let len = usize::try_from(value).into_diagnostic()?; - return Ok(Some(len)); + return Ok(Some((usize::try_from(value).into_diagnostic()?, idx + 1))); } } + if data.len() >= 10 { + return Err(miette::miette!("uvarint overflow in CAR file")); + } + Ok(None) +} - Err(miette::miette!("uvarint overflow in CAR file")) +fn read_uvarint(data: &[u8], offset: &mut usize) -> Result> { + let rest = &data[(*offset).min(data.len())..]; + if rest.is_empty() { + return Ok(None); + } + let (value, len) = + try_read_uvarint(rest)?.ok_or_else(|| miette::miette!("truncated uvarint in CAR file"))?; + *offset += len; + Ok(Some(value)) } #[cfg(test)] @@ -362,6 +451,148 @@ mod tests { buf.into() } + /// a rooted car whose blocks are big enough for some sections to need a two-byte length + async fn stream_fixture() -> (IpldCid, Vec<(IpldCid, Vec)>, Bytes) { + let blocks = [ + b"a".to_vec(), + vec![7u8; 300], + b"record".to_vec(), + vec![1u8; 129], + ] + .into_iter() + .map(|bytes| { + ( + jacquard_repo::mst::util::compute_cid(&bytes).unwrap(), + bytes, + ) + }) + .collect::>(); + let root = blocks[0].0; + let mut buf = Vec::new(); + let mut writer = + iroh_car::CarWriter::new(iroh_car::CarHeader::new_v1(vec![root]), &mut buf); + for (cid, bytes) in &blocks { + writer.write(*cid, bytes).await.unwrap(); + } + writer.finish().await.unwrap(); + (root, blocks, buf.into()) + } + + fn stream_in_chunks( + car: &[u8], + chunk: usize, + verify_cids: bool, + ) -> Result<(CommitCar, Vec)> { + let mut stream = CarStream::new(verify_cids); + let mut kept = CommitBlocks::new(); + let mut order = Vec::new(); + for piece in car.chunks(chunk) { + stream.push(piece, |_, cid, bytes| { + order.push(cid); + kept.insert(cid, Bytes::copy_from_slice(bytes)); + Ok(()) + })?; + } + Ok((stream.finish_with(kept)?, order)) + } + + #[tokio::test] + async fn stream_yields_what_a_whole_parse_does_however_the_car_is_chunked() { + let (root, blocks, car) = stream_fixture().await; + let whole = parse_car(car.clone(), true).unwrap(); + for chunk in [1, 2, 3, 7, 64, 333, car.len()] { + let (streamed, order) = stream_in_chunks(&car, chunk, true).unwrap(); + + assert_eq!(streamed.root, root, "chunk {chunk}"); + assert_eq!( + order, + blocks.iter().map(|(cid, _)| *cid).collect::>(), + "chunk {chunk}" + ); + for (cid, bytes) in &blocks { + assert_eq!( + streamed.blocks.get(cid).unwrap().as_ref(), + bytes.as_slice(), + "chunk {chunk}" + ); + assert_eq!( + whole.blocks.get(cid), + streamed.blocks.get(cid), + "chunk {chunk}" + ); + } + assert!(streamed.blocks.block(&root).unwrap().verify().is_ok()); + } + } + + #[tokio::test] + async fn stream_cut_short_is_truncated() { + let (_, blocks, car) = stream_fixture().await; + let err = stream_in_chunks(&[], 1, true).unwrap_err(); + assert!(err.to_string().contains("empty CAR file"), "{err}"); + + // cutting exactly between sections leaves a complete, shorter car + let mut boundaries = Vec::new(); + let mut offset = 0; + while let Some(len) = read_uvarint(&car, &mut offset).unwrap() { + offset += len; + boundaries.push(offset); + } + assert_eq!(boundaries.len(), blocks.len() + 1); + for cut in 1..car.len() { + let result = stream_in_chunks(&car[..cut], 5, true); + if let Some(sections) = boundaries.iter().position(|&end| end == cut) { + let (streamed, order) = result.unwrap(); + assert_eq!(order.len(), sections, "cut at {cut}"); + assert_eq!(streamed.blocks.map().len(), sections, "cut at {cut}"); + continue; + } + let err = result + .err() + .unwrap_or_else(|| panic!("a car cut at {cut} parsed")); + let want = if cut < boundaries[0] { + "truncated CAR header" + } else { + "truncated CAR block" + }; + assert!(err.to_string().contains(want), "cut at {cut}: {err}"); + } + } + + #[tokio::test] + async fn stream_rejects_a_forged_block_when_verifying() { + let real = jacquard_repo::mst::util::compute_cid(b"trusted").unwrap(); + let car = car_with_blocks(&[(real, b"trusted"), (cid(1), b"forged")]).await; + + let err = stream_in_chunks(&car, 3, true).unwrap_err(); + assert!(err.to_string().contains("CAR block CID mismatch"), "{err}"); + let (unverified, _) = stream_in_chunks(&car, 3, false).unwrap(); + assert_eq!(unverified.blocks.get(&cid(1)).unwrap().as_ref(), b"forged"); + assert!(unverified.blocks.block(&cid(1)).unwrap().verify().is_err()); + } + + #[tokio::test] + async fn stream_rejects_a_rootless_car() { + let mut buf = Vec::new(); + let mut writer = + iroh_car::CarWriter::new(iroh_car::CarHeader::new_v1(Vec::new()), &mut buf); + writer.write(cid(1), b"block".to_vec()).await.unwrap(); + writer.finish().await.unwrap(); + + let streamed = stream_in_chunks(&buf, 4, false).unwrap_err(); + let whole = parse_car(buf.into(), false).unwrap_err(); + assert_eq!(streamed.to_string(), whole.to_string()); + } + + #[test] + fn partial_uvarint_waits_for_more_bytes() { + assert_eq!(try_read_uvarint(&[]).unwrap(), None); + assert_eq!(try_read_uvarint(&[0x80]).unwrap(), None); + assert_eq!(try_read_uvarint(&[0xac, 0x02]).unwrap(), Some((300, 2))); + assert_eq!(try_read_uvarint(&[0x05, 0xff]).unwrap(), Some((5, 1))); + assert!(try_read_uvarint(&[0x80; 10]).is_err()); + } + #[tokio::test] async fn verifying_parse_rejects_forged_block() { let real = jacquard_repo::mst::util::compute_cid(b"trusted").unwrap(); diff --git a/src/sparse_mst.rs b/src/sparse_mst.rs index b26129a..6ee1071 100644 --- a/src/sparse_mst.rs +++ b/src/sparse_mst.rs @@ -1,4 +1,4 @@ -use std::collections::BTreeMap; +use std::collections::{BTreeMap, HashSet}; use cid::Cid as IpldCid; use jacquard_repo::mst::util::layer_for_key; @@ -141,9 +141,57 @@ impl NodeBounds { } } +/// whether a node seen on its own, its place in the tree unknown, could hold a key in +/// `ranges`: a leaf in one, or a subtree that could reach one whatever the node's outer bounds +/// turn out to be. a scan of `ranges` takes nothing from a node this rules out, from any +/// ancestor, because real bounds only narrow the intervals checked here +pub(crate) fn node_may_reach(ranges: &[KeyRange], entries: &[FlatEntry]) -> bool { + entries.iter().enumerate().any(|(idx, entry)| match entry { + FlatEntry::Leaf { key, .. } => ranges.iter().any(|range| range.contains(key)), + FlatEntry::Tree { .. } => { + let (lower, upper) = (previous_leaf(entries, idx), next_leaf(entries, idx)); + ranges + .iter() + .any(|range| range.intersects_subtree(lower, upper)) + } + }) +} + +/// mst nodes [`node_may_reach`] ruled out, kept as bare sha-256 digests since a car can rule +/// out most of its nodes and a whole cid is several times bigger. a digest names the same +/// bytes whatever codec a cid puts in front of it +#[derive(Debug, Default)] +pub(crate) struct RuledOut(HashSet<[u8; 32], ahash::RandomState>); + +impl RuledOut { + /// remembers `cid` as ruled out. false for a cid that isn't sha-256, which can't be ruled out + pub(crate) fn insert(&mut self, cid: &IpldCid) -> bool { + sha256_digest(cid) + .map(|digest| self.0.insert(digest)) + .is_some() + } + + pub(crate) fn contains(&self, cid: &IpldCid) -> bool { + sha256_digest(cid).is_some_and(|digest| self.0.contains(&digest)) + } + + pub(crate) fn len(&self) -> usize { + self.0.len() + } +} + +fn sha256_digest(cid: &IpldCid) -> Option<[u8; 32]> { + let hash = cid.hash(); + if hash.code() != jacquard_common::types::crypto::SHA2_256 { + return None; + } + hash.digest().try_into().ok() +} + pub(crate) struct SparseScanner { ranges: Vec, blocks: Blocks, + ruled_out: RuledOut, root: Option, pending_missing: BTreeMap>, visited: BTreeMap>, @@ -156,6 +204,7 @@ impl SparseScanner { Self { ranges, blocks, + ruled_out: RuledOut::default(), root: None, pending_missing: BTreeMap::new(), visited: BTreeMap::new(), @@ -168,6 +217,12 @@ impl SparseScanner { self.blocks.extend(blocks); } + /// nodes a streamed car already ruled out, which the scan then reads as empty instead of + /// missing + pub(crate) fn rule_out(&mut self, ruled_out: RuledOut) { + self.ruled_out = ruled_out; + } + pub(crate) fn scan(&mut self, root: IpldCid) -> Result>> { if let Some(existing_root) = self.root { if existing_root != root { @@ -225,6 +280,9 @@ impl SparseScanner { continue; } + if self.ruled_out.contains(&cid) { + continue; + } let Some(bytes) = self.blocks.get(&cid) else { let pending = self.pending_missing.entry(cid).or_default(); if !pending.iter().any(|prior| prior.covers(&bounds)) { @@ -382,6 +440,103 @@ mod tests { assert_eq!(scan.node_blocks_seen, 2); } + #[test] + fn node_may_reach_checks_leaves_and_every_subtree_interval() { + let wanted = sparse_ranges(&[SmolStr::new("mm.x")]); + let leaf = |key: &str| FlatEntry::Leaf { + key: SmolStr::new(key), + cid: cid(9), + }; + let tree = FlatEntry::Tree { cid: cid(8) }; + + assert!(node_may_reach(&wanted, &[leaf("mm.x/1")])); + assert!(!node_may_reach(&wanted, &[leaf("aa.a/1"), leaf("zz.z/1")])); + assert!(!node_may_reach(&wanted, &[])); + // between two keys the interval is known exactly + assert!(node_may_reach( + &wanted, + &[leaf("aa.a/1"), tree.clone(), leaf("zz.z/1")] + )); + assert!(!node_may_reach( + &wanted, + &[leaf("aa.a/1"), tree.clone(), leaf("bb.b/1")] + )); + // at either edge the node's real bound is unknown, so it could reach anything past it + assert!(node_may_reach(&wanted, &[tree.clone(), leaf("zz.z/1")])); + assert!(!node_may_reach(&wanted, &[tree.clone(), leaf("aa.a/1")])); + assert!(node_may_reach(&wanted, &[leaf("aa.a/1"), tree.clone()])); + assert!(!node_may_reach(&wanted, &[leaf("zz.z/1"), tree])); + } + + #[tokio::test] + async fn scan_over_kept_nodes_finds_every_wanted_leaf() { + use jacquard_repo::{MemoryBlockStore, Mst}; + + let mut mst = Mst::new(std::sync::Arc::new(MemoryBlockStore::new())); + for i in 0..3000usize { + // mostly likes and posts with a few records of everything else, the way a real + // account's repo looks + let collection = match i % 300 { + 0 => "sh.tangled.repo", + 150 => "sh.tangled.repo.collaborator", + 7 | 107 | 207 => "io.atcr.hold.layer", + n if n % 4 == 0 => "app.bsky.feed.post", + n if n % 20 == 1 => "zz.example.record", + _ => "app.bsky.feed.like", + }; + let key = format!("{collection}/{i:013}"); + mst = mst.add(&key, cid((i % 251) as u8)).await.unwrap(); + } + let (root, nodes) = mst.collect_blocks().await.unwrap(); + let all: crate::car::CarBlocks = nodes.into_iter().collect(); + let every_leaf = crate::mst::leaves(&all, root).unwrap(); + + // whether the patterns want few enough keys that the leaf-level nodes, most of a tree, + // should all be ruled out + for (patterns, selective) in [ + ( + vec!["sh.tangled.repo", "sh.tangled.repo.collaborator"], + true, + ), + (vec!["sh.tangled.*"], true), + (vec!["app.bsky.feed.post"], false), + (vec!["io.atcr.*", "zz.example.record"], true), + (vec!["nothing.here"], true), + (vec!["app.bsky.*", "sh.tangled.repo"], false), + ] { + let ranges = sparse_ranges(&patterns.iter().map(SmolStr::new).collect::>()); + let mut kept = BTreeMap::new(); + let mut ruled_out = RuledOut::default(); + for (cid, bytes) in &all { + if node_may_reach(&ranges, &decode_node(bytes).unwrap()) { + kept.insert(*cid, bytes.clone()); + } else { + assert!(ruled_out.insert(cid)); + } + } + let mut scanner = SparseScanner::new(ranges.clone(), Blocks::unverified(kept.clone())); + scanner.rule_out(ruled_out); + + let scan = scanner + .scan(root) + .unwrap() + .unwrap_or_else(|missing| panic!("{patterns:?}: missing {missing:?}")); + + let want = every_leaf + .iter() + .filter(|(key, _)| ranges.iter().any(|range| range.contains(key))) + .cloned() + .collect::>(); + assert_eq!(scan.leaves, want, "{patterns:?}"); + assert!( + !selective || kept.len() * 3 < all.len(), + "{patterns:?}: kept {} of {} nodes", + kept.len(), + all.len() + ); + } + } + #[test] fn estimates_node_layer_from_first_leaf() { let root_node = node(vec![("sh.tangled.repo/1", cid(1), None)], None); -- 2.51.2