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);