diff --git a/src/backfill/sparse.rs b/src/backfill/sparse.rs index 354132e..0f07ba0 100644 --- a/src/backfill/sparse.rs +++ b/src/backfill/sparse.rs @@ -139,12 +139,7 @@ pub(crate) async fn process_did_sparse( Err(e) => Err(e).into_diagnostic()?, }; - let parsed = jacquard_repo::car::reader::parse_car_bytes(&proof.body) - .await - .into_diagnostic()?; - if app_state.verify_cids { - crate::car::validate_block_cids(&parsed.blocks)?; - } + let parsed = crate::car::parse_commit_car(proof.body, app_state.verify_cids)?; let root_bytes = parsed .blocks .get(&parsed.root) diff --git a/src/bin/backfill_strategy_bench.rs b/src/bin/backfill_strategy_bench.rs index 21fc7fe..65f5101 100644 --- a/src/bin/backfill_strategy_bench.rs +++ b/src/bin/backfill_strategy_bench.rs @@ -354,15 +354,8 @@ async fn bench_sparse( .into_output() .map_err(|err| miette::miette!("getRecord seed failed for {did}: {err}"))?; let seed_bytes = seed.body.len(); - let parsed = jacquard_repo::car::reader::parse_car_bytes(&seed.body) - .await - .into_diagnostic() - .wrap_err_with(|| { - format!( - "parse getRecord seed CAR for {did} ({} bytes)", - seed.body.len() - ) - })?; + let parsed = car::parse_commit_car(seed.body, false) + .wrap_err_with(|| format!("parse getRecord seed CAR for {did} ({seed_bytes} bytes)"))?; let root_bytes = parsed .blocks .get(&parsed.root) diff --git a/src/car.rs b/src/car.rs index 275e764..9f4bbbc 100644 --- a/src/car.rs +++ b/src/car.rs @@ -21,14 +21,7 @@ pub(crate) struct ParsedCar { #[cfg_attr(not(feature = "indexer"), allow(dead_code))] 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()?; - let root = header - .roots() - .first() - .copied() - .ok_or_else(|| miette::miette!("CAR file has no roots"))?; + let (root, blocks_start) = car_root(&data)?; let mut blocks = CarBlocks::with_capacity_and_hasher( data.len() / EXPECTED_SECTION_BYTES, ahash::RandomState::new(), @@ -40,16 +33,47 @@ pub(crate) fn parse_car(data: Bytes, verify_cids: bool) -> Result { Ok(ParsedCar { root, blocks }) } +/// for commit-sized cars (firehose `#commit` and `#sync`, `getRecord` proofs): the root and an +/// ordered map, the shape the commit pipeline and jacquard's `MemoryBlockStore` take. +pub(crate) fn parse_commit_car( + data: Bytes, + verify_cids: bool, +) -> Result { + let (root, blocks_start) = car_root(&data)?; + let blocks = ordered_blocks(&data, blocks_start, verify_cids)?; + Ok(jacquard_repo::car::reader::ParsedCar { root, blocks }) +} + #[cfg_attr(not(feature = "indexer"), allow(dead_code))] pub(crate) fn parse_car_blocks(data: Bytes, verify_cids: bool) -> Result> { let (_, blocks_start) = car_header_bounds(&data)?; + ordered_blocks(&data, blocks_start, verify_cids) +} + +fn ordered_blocks( + data: &Bytes, + blocks_start: usize, + verify_cids: bool, +) -> Result> { let mut blocks = BTreeMap::new(); - for_each_block(&data, blocks_start, verify_cids, |cid, payload| { + for_each_block(data, blocks_start, verify_cids, |cid, payload| { blocks.insert(cid, payload); })?; Ok(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 + .roots() + .first() + .copied() + .ok_or_else(|| miette::miette!("CAR file has no roots"))?; + Ok((root, blocks_start)) +} + fn car_header_bounds(data: &[u8]) -> Result<(usize, usize)> { let mut offset = 0; let Some(header_len) = read_uvarint(data, &mut offset)? else { @@ -102,15 +126,7 @@ fn for_each_block( Ok(()) } -/// validates that each block's payload hashes to its claimed CID (sha2-256, dag-cbor). -pub(crate) fn validate_block_cids<'a>( - blocks: impl IntoIterator, -) -> Result<()> { - blocks - .into_iter() - .try_for_each(|(claimed_cid, bytes)| validate_block_cid(claimed_cid, bytes)) -} - +/// 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()?; if computed_cid != *claimed_cid { @@ -121,7 +137,6 @@ fn validate_block_cid(claimed_cid: &IpldCid, bytes: &[u8]) -> Result<()> { Ok(()) } -#[cfg_attr(not(feature = "indexer"), allow(dead_code))] fn read_uvarint(data: &[u8], offset: &mut usize) -> Result> { if *offset == data.len() { return Ok(None); @@ -212,6 +227,8 @@ mod tests { 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_commit_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")); } @@ -222,9 +239,12 @@ mod tests { 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(); + let parsed = parse_car(car.clone(), true).unwrap(); assert_eq!(parsed.root, root); assert_eq!(parsed.blocks.get(&leaf).unwrap().as_ref(), b"leaf"); + let commit = parse_commit_car(car, true).unwrap(); + assert_eq!(commit.root, root); + assert_eq!(commit.blocks[&leaf].as_ref(), b"leaf"); } #[tokio::test] @@ -234,32 +254,53 @@ mod tests { 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(); + let err = parse_car(car.clone(), true).unwrap_err(); + assert!(err.to_string().contains("CAR block CID mismatch")); + let err = parse_commit_car(car, true).unwrap_err(); assert!(err.to_string().contains("CAR block CID mismatch")); } #[tokio::test] - async fn repo_and_block_parses_keep_the_last_duplicate_copy() { + async fn every_parse_keeps_the_last_duplicate_copy() { let car = car_with_blocks(&[(cid(1), b"first"), (cid(1), b"second")]).await; let repo = parse_car(car.clone(), false).unwrap(); assert_eq!(repo.blocks[&cid(1)].as_ref(), b"second"); + let commit = parse_commit_car(car.clone(), false).unwrap(); + assert_eq!(commit.blocks[&cid(1)].as_ref(), b"second"); let blocks = parse_car_blocks(car, false).unwrap(); assert_eq!(blocks[&cid(1)].as_ref(), b"second"); } - #[test] - fn validate_block_cids_rejects_mismatched_cid() { - let blocks = BTreeMap::from([(cid(1), Bytes::from_static(b"forged"))]); - let err = validate_block_cids(&blocks).unwrap_err(); - assert!(err.to_string().contains("CAR block CID mismatch")); + // iroh-car's reader, behind jacquard's `parse_car_bytes`, accepts each of the next three + + #[tokio::test] + async fn rejects_a_trailing_partial_varint() { + // iroh-car takes a varint cut off by the end of the car for a clean end + let mut car = car_with_blocks(&[(cid(1), b"block")]).await.to_vec(); + car.push(0x80); + + let err = parse_commit_car(car.into(), false).unwrap_err(); + assert!(err.to_string().contains("truncated uvarint")); + } + + #[tokio::test] + async fn rejects_a_section_longer_than_the_car() { + // iroh-car zero-fills a buffer of the claimed length before reading into it + let mut car = car_with_blocks(&[(cid(1), b"block")]).await.to_vec(); + car.extend_from_slice(&[0xff, 0xff, 0xff, 0x01, 0x00]); + + let err = parse_commit_car(car.into(), false).unwrap_err(); + assert!(err.to_string().contains("truncated CAR block")); } - #[test] - fn validate_block_cids_accepts_matching_cid() { - let bytes = Bytes::from_static(b"trusted"); - let claimed = jacquard_repo::mst::util::compute_cid(bytes.as_ref()).unwrap(); - let blocks = BTreeMap::from([(claimed, bytes)]); - validate_block_cids(&blocks).unwrap(); + #[tokio::test] + async fn rejects_a_block_without_payload() { + // iroh-car accepts it, and the cid of empty bytes even passes verification + let empty = jacquard_repo::mst::util::compute_cid(b"").unwrap(); + let car = car_with_blocks(&[(empty, b"")]).await; + + let err = parse_commit_car(car, true).unwrap_err(); + assert!(err.to_string().contains("has no payload")); } } diff --git a/src/control/stream/jetstream.rs b/src/control/stream/jetstream.rs index f3fecd6..61be018 100644 --- a/src/control/stream/jetstream.rs +++ b/src/control/stream/jetstream.rs @@ -273,11 +273,7 @@ fn stored_event_to_bytes( let cid = if matches!(action, "create" | "update") { let cid = op.cid.as_ref()?; let cid_ipld = cid.to_ipld().ok()?; - let parsed = tokio::runtime::Handle::current() - .block_on(jacquard_repo::car::reader::parse_car_bytes( - commit.blocks.as_ref(), - )) - .ok()?; + let parsed = crate::car::parse_commit_car(commit.blocks.clone(), false).ok()?; let block = parsed.blocks.get(&cid_ipld)?; let val = serde_ipld_dagcbor::from_slice::(block).ok()?; record_raw = serde_json::value::to_raw_value(&val).ok(); diff --git a/src/ingest/indexer/message.rs b/src/ingest/indexer/message.rs index f83c088..de3eadc 100644 --- a/src/ingest/indexer/message.rs +++ b/src/ingest/indexer/message.rs @@ -10,7 +10,7 @@ pub struct IndexerCommitData { /// true if the relay detected a gap (missing seq) before this commit /// and the indexer should trigger a backfill. pub chain_break: bool, - /// result of parse_car_bytes, already done by relay so indexer does not re-parse. + /// the car the relay already parsed and validated, so the indexer does not re-parse it. pub parsed_blocks: jacquard_repo::car::reader::ParsedCar, } diff --git a/src/ingest/indexer/shard.rs b/src/ingest/indexer/shard.rs index b07f0aa..38c2215 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -463,13 +463,10 @@ impl FirehoseWorker { let (key, value) = guard.into_inner().into_diagnostic()?; let commit: Commit = rmp_serde::from_slice(&value).into_diagnostic()?; - let parsed_blocks = TokioHandle::current() - .block_on(jacquard_repo::car::reader::parse_car_bytes( - commit.blocks.as_ref(), - )) + // buffered commits were validated (cids included) and source-checked on arrival, + // so neither the cid check nor the host check is repeated here + let parsed_blocks = crate::car::parse_commit_car(commit.blocks.clone(), false) .map_err(|e| IngestError::Generic(miette::miette!("malformed CAR: {e}")))?; - - // buffered commits have already been source-checked on arrival; skip host check let res = Self::handle_commit(ctx, repo_state, &commit, false, parsed_blocks); let res = match res { Ok(r) => r, diff --git a/src/ingest/validation.rs b/src/ingest/validation.rs index ee47094..94d9090 100644 --- a/src/ingest/validation.rs +++ b/src/ingest/validation.rs @@ -2,7 +2,7 @@ use jacquard_common::IntoStatic; use jacquard_common::types::crypto::PublicKey; use jacquard_repo::MemoryBlockStore; use jacquard_repo::Mst; -use jacquard_repo::car::reader::{ParsedCar, parse_car_bytes}; +use jacquard_repo::car::reader::ParsedCar; use jacquard_repo::commit::Commit as AtpCommit; use jacquard_repo::mst::VerifiedWriteOp; use miette::IntoDiagnostic; @@ -100,7 +100,7 @@ impl ChainBreak { pub struct ValidatedCommit<'c> { #[cfg_attr(not(feature = "indexer"), allow(dead_code))] pub commit: &'c Commit<'c>, - /// result of parse_car_bytes, already done so apply_commit does not re-parse + /// the parsed car, kept so apply_commit does not re-parse it pub parsed_blocks: ParsedCar, pub commit_obj: AtpCommit, pub chain_break: ChainBreak, @@ -130,7 +130,7 @@ impl From<&crate::config::Config> for ValidationOptions { } } -/// all methods panic if called outside a tokio runtime context. +/// `validate_commit` panics if called outside a tokio runtime context. pub struct ValidationContext<'a> { pub opts: &'a ValidationOptions, } @@ -204,15 +204,9 @@ pub fn validate_commit<'c>( } } - // 4. CAR parse - let parsed = handle - .block_on(parse_car_bytes(msg.blocks.as_ref())) - .map_err(|e| CommitValidationError::MalformedCar(miette::miette!("{e}")))?; - - if opts.verify_cids { - crate::car::validate_block_cids(&parsed.blocks) - .map_err(CommitValidationError::MalformedCar)?; - } + // 4. CAR parse, checking each block's cid as it is read + let parsed = crate::car::parse_commit_car(msg.blocks.clone(), opts.verify_cids) + .map_err(CommitValidationError::MalformedCar)?; let root_bytes = parsed.blocks.get(&parsed.root).ok_or_else(|| { CommitValidationError::MalformedCar(miette::miette!("root block missing from CAR")) @@ -289,13 +283,11 @@ pub fn validate_commit<'c>( }) } -/// panics if called outside a tokio runtime context. pub fn validate_sync<'c>( msg: &'c Sync<'c>, signing_key: Option<&PublicKey>, opts: &ValidationOptions, ) -> Result { - let handle = tokio::runtime::Handle::current(); const MAX_BLOCKS_BYTES: usize = 2_097_152; // 1. size limit @@ -303,15 +295,9 @@ pub fn validate_sync<'c>( return Err(SyncValidationError::SizeLimitExceeded); } - // 2. CAR parse - let parsed = handle - .block_on(parse_car_bytes(msg.blocks.as_ref())) - .map_err(|e| SyncValidationError::MalformedCar(miette::miette!("{e}")))?; - - if opts.verify_cids { - crate::car::validate_block_cids(&parsed.blocks) - .map_err(SyncValidationError::MalformedCar)?; - } + // 2. CAR parse, checking each block's cid as it is read + let parsed = crate::car::parse_commit_car(msg.blocks.clone(), opts.verify_cids) + .map_err(SyncValidationError::MalformedCar)?; let root_bytes = parsed.blocks.get(&parsed.root).ok_or_else(|| { SyncValidationError::MalformedCar(miette::miette!("root block missing from CAR")) @@ -450,6 +436,7 @@ fn verify_mst( mod tests { use super::*; use crate::ingest::stream::types::Datetime; + use jacquard_common::types::cid::CidLink; use jacquard_common::types::string::Did; #[test] @@ -459,24 +446,50 @@ mod tests { let rt = tokio::runtime::Runtime::new().unwrap(); let _guard = rt.enter(); let claimed = jacquard_repo::mst::util::compute_cid(b"trusted").unwrap(); - let blocks = rt.block_on(async { - let mut car = Vec::new(); - let header = iroh_car::CarHeader::new_v1(vec![claimed]); - let mut writer = iroh_car::CarWriter::new(header, &mut car); - writer.write(claimed, b"forged".to_vec()).await.unwrap(); - writer.finish().await.unwrap(); - car - }); - let msg = Sync { - blocks: blocks.into(), - did: Did::new_static("did:plc:aaaaaaaaaaaaaaaaaaaaaaaa").unwrap(), - rev: "3l2abcdefgh22".into(), + let blocks: bytes::Bytes = rt + .block_on(async { + let mut car = Vec::new(); + let header = iroh_car::CarHeader::new_v1(vec![claimed]); + let mut writer = iroh_car::CarWriter::new(header, &mut car); + writer.write(claimed, b"forged".to_vec()).await.unwrap(); + writer.finish().await.unwrap(); + car + }) + .into(); + let did = Did::new_static("did:plc:aaaaaaaaaaaaaaaaaaaaaaaa").unwrap(); + let rev = "3l2abcdefgh22"; + let time = Datetime(chrono::Utc::now().fixed_offset()); + + let sync = Sync { + blocks: blocks.clone(), + did: did.clone(), + rev: rev.into(), seq: 1, - time: Datetime(chrono::Utc::now().fixed_offset()), + time: time.clone(), + }; + let Err(SyncValidationError::MalformedCar(err)) = validate_sync(&sync, None, &opts) else { + panic!("forged #sync block was accepted"); }; + assert!(err.to_string().contains("CAR block CID mismatch"), "{err}"); - let Err(SyncValidationError::MalformedCar(err)) = validate_sync(&msg, None, &opts) else { - panic!("forged block was accepted"); + let commit = Commit { + blobs: Vec::new(), + blocks, + commit: CidLink::from(claimed), + ops: Vec::new(), + prev_data: None, + rebase: false, + repo: did, + rev: rev.parse().unwrap(), + seq: 2, + since: None, + time, + too_big: false, + }; + let Err(CommitValidationError::MalformedCar(err)) = + validate_commit(&commit, &RepoState::synced(), None, &opts) + else { + panic!("forged #commit block was accepted"); }; assert!(err.to_string().contains("CAR block CID mismatch"), "{err}"); }