From 3a511b3f69445e3b467bb748e66a9a0a392604c2 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Wed, 30 Sep 2026 07:52:11 +0300 Subject: [PATCH] [car] keep full repo car blocks in a seeded hash map --- src/backfill/worker/process.rs | 9 +- src/car.rs | 161 ++++++++++++++++++++++++++++++--- 2 files changed, 153 insertions(+), 17 deletions(-) diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index 4c0dfa8..abd331a 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -13,13 +13,14 @@ 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::did::Did; +use jacquard_repo::BlockStore; use jacquard_repo::mst::Mst; -use jacquard_repo::{BlockStore, MemoryBlockStore}; use crate::backfill::admission::BackfillAdmission; use crate::backfill::client::{ThrottledHttpClient, collect_body_bounded}; use crate::backfill::error::BackfillError; use crate::backfill::sparse::{SparseBackfillResult, process_did_sparse}; +use crate::car::CarBlockStore; use crate::config::{BackfillStrategy, RateTier}; use crate::db::types::{DbAction, DbRkey}; use crate::db::{self, Txn as DbTxn, keys}; @@ -281,7 +282,7 @@ pub(super) async fn process_did( let start = Instant::now(); let root_cid = parsed.root; - let store = Arc::new(MemoryBlockStore::new_from_blocks(parsed.blocks)); + let store = Arc::new(CarBlockStore::new(parsed.blocks)); trace!( blocks = store.len(), elapsed = ?start.elapsed(), @@ -315,9 +316,9 @@ pub(super) async fn process_did( let root_commit = Commit::from(root_commit); - // 5. walk mst and fetch every record block under one store lock + // 5. walk mst and fetch every record block in one batch let start = Instant::now(); - let mst: Mst = Mst::load(store, root_commit.data, None); + let mst: Mst = Mst::load(store, root_commit.data, None); let leaves = mst.leaves().await.into_diagnostic()?; let leaf_cids = leaves.iter().map(|(_, cid)| *cid).collect::>(); let leaf_blocks = mst.storage().get_many(&leaf_cids).await.into_diagnostic()?; diff --git a/src/car.rs b/src/car.rs index 0c4065e..4f8ec61 100644 --- a/src/car.rs +++ b/src/car.rs @@ -1,15 +1,24 @@ -use std::collections::BTreeMap; +use std::collections::{BTreeMap, HashMap}; use std::io::Cursor; +use std::sync::Arc; use bytes::Bytes; use cid::Cid as IpldCid; +use jacquard_repo::{BlockStore, CommitData, RepoError, RepoErrorKind}; use miette::{IntoDiagnostic, Result}; +/// blocks of a whole-repo CAR. hashed rather than ordered: the MST walk looks up every +/// block, and cids are chosen by the PDS, so the hasher is randomly seeded. +pub(crate) type CarBlocks = HashMap; + +/// repo blocks average a few hundred bytes; sizing for that avoids most rehashing. +const EXPECTED_SECTION_BYTES: usize = 256; + #[cfg_attr(not(feature = "indexer"), allow(dead_code))] #[derive(Debug)] pub(crate) struct ParsedCar { pub root: IpldCid, - pub blocks: BTreeMap, + pub blocks: CarBlocks, } #[cfg_attr(not(feature = "indexer"), allow(dead_code))] @@ -22,7 +31,13 @@ pub(crate) fn parse_car(data: Bytes, verify_cids: bool) -> Result { .first() .copied() .ok_or_else(|| miette::miette!("CAR file has no roots"))?; - let blocks = parse_car_blocks_from(data, blocks_start, verify_cids)?; + let mut blocks = CarBlocks::with_capacity_and_hasher( + data.len() / EXPECTED_SECTION_BYTES, + ahash::RandomState::new(), + ); + for_each_block(&data, blocks_start, verify_cids, |cid, payload| { + blocks.insert(cid, payload); + })?; Ok(ParsedCar { root, blocks }) } @@ -30,7 +45,11 @@ pub(crate) fn parse_car(data: Bytes, verify_cids: bool) -> Result { #[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)?; - parse_car_blocks_from(data, blocks_start, verify_cids) + let mut blocks = BTreeMap::new(); + for_each_block(&data, blocks_start, verify_cids, |cid, payload| { + blocks.insert(cid, payload); + })?; + Ok(blocks) } fn car_header_bounds(data: &[u8]) -> Result<(usize, usize)> { @@ -49,15 +68,16 @@ fn car_header_bounds(data: &[u8]) -> Result<(usize, usize)> { Ok((header_start, header_end)) } -/// 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, +/// yields each section's cid and zero-copy payload in file order. with `verify_cids`, every +/// section is hashed as it is read, while its bytes are still in cache; checking a finished +/// map afterwards would walk memory in cid order. +fn for_each_block( + data: &Bytes, mut offset: usize, verify_cids: bool, -) -> Result> { - let mut blocks = BTreeMap::new(); - while let Some(section_len) = read_uvarint(&data, &mut offset)? { + mut on_block: impl FnMut(IpldCid, Bytes), +) -> Result<()> { + while let Some(section_len) = read_uvarint(data, &mut offset)? { let section_end = offset .checked_add(section_len) .ok_or_else(|| miette::miette!("CAR block length overflow"))?; @@ -78,10 +98,10 @@ fn parse_car_blocks_from( if verify_cids { validate_block_cid(&cid, &payload)?; } - blocks.insert(cid, payload); + on_block(cid, payload); } - Ok(blocks) + Ok(()) } /// validates that each block's payload hashes to its claimed CID (sha2-256, dag-cbor). @@ -103,6 +123,58 @@ fn validate_block_cid(claimed_cid: &IpldCid, bytes: &[u8]) -> Result<()> { Ok(()) } +/// read-only [`BlockStore`] over the blocks of a parsed CAR, for walking an imported repo. +#[cfg_attr(not(feature = "indexer"), allow(dead_code))] +#[derive(Clone)] +pub(crate) struct CarBlockStore(Arc); + +#[cfg_attr(not(feature = "indexer"), allow(dead_code))] +impl CarBlockStore { + pub(crate) fn new(blocks: CarBlocks) -> Self { + Self(Arc::new(blocks)) + } + + pub(crate) fn len(&self) -> usize { + self.0.len() + } +} + +fn read_only_store() -> RepoError { + RepoError::new( + RepoErrorKind::Storage, + Some("CarBlockStore is read-only".into()), + ) +} + +impl BlockStore for CarBlockStore { + async fn get(&self, cid: &IpldCid) -> jacquard_repo::Result> { + Ok(self.0.get(cid).cloned()) + } + + async fn put(&self, _data: &[u8]) -> jacquard_repo::Result { + Err(read_only_store()) + } + + async fn has(&self, cid: &IpldCid) -> jacquard_repo::Result { + Ok(self.0.contains_key(cid)) + } + + async fn put_many( + &self, + _blocks: impl IntoIterator + Send, + ) -> jacquard_repo::Result<()> { + Err(read_only_store()) + } + + async fn get_many(&self, cids: &[IpldCid]) -> jacquard_repo::Result>> { + Ok(cids.iter().map(|cid| self.0.get(cid).cloned()).collect()) + } + + async fn apply_commit(&self, _commit: CommitData) -> jacquard_repo::Result<()> { + Err(read_only_store()) + } +} + #[cfg_attr(not(feature = "indexer"), allow(dead_code))] fn read_uvarint(data: &[u8], offset: &mut usize) -> Result> { if *offset == data.len() { @@ -220,6 +292,69 @@ mod tests { assert!(err.to_string().contains("CAR block CID mismatch")); } + #[tokio::test] + async fn repo_and_block_parses_keep_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 blocks = parse_car_blocks(car, false).unwrap(); + assert_eq!(blocks[&cid(1)].as_ref(), b"second"); + } + + #[tokio::test] + async fn car_block_store_serves_reads_and_rejects_writes() { + let present = cid(1); + let store = CarBlockStore::new(CarBlocks::from_iter([( + present, + Bytes::from_static(b"block"), + )])); + + assert_eq!( + store.get(&present).await.unwrap().as_deref(), + Some(&b"block"[..]) + ); + assert_eq!(store.get(&cid(2)).await.unwrap(), None); + assert!(store.has(&present).await.unwrap()); + assert!(!store.has(&cid(2)).await.unwrap()); + assert_eq!( + store.get_many(&[cid(2), present]).await.unwrap(), + vec![None, Some(Bytes::from_static(b"block"))] + ); + assert!(store.put(b"new").await.is_err()); + assert!(store.put_many([(cid(3), Bytes::new())]).await.is_err()); + assert_eq!(store.len(), 1); + } + + #[tokio::test] + async fn car_block_store_walks_a_parsed_repo_like_memory_store() { + use jacquard_repo::{MemoryBlockStore, Mst}; + + let mut mst = Mst::new(Arc::new(MemoryBlockStore::new())); + let mut records = Vec::new(); + for i in 0..500 { + let body = format!("record {i}").into_bytes(); + let record_cid = jacquard_repo::mst::util::compute_cid(&body).unwrap(); + let key = format!("app.bsky.feed.post/{i:013}"); + mst = mst.add(&key, record_cid).await.unwrap(); + records.push((record_cid, body)); + } + let (root, nodes) = mst.collect_blocks().await.unwrap(); + let blocks: Vec<(IpldCid, &[u8])> = nodes + .iter() + .map(|(cid, bytes)| (*cid, bytes.as_ref())) + .chain(records.iter().map(|(cid, body)| (*cid, body.as_slice()))) + .collect(); + let parsed = parse_car(car_with_blocks(&blocks).await, true).unwrap(); + + let loaded = Mst::load(Arc::new(CarBlockStore::new(parsed.blocks)), root, None); + let expected = mst.leaves().await.unwrap(); + assert_eq!(loaded.leaves().await.unwrap(), expected); + let leaf_cids: Vec<_> = expected.iter().map(|(_, cid)| *cid).collect(); + let bodies = loaded.storage().get_many(&leaf_cids).await.unwrap(); + assert!(bodies.iter().all(Option::is_some)); + } + #[test] fn validate_block_cids_rejects_mismatched_cid() { let blocks = BTreeMap::from([(cid(1), Bytes::from_static(b"forged"))]); -- 2.51.2