diff --git a/src/backfill/sparse.rs b/src/backfill/sparse.rs index 258d7f0..354132e 100644 --- a/src/backfill/sparse.rs +++ b/src/backfill/sparse.rs @@ -396,12 +396,8 @@ async fn fetch_block_chunk( "getBlocks response for {did} exceeded max body size of {max_body_bytes} bytes" ))); } - let blocks = - crate::car::parse_car_blocks(bytes::Bytes::from(body)).map_err(BackfillError::from)?; - if verify_cids { - crate::car::validate_block_cids(&blocks)?; - } - Ok(blocks) + crate::car::parse_car_blocks(bytes::Bytes::from(body), verify_cids) + .map_err(BackfillError::from) } else { let retry_after = (status == StatusCode::TOO_MANY_REQUESTS) .then_some(()) diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index c5e2ef4..ead92cf 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -273,14 +273,11 @@ pub(super) async fn process_did( // 3. import repo let start = Instant::now(); - let parsed = crate::car::parse_car(car_bytes)?; + let parsed = crate::car::parse_car(car_bytes, app_state.verify_cids)?; trace!(elapsed = %start.elapsed().as_secs_f32(), "parsed car"); let start = Instant::now(); let root_cid = parsed.root; - if app_state.verify_cids { - crate::car::validate_block_cids(&parsed.blocks)?; - } let store = Arc::new(MemoryBlockStore::new_from_blocks(parsed.blocks)); trace!( blocks = store.len(), diff --git a/src/bin/backfill_strategy_bench.rs b/src/bin/backfill_strategy_bench.rs index e618ac9..f533cf6 100644 --- a/src/bin/backfill_strategy_bench.rs +++ b/src/bin/backfill_strategy_bench.rs @@ -467,7 +467,7 @@ async fn fetch_block_chunk( .into_output() .map_err(|err: XrpcError<_>| miette::miette!("getBlocks failed for {did}: {err}"))?; let bytes = car.body.len(); - let parsed = car::parse_car_blocks(car.body).wrap_err_with(|| { + let parsed = car::parse_car_blocks(car.body, false).wrap_err_with(|| { let cids = cids .iter() .map(ToString::to_string) diff --git a/src/car.rs b/src/car.rs index b7ac940..0c4065e 100644 --- a/src/car.rs +++ b/src/car.rs @@ -13,7 +13,7 @@ pub(crate) struct ParsedCar { } #[cfg_attr(not(feature = "indexer"), allow(dead_code))] -pub(crate) fn parse_car(data: Bytes) -> Result { +pub(crate) fn parse_car(data: Bytes, verify_cids: bool) -> Result { let (header_start, blocks_start) = car_header_bounds(&data)?; let header = iroh_car::CarHeader::decode(&data[header_start..blocks_start]).into_diagnostic()?; @@ -22,15 +22,15 @@ pub(crate) fn parse_car(data: Bytes) -> Result { .first() .copied() .ok_or_else(|| miette::miette!("CAR file has no roots"))?; - let blocks = parse_car_blocks_from(data, blocks_start)?; + let blocks = parse_car_blocks_from(data, blocks_start, verify_cids)?; Ok(ParsedCar { root, blocks }) } #[cfg_attr(not(feature = "indexer"), allow(dead_code))] -pub(crate) fn parse_car_blocks(data: Bytes) -> Result> { +pub(crate) fn parse_car_blocks(data: Bytes, verify_cids: bool) -> Result> { let (_, blocks_start) = car_header_bounds(&data)?; - parse_car_blocks_from(data, blocks_start) + parse_car_blocks_from(data, blocks_start, verify_cids) } fn car_header_bounds(data: &[u8]) -> Result<(usize, usize)> { @@ -49,7 +49,13 @@ fn car_header_bounds(data: &[u8]) -> Result<(usize, usize)> { Ok((header_start, header_end)) } -fn parse_car_blocks_from(data: Bytes, mut offset: usize) -> Result> { +/// with `verify_cids`, every section is hashed as it is read, while its bytes are still in +/// cache and in file order; checking the finished map afterwards walks memory in cid order. +fn parse_car_blocks_from( + data: Bytes, + mut offset: usize, + verify_cids: bool, +) -> Result> { let mut blocks = BTreeMap::new(); while let Some(section_len) = read_uvarint(&data, &mut offset)? { let section_end = offset @@ -68,7 +74,11 @@ fn parse_car_blocks_from(data: Bytes, mut offset: usize) -> Result= section.len() { return Err(miette::miette!("CAR block has no payload for {cid}")); } - blocks.insert(cid, section.slice(block_start..)); + let payload = section.slice(block_start..); + if verify_cids { + validate_block_cid(&cid, &payload)?; + } + blocks.insert(cid, payload); } Ok(blocks) @@ -78,14 +88,17 @@ fn parse_car_blocks_from(data: Bytes, mut offset: usize) -> Result( blocks: impl IntoIterator, ) -> Result<()> { - for (claimed_cid, bytes) in blocks { - let computed_cid = - jacquard_repo::mst::util::compute_cid(bytes.as_ref()).into_diagnostic()?; - if computed_cid != *claimed_cid { - return Err(miette::miette!( - "CAR block CID mismatch: claimed {claimed_cid}, computed {computed_cid}" - )); - } + blocks + .into_iter() + .try_for_each(|(claimed_cid, bytes)| validate_block_cid(claimed_cid, bytes)) +} + +fn validate_block_cid(claimed_cid: &IpldCid, bytes: &[u8]) -> Result<()> { + let computed_cid = jacquard_repo::mst::util::compute_cid(bytes).into_diagnostic()?; + if computed_cid != *claimed_cid { + return Err(miette::miette!( + "CAR block CID mismatch: claimed {claimed_cid}, computed {computed_cid}" + )); } Ok(()) } @@ -136,7 +149,7 @@ mod tests { writer.write(cid, b"block".to_vec()).await.unwrap(); writer.finish().await.unwrap(); - let blocks = parse_car_blocks(buf.into()).unwrap(); + let blocks = parse_car_blocks(buf.into(), false).unwrap(); assert_eq!(blocks.get(&cid).unwrap().as_ref(), b"block"); } @@ -152,7 +165,7 @@ mod tests { let data = Bytes::from(buf); let data_range = data.as_ptr() as usize..data.as_ptr() as usize + data.len(); - let parsed = parse_car(data).unwrap(); + let parsed = parse_car(data, false).unwrap(); let payload = parsed.blocks.get(&block).unwrap(); let payload_range = payload.as_ptr() as usize..payload.as_ptr() as usize + payload.len(); @@ -162,6 +175,51 @@ mod tests { assert!(payload_range.end <= data_range.end); } + async fn car_with_blocks(blocks: &[(IpldCid, &[u8])]) -> Bytes { + let mut buf = Vec::new(); + let header = iroh_car::CarHeader::new_v1(vec![blocks[0].0]); + let mut writer = iroh_car::CarWriter::new(header, &mut buf); + for (cid, bytes) in blocks { + writer.write(*cid, bytes.to_vec()).await.unwrap(); + } + writer.finish().await.unwrap(); + buf.into() + } + + #[tokio::test] + async fn verifying_parse_rejects_forged_block() { + let real = jacquard_repo::mst::util::compute_cid(b"trusted").unwrap(); + let car = car_with_blocks(&[(real, b"trusted"), (cid(1), b"forged")]).await; + + assert!(parse_car(car.clone(), false).is_ok()); + let err = parse_car(car.clone(), true).unwrap_err(); + assert!(err.to_string().contains("CAR block CID mismatch")); + let err = parse_car_blocks(car, true).unwrap_err(); + assert!(err.to_string().contains("CAR block CID mismatch")); + } + + #[tokio::test] + async fn verifying_parse_accepts_matching_blocks() { + let root = jacquard_repo::mst::util::compute_cid(b"root").unwrap(); + let leaf = jacquard_repo::mst::util::compute_cid(b"leaf").unwrap(); + let car = car_with_blocks(&[(root, b"root"), (leaf, b"leaf")]).await; + + let parsed = parse_car(car, true).unwrap(); + assert_eq!(parsed.root, root); + assert_eq!(parsed.blocks.get(&leaf).unwrap().as_ref(), b"leaf"); + } + + #[tokio::test] + async fn verifying_parse_rejects_forged_copy_shadowed_by_valid_duplicate() { + // the map keeps only the last copy, so the forged one is caught only while parsing; + // no honest pds sends it, so the whole car is rejected + let real = jacquard_repo::mst::util::compute_cid(b"trusted").unwrap(); + let car = car_with_blocks(&[(real, b"forged"), (real, b"trusted")]).await; + + let err = parse_car(car, true).unwrap_err(); + assert!(err.to_string().contains("CAR block CID mismatch")); + } + #[test] fn validate_block_cids_rejects_mismatched_cid() { let blocks = BTreeMap::from([(cid(1), Bytes::from_static(b"forged"))]);