diff --git a/benches/repo.rs b/benches/repo.rs index d9d1fd1..839cc94 100644 --- a/benches/repo.rs +++ b/benches/repo.rs @@ -291,8 +291,12 @@ fn record_check(c: &mut Criterion) { group.finish(); } -/// a full backfill without the database: parse, walk, and check every record for storage. -/// verify against trust here is the whole cpu cost of `verify_cids` on a backfill. +/// how big the pieces of a car fed to [`car::WholeCar`] are, about what a getRepo body arrives +/// in +const BODY_CHUNK: usize = 16 * 1024; + +/// a full backfill without the database: stream the car in, walk, and check every record for +/// storage. verify against trust here is the whole cpu cost of `verify_cids` on a backfill. fn backfill_cpu(c: &mut Criterion) { let mut group = group(c, "backfill_cpu"); for repo in repos() { @@ -300,7 +304,12 @@ fn backfill_cpu(c: &mut Criterion) { for (check, verify_cids) in CHECKS { group.bench_function(BenchmarkId::new(check, repo.name), |b| { b.iter_with_large_drop(|| { - let parsed = car::parse_car(repo.car.clone(), verify_cids).unwrap(); + let mut whole = car::WholeCar::new(verify_cids); + whole.reserve(repo.car.len()); + for chunk in repo.car.chunks(BODY_CHUNK) { + whole.push(chunk).unwrap(); + } + let parsed = whole.finish().unwrap(); let leaves = mst::leaves(parsed.blocks.map(), commit_data(&parsed)).unwrap(); for (_, cid) in &leaves { black_box(parsed.blocks.block(cid).unwrap().verify().unwrap()); diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index 5b123dd..26a20b2 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -22,7 +22,7 @@ use crate::backfill::sparse::{ SparseBackfillResult, check_commit_did, fetch_blocks, process_did_sparse, }; use crate::backfill::{Persisted, backfilled_state}; -use crate::car::{Blocks, CarBlock}; +use crate::car::{Blocks, CarBlock, CarBlocks, ParsedCar, WholeCar}; use crate::config::{BackfillStrategy, RateTier}; use crate::db::types::{DbAction, DbRkey}; use crate::db::{self, Txn as DbTxn, keys}; @@ -210,6 +210,10 @@ pub(super) async fn process_did( 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 = FilteredCar::new(filter.clone(), app_state.verify_cids).map_or_else( + || RepoCar::Whole(WholeCar::new(app_state.verify_cids)), + RepoCar::Filtered, + ); let car = match fetch_full_repo_car( http, &pds, @@ -217,12 +221,12 @@ pub(super) async fn process_did( &throttle, &tier, app_state.max_car_body_bytes, - FilteredCar::new(filter.clone(), app_state.verify_cids), + car, ) .await? { - FullRepoOutcome::Car(body) => FetchedCar::Whole(body), - FullRepoOutcome::Filtered(repo) => FetchedCar::Filtered(repo), + FullRepoOutcome::Car(parsed) => ParsedRepo::Whole(parsed), + FullRepoOutcome::Filtered(repo) => ParsedRepo::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 @@ -281,36 +285,20 @@ pub(super) async fn process_did( // a filtered car can still go back to the pds for records it dropped, so it keeps // its admission through the mst walk let admission_permit = match &car { - FetchedCar::Whole(_) => { + ParsedRepo::Whole(_) => { drop(admission_permit); None } - FetchedCar::Filtered(_) => Some(admission_permit), + ParsedRepo::Filtered(_) => Some(admission_permit), }; - // 3. import repo - 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) => { + match &car { + ParsedRepo::Whole(parsed) => trace!( + blocks = parsed.blocks.map().len(), + elapsed = ?start.elapsed(), + "streamed whole car" + ), + ParsedRepo::Filtered(repo) => { let stats = repo.stats; trace!( bytes = stats.bytes, @@ -321,9 +309,8 @@ pub(super) async fn process_did( elapsed = ?start.elapsed(), "streamed car" ); - ParsedRepo::Filtered(repo) } - }; + } // 4. parse root commit to get mst root let root_bytes = match &car { @@ -577,37 +564,54 @@ 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), + Whole(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; +/// what a getRepo body is read into as it streams +enum RepoCar { + Whole(WholeCar), + Filtered(FilteredCar), +} + +impl RepoCar { + fn push(&mut self, chunk: &[u8]) -> Result<()> { + match self { + RepoCar::Whole(car) => car.push(chunk), + RepoCar::Filtered(car) => car.push(chunk), + } + } + + fn finish(self) -> Result { + Ok(match self { + RepoCar::Whole(car) => FullRepoOutcome::Car(car.finish()?), + RepoCar::Filtered(car) => FullRepoOutcome::Filtered(Box::new(car.finish()?)), + }) + } +} + +/// how many chunks the car reader can fall behind the download before the download waits +const 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( +async fn stream_car( 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) - { + mut car: RepoCar, +) -> Result, BackfillError> { + let content_length = resp.content_length(); + if 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 || { + if let (RepoCar::Whole(whole), Some(len)) = (&mut car, content_length) { + whole.reserve(len as usize); + } + let (tx, mut rx) = tokio::sync::mpsc::channel::(CHUNKS_IN_FLIGHT); + let reader = tokio::task::spawn_blocking(move || { while let Some(chunk) = rx.blocking_recv() { car.push(&chunk)?; } @@ -615,12 +619,19 @@ async fn stream_filtered( }); let mut received = 0usize; let read = async { - while let Some(chunk) = resp.chunk().await? { + loop { + // the reader stopping on an error ends the download now, not when the next chunk + // shows up, which a stalled pds may never send + let chunk = tokio::select! { + chunk = resp.chunk() => chunk?, + () = tx.closed() => break, + }; + let Some(chunk) = chunk else { break }; received += chunk.len(); if received > max_bytes { return Ok(false); } - // a closed channel means the filter stopped on an error, which joining it reports + // a closed channel means the reader stopped on an error, which joining it reports if tx.send(chunk).await.is_err() { break; } @@ -629,10 +640,10 @@ async fn stream_filtered( } .await; drop(tx); - let filtered = filter.await.into_diagnostic()?; + let outcome = reader.await.into_diagnostic()?; match read.map_err(|e| BackfillError::Transport(e.to_string().into()))? { false => Ok(None), - true => Ok(Some(filtered?)), + true => Ok(Some(outcome?)), } } @@ -753,8 +764,8 @@ async fn blocks_from_whole_car( cids: HashSet, ) -> Result, BackfillError> { let max = app_state.max_car_body_bytes; - let only = FilteredCar::only(root, cids, app_state.verify_cids); - match fetch_full_repo_car(http, pds, did, throttle, tier, max, Some(only)).await? { + let only = RepoCar::Filtered(FilteredCar::only(root, cids, app_state.verify_cids)); + match fetch_full_repo_car(http, pds, did, throttle, tier, max, only).await? { FullRepoOutcome::Filtered(repo) => Ok(repo.blocks), FullRepoOutcome::NotFound => Err(BackfillError::RepoNotFound), // the retry's first getRepo sees the same status and records it @@ -775,8 +786,8 @@ const ERROR_BODY_MAX_BYTES: usize = 64 * 1024; /// rate-limit failures are surfaced as [`BackfillError`] instead. #[derive(Debug)] enum FullRepoOutcome { - /// the streamed CAR body, within the configured size ceiling. - Car(Bytes), + /// every block of the car, within the configured size ceiling. + Car(ParsedCar), /// 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 @@ -799,7 +810,7 @@ async fn fetch_full_repo_car( throttle: &ThrottleHandle, tier: &RateTier, max_body_bytes: usize, - filtered: Option, + car: RepoCar, ) -> Result { let pds_endpoint = PublicEndpoint::parse_http(pds).map_err(|error| { BackfillError::Generic(miette::miette!("unsafe public PDS endpoint {pds}: {error}")) @@ -850,16 +861,7 @@ async fn fetch_full_repo_car( let status = resp.status(); if status.is_success() { - 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(|| { + return stream_car(resp, max_body_bytes, car).await?.ok_or_else(|| { BackfillError::Generic(miette::miette!( "getRepo response for {did} exceeded max body size of {max_body_bytes} bytes" )) @@ -914,6 +916,13 @@ mod tests { use axum::Router; use axum::routing::get; + async fn serve(app: Router) -> url::Url { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = url::Url::parse(&format!("http://{}/", listener.local_addr().unwrap())).unwrap(); + tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + url + } + /// spawns a server answering `com.atproto.sync.getRepo` with a fixed status and body. async fn spawn_get_repo(status: StatusCode, body: Vec) -> url::Url { let app = Router::new().route( @@ -923,10 +932,7 @@ mod tests { async move { (status, body) } }), ); - let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - let url = url::Url::parse(&format!("http://{}/", listener.local_addr().unwrap())).unwrap(); - tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); - url + serve(app).await } /// spawns a server answering getRepo with `chunks` as a chunked body, then dropping the @@ -946,10 +952,7 @@ mod tests { async move { axum::body::Body::from_stream(parts.chain(cut)) } }), ); - let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - let url = url::Url::parse(&format!("http://{}/", listener.local_addr().unwrap())).unwrap(); - tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); - url + serve(app).await } /// a car of just its commit, or of a commit whose bytes don't match its cid @@ -975,13 +978,13 @@ mod tests { pds: &url::Url, max_body_bytes: usize, ) -> Result { - fetch_with(pds, max_body_bytes, None).await + fetch_with(pds, max_body_bytes, RepoCar::Whole(WholeCar::new(true))).await } async fn fetch_with( pds: &url::Url, max_body_bytes: usize, - filtered: Option, + car: RepoCar, ) -> Result { let address = std::net::SocketAddr::new("127.0.0.1".parse().unwrap(), pds.port().unwrap()); let public_pds = url::Url::parse("http://pds.example/").unwrap(); @@ -1003,7 +1006,7 @@ mod tests { &throttle, &RateTier::trusted(), max_body_bytes, - filtered, + car, ) .await } @@ -1014,7 +1017,10 @@ mod tests { let pds = spawn_streamed_get_repo(car.chunks(7).map(Bytes::copy_from_slice).collect(), false) .await; - match fetch_with(&pds, 1024, Some(tangled_only())).await.unwrap() { + match fetch_with(&pds, 1024, RepoCar::Filtered(tangled_only())) + .await + .unwrap() + { FullRepoOutcome::Filtered(repo) => assert_eq!(repo.stats.kept_blocks, 1), other => panic!("expected Filtered, got {other:?}"), } @@ -1025,7 +1031,7 @@ mod tests { let car = commit_only_car(false).await; let half = car.slice(..car.len() / 2); let pds = spawn_streamed_get_repo(vec![half], true).await; - let err = fetch_with(&pds, 1024, Some(tangled_only())) + let err = fetch_with(&pds, 1024, RepoCar::Filtered(tangled_only())) .await .unwrap_err(); assert!(matches!(err, BackfillError::Transport(_)), "{err}"); @@ -1034,7 +1040,7 @@ mod tests { #[tokio::test] async fn streamed_get_repo_rejects_a_forged_block_without_blaming_the_network() { let pds = spawn_streamed_get_repo(vec![commit_only_car(true).await], false).await; - let err = fetch_with(&pds, 1024, Some(tangled_only())) + let err = fetch_with(&pds, 1024, RepoCar::Filtered(tangled_only())) .await .unwrap_err(); assert!(matches!(err, BackfillError::Generic(_)), "{err}"); @@ -1050,7 +1056,7 @@ mod tests { spawn_streamed_get_repo(car.chunks(7).map(Bytes::copy_from_slice).collect(), false) .await; for pds in [sized, chunked] { - let err = fetch_with(&pds, max, Some(tangled_only())) + let err = fetch_with(&pds, max, RepoCar::Filtered(tangled_only())) .await .unwrap_err(); assert!( @@ -1094,13 +1100,46 @@ mod tests { #[tokio::test] async fn full_get_repo_streams_body_within_ceiling() { - let pds = spawn_get_repo(StatusCode::OK, b"car-bytes".to_vec()).await; + let car = commit_only_car(false).await; + let pds = spawn_get_repo(StatusCode::OK, car.to_vec()).await; + let whole = crate::car::parse_car(car, true).unwrap(); match fetch(&pds, 1024).await.unwrap() { - FullRepoOutcome::Car(body) => assert_eq!(body.as_ref(), b"car-bytes"), + FullRepoOutcome::Car(parsed) => { + assert_eq!(parsed.root, whole.root); + assert_eq!(parsed.blocks.map(), whole.blocks.map()); + } other => panic!("expected Car, got {other:?}"), } } + #[tokio::test] + async fn full_get_repo_checks_blocks_as_they_arrive() { + // a forged block, then a body that never ends: only a reader that checks each block + // as it comes in gets to answer + let car = commit_only_car(true).await; + let app = Router::new().route( + "/xrpc/com.atproto.sync.getRepo", + get(move || { + use futures::StreamExt; + let first = futures::stream::iter([Ok::<_, std::io::Error>(car.clone())]); + async move { axum::body::Body::from_stream(first.chain(futures::stream::pending())) } + }), + ); + let pds = serve(app).await; + + for car in [ + RepoCar::Whole(WholeCar::new(true)), + RepoCar::Filtered(tangled_only()), + ] { + let fetch = fetch_with(&pds, 1024, car); + let fetched = tokio::time::timeout(std::time::Duration::from_secs(5), fetch) + .await + .expect("the forged block ends the fetch before the body does"); + let err = fetched.unwrap_err(); + assert!(err.to_string().contains("CAR block CID mismatch"), "{err}"); + } + } + #[tokio::test] async fn full_get_repo_rejects_oversized_body() { let pds = spawn_get_repo(StatusCode::OK, vec![0u8; 4096]).await; diff --git a/src/backfill/worker/task.rs b/src/backfill/worker/task.rs index 614270e..964a2cf 100644 --- a/src/backfill/worker/task.rs +++ b/src/backfill/worker/task.rs @@ -1184,6 +1184,53 @@ mod tests { Ok(()) } + #[tokio::test] + async fn full_backfill_stores_every_record_a_whole_parse_finds() -> miette::Result<()> { + let core = tangled_repo("core"); + let posts = (0..500) + .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())); + // the same body under two keys is one block in the car + records.push((format!("{WANTED}/copy"), core.as_slice())); + let car = build_car_keyed(DID, "3jzfcijpj2z2a", &records).await?; + + let parsed = crate::car::parse_car(car.bytes.clone(), true)?; + let want = crate::mst::leaves(parsed.blocks.map(), car.mst_root)? + .into_iter() + .map(|(key, cid)| (key.to_string(), parsed.blocks.get(&cid).unwrap().to_vec())) + .collect::>(); + assert_eq!(want.len(), records.len()); + + let fixture = Fixture::new().await?; + 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 + ); + + let prefix = keys::record_prefix_did(&fixture.did); + let stored = fixture + .state + .db + .indexer + .record_prefix(&prefix) + .map(|guard| { + let (key, body) = guard.into_inner().into_diagnostic()?; + let (collection, rkey) = keys::split_record_suffix(&key[prefix.len()..])?; + Ok((format!("{collection}/{rkey}"), body.to_vec())) + }) + .collect::>>()?; + assert_eq!(stored, want); + Ok(()) + } + /// a repo with a record filed under a collection other than its `$type`, which a filtered /// car drops and the scan then finds missing async fn car_with_a_misfiled_record() -> miette::Result<(BuiltCar, Vec)> { diff --git a/src/car.rs b/src/car.rs index dc6152d..3b28397 100644 --- a/src/car.rs +++ b/src/car.rs @@ -169,7 +169,8 @@ impl<'a> VerifiedBlock<'a> { } } -#[cfg_attr(not(feature = "indexer"), allow(dead_code))] +// backfill streams whole cars with WholeCar, so only tests and the benches parse one in a piece +#[cfg_attr(not(test), allow(dead_code))] pub(crate) fn parse_car(data: Bytes, verify_cids: bool) -> Result> { let (root, blocks_start) = car_root(&data)?; let mut map = CarBlocks::with_capacity_and_hasher( @@ -289,9 +290,9 @@ fn split_section(section: &[u8]) -> Result<(IpldCid, usize)> { 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 +/// a car read as its chunks arrive, so a caller can copy out the blocks it wants without the +/// whole file ever being in memory as one piece. 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 { @@ -347,7 +348,7 @@ impl CarStream { /// 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 { + pub(crate) fn finish_with(self, kept: M) -> 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")), @@ -364,6 +365,42 @@ impl CarStream { } } +/// a whole car read as its chunks arrive, every block copied out on its own. the body is never +/// held in one piece, and blocks the mst walk leaves behind are freed with the map instead of +/// pinning the body until the last record is stored +#[cfg_attr(not(feature = "indexer"), allow(dead_code))] +pub(crate) struct WholeCar { + stream: CarStream, + blocks: CarBlocks, +} + +#[cfg_attr(not(feature = "indexer"), allow(dead_code))] +impl WholeCar { + pub(crate) fn new(verify_cids: bool) -> Self { + Self { + stream: CarStream::new(verify_cids), + blocks: CarBlocks::default(), + } + } + + /// makes room for a car of about `car_bytes`, so the map isn't rehashed as it fills + pub(crate) fn reserve(&mut self, car_bytes: usize) { + self.blocks.reserve(car_bytes / EXPECTED_SECTION_BYTES); + } + + pub(crate) fn push(&mut self, chunk: &[u8]) -> Result<()> { + let blocks = &mut self.blocks; + self.stream.push(chunk, |_, cid, bytes| { + blocks.insert(cid, Bytes::copy_from_slice(bytes)); + Ok(()) + }) + } + + pub(crate) fn finish(self) -> Result> { + self.stream.finish_with(self.blocks) + } +} + /// 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()?; @@ -534,6 +571,29 @@ mod tests { } } + #[tokio::test] + async fn whole_car_keeps_every_block_a_whole_parse_does() { + let (root, _, car) = stream_fixture().await; + let parsed = parse_car(car.clone(), true).unwrap(); + for chunk in [1, 3, 64, car.len()] { + let mut whole = WholeCar::new(true); + whole.reserve(car.len()); + for piece in car.chunks(chunk) { + whole.push(piece).unwrap(); + } + let streamed = whole.finish().unwrap(); + assert_eq!(streamed.root, root, "chunk {chunk}"); + assert_eq!(streamed.blocks.map(), parsed.blocks.map(), "chunk {chunk}"); + assert!(streamed.blocks.block(&root).unwrap().cid_checked); + } + + let real = jacquard_repo::mst::util::compute_cid(b"trusted").unwrap(); + let forged = car_with_blocks(&[(real, b"trusted"), (cid(1), b"forged")]).await; + let mut whole = WholeCar::new(true); + let err = whole.push(&forged).unwrap_err(); + assert!(err.to_string().contains("CAR block CID mismatch"), "{err}"); + } + #[tokio::test] async fn stream_cut_short_is_truncated() { let (_, blocks, car) = stream_fixture().await; diff --git a/tests/representation_inventory.tsv b/tests/representation_inventory.tsv index 64a58b4..56521b5 100644 --- a/tests/representation_inventory.tsv +++ b/tests/representation_inventory.tsv @@ -96,6 +96,10 @@ car-codec:src/car.rs::CommitBlocks::get rust:src/car.rs::verifying_parse_accepts car-codec:src/car.rs::HashMap::get rust:src/car.rs::parses_rooted_car_with_zero_copy_payloads car-codec:src/car.rs::VerifiedBlock::bytes rust:src/car.rs::blocks_from_a_verifying_parse_verify_as_parsed car-codec:src/car.rs::VerifiedBlock::cid rust:src/car.rs::blocks_from_a_verifying_parse_verify_as_parsed +car-codec:src/car.rs::WholeCar::finish rust:src/car.rs::whole_car_keeps_every_block_a_whole_parse_does +car-codec:src/car.rs::WholeCar::new rust:src/car.rs::whole_car_keeps_every_block_a_whole_parse_does +car-codec:src/car.rs::WholeCar::push rust:src/car.rs::whole_car_keeps_every_block_a_whole_parse_does +car-codec:src/car.rs::WholeCar::reserve rust:src/car.rs::whole_car_keeps_every_block_a_whole_parse_does car-codec:src/car.rs::car_header_bounds rust:src/car.rs::parses_rooted_car_with_zero_copy_payloads car-codec:src/car.rs::car_root rust:src/car.rs::stream_rejects_a_rootless_car car-codec:src/car.rs::for_each_block rust:src/car.rs::verifying_parse_rejects_forged_copy_shadowed_by_valid_duplicate