diff --git a/Cargo.lock b/Cargo.lock index 7db6fe16..d529cbb8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2628,6 +2628,7 @@ dependencies = [ "ed25519-dalek", "hex", "iroh-car", + "jacquard-api", "jacquard-common", "jacquard-derive", "k256", diff --git a/crates/jacquard-repo/Cargo.toml b/crates/jacquard-repo/Cargo.toml index 157cb313..e7b47df4 100644 --- a/crates/jacquard-repo/Cargo.toml +++ b/crates/jacquard-repo/Cargo.toml @@ -18,6 +18,7 @@ default = [] # Internal jacquard-common = { path = "../jacquard-common", version = "0.9", features = ["crypto-ed25519", "crypto-k256", "crypto-p256"] } jacquard-derive = { path = "../jacquard-derive", version = "0.9" } +jacquard-api = { path = "../jacquard-api", version = "0.9", features = ["streaming"] } # Serialization serde.workspace = true diff --git a/crates/jacquard-repo/src/commit/firehose.rs b/crates/jacquard-repo/src/commit/firehose.rs index b65dd0e7..a4848469 100644 --- a/crates/jacquard-repo/src/commit/firehose.rs +++ b/crates/jacquard-repo/src/commit/firehose.rs @@ -4,209 +4,67 @@ //! to avoid a dependency on the full API crate. They represent firehose protocol messages, //! which are DISTINCT from repository commit objects. -use bytes::Bytes; -use jacquard_common::types::cid::CidLink; +pub use jacquard_api::com_atproto::sync::subscribe_repos::Commit as FirehoseCommit; +pub use jacquard_api::com_atproto::sync::subscribe_repos::RepoOp; +use jacquard_api::com_atproto::sync::subscribe_repos::{Commit, RepoOpAction}; use jacquard_common::types::crypto::PublicKey; -use jacquard_common::types::string::{Datetime, Did, Tid}; -use jacquard_common::{CowStr, IntoStatic}; use smol_str::ToSmolStr; -/// Firehose commit message (sync v1.0 and v1.1) +/// Convert to VerifiedWriteOp for v1.1 validation /// -/// Represents an update of repository state in the firehose stream. -/// This is the message format sent over `com.atproto.sync.subscribeRepos`. -/// -/// **Sync v1.0 vs v1.1:** -/// - v1.0: `prev_data` is None/skipped, consumers must have sufficient previous repository state to validate -/// - v1.1: `prev_data` includes previous MST root for inductive validation -#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] -#[serde(rename_all = "camelCase")] -pub struct FirehoseCommit<'a> { - /// The repo this event comes from - #[serde(borrow)] - pub repo: Did<'a>, - - /// The rev of the emitted commit - pub rev: Tid, - - /// The stream sequence number of this message - pub seq: i64, - - /// The rev of the last emitted commit from this repo (if any) - pub since: Tid, - - /// Timestamp of when this message was originally broadcast - pub time: Datetime, - - /// Repo commit object CID - /// - /// This CID points to the repository commit block (with did, version, data, rev, prev, sig). - /// It must be the first entry in the CAR header 'roots' list. - #[serde(borrow)] - pub commit: CidLink<'a>, - - /// CAR file containing relevant blocks - /// - /// Contains blocks as a diff since the previous repo state. The commit block - /// must be included, and its CID must be the first root in the CAR header. - /// - /// For sync v1.1, may include additional MST node blocks needed for operation inversion. - #[serde(with = "super::serde_bytes_helper")] - pub blocks: Bytes, - - /// Operations in this commit - #[serde(borrow)] - pub ops: Vec>, - - /// Previous MST root CID (sync v1.1 only) - /// - /// The root CID of the MST tree for the previous commit (indicated by the 'since' field). - /// Corresponds to the 'data' field in the previous repo commit object. - /// - /// **Sync v1.1 inductive validation:** - /// - Enables validation without local MST state - /// - Operations can be inverted (creates→deletes, deletes→creates with prev values) - /// - Required for "inductive firehose" consumption - /// - /// **Sync v1.0:** - /// - This field is None - /// - Consumers must have previous repository state - #[serde(skip_serializing_if = "Option::is_none")] - #[serde(borrow)] - pub prev_data: Option>, - - /// Blob CIDs referenced in this commit - #[serde(borrow)] - pub blobs: Vec>, - - /// DEPRECATED: Replaced by #sync event and data limits - /// - /// Indicates that this commit contained too many ops, or data size was too large. - /// Consumers will need to make a separate request to get missing data. - pub too_big: bool, - - /// DEPRECATED: Unused - pub rebase: bool, -} - -/// A repository operation (mutation of a single record) -#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] -#[serde(rename_all = "camelCase")] -pub struct RepoOp<'a> { - /// Operation type: "create", "update", or "delete" - #[serde(borrow)] - pub action: CowStr<'a>, - - /// Collection/rkey path (e.g., "app.bsky.feed.post/abc123") - #[serde(borrow)] - pub path: CowStr<'a>, - - /// For creates and updates, the new record CID. For deletions, None (null). - #[serde(skip_serializing_if = "Option::is_none")] - #[serde(borrow)] - pub cid: Option>, - - /// For updates and deletes, the previous record CID - /// - /// Required for sync v1.1 inductive firehose validation. - /// For creates, this field should not be defined. - #[serde(skip_serializing_if = "Option::is_none")] - #[serde(borrow)] - pub prev: Option>, -} - -impl<'a> RepoOp<'a> { - /// Convert to VerifiedWriteOp for v1.1 validation - /// - /// Validates that all required fields are present for inversion. - pub fn to_invertible_op(&self) -> Result { - let key = self.path.to_smolstr(); - - match self.action.as_ref() { - "create" => { - let cid = self - .cid - .as_ref() - .ok_or_else(|| RepoError::invalid_commit("create operation missing cid field"))? - .to_ipld() - .map_err(|e| RepoError::invalid_cid_conversion(e, "create cid"))?; - - Ok(VerifiedWriteOp::Create { key, cid }) - } - "update" => { - let cid = self - .cid - .as_ref() - .ok_or_else(|| RepoError::invalid_commit("update operation missing cid field"))? - .to_ipld() - .map_err(|e| RepoError::invalid_cid_conversion(e, "update cid"))?; - - let prev = self - .prev - .as_ref() - .ok_or_else(|| { - RepoError::invalid_commit( - "update operation missing prev field for v1.1 validation", - ) - })? - .to_ipld() - .map_err(|e| RepoError::invalid_cid_conversion(e, "update prev"))?; - - Ok(VerifiedWriteOp::Update { key, cid, prev }) - } - "delete" => { - let prev = self - .prev - .as_ref() - .ok_or_else(|| { - RepoError::invalid_commit( - "delete operation missing prev field for v1.1 validation", - ) - })? - .to_ipld() - .map_err(|e| RepoError::invalid_cid_conversion(e, "delete prev"))?; - - Ok(VerifiedWriteOp::Delete { key, prev }) - } - action => Err(RepoError::invalid_commit(format!( - "unknown action type: {}", - action - ))), +/// Validates that all required fields are present for inversion. +pub fn to_invertible_op(op: &RepoOp<'_>) -> Result { + let key = op.path.to_smolstr(); + match op.action { + RepoOpAction::Create => { + let cid = op + .cid + .as_ref() + .ok_or_else(|| RepoError::invalid_commit("create operation missing cid field"))? + .to_ipld() + .map_err(|e| RepoError::invalid_cid_conversion(e, "create cid"))?; + + Ok(VerifiedWriteOp::Create { key, cid }) } - } -} - -impl IntoStatic for FirehoseCommit<'_> { - type Output = FirehoseCommit<'static>; - - fn into_static(self) -> Self::Output { - FirehoseCommit { - repo: self.repo.into_static(), - rev: self.rev, - seq: self.seq, - since: self.since, - time: self.time, - commit: self.commit.into_static(), - blocks: self.blocks, - ops: self.ops.into_iter().map(|op| op.into_static()).collect(), - prev_data: self.prev_data.map(|pd| pd.into_static()), - blobs: self.blobs.into_iter().map(|b| b.into_static()).collect(), - too_big: self.too_big, - rebase: self.rebase, + RepoOpAction::Update => { + let cid = op + .cid + .as_ref() + .ok_or_else(|| RepoError::invalid_commit("update operation missing cid field"))? + .to_ipld() + .map_err(|e| RepoError::invalid_cid_conversion(e, "update cid"))?; + + let prev = op + .prev + .as_ref() + .ok_or_else(|| { + RepoError::invalid_commit( + "update operation missing prev field for v1.1 validation", + ) + })? + .to_ipld() + .map_err(|e| RepoError::invalid_cid_conversion(e, "update prev"))?; + + Ok(VerifiedWriteOp::Update { key, cid, prev }) } - } -} - -impl IntoStatic for RepoOp<'_> { - type Output = RepoOp<'static>; - - fn into_static(self) -> Self::Output { - RepoOp { - action: self.action.into_static(), - path: self.path.into_static(), - cid: self.cid.into_static(), - prev: self.prev.map(|p| p.into_static()), + RepoOpAction::Delete => { + let prev = op + .prev + .as_ref() + .ok_or_else(|| { + RepoError::invalid_commit( + "delete operation missing prev field for v1.1 validation", + ) + })? + .to_ipld() + .map_err(|e| RepoError::invalid_cid_conversion(e, "delete prev"))?; + + Ok(VerifiedWriteOp::Delete { key, prev }) } + RepoOpAction::Other(ref action) => Err(RepoError::invalid_commit(format!( + "unknown action type: {}", + action + ))), } } @@ -220,196 +78,192 @@ use crate::storage::{BlockStore, LayeredBlockStore, MemoryBlockStore}; use cid::Cid as IpldCid; use std::sync::Arc; -impl<'a> FirehoseCommit<'a> { - /// Validate a sync v1.0 commit - /// - /// **Requirements:** - /// - Must have previous MST state (potentially full repository) - /// - All blocks needed for validation must be in `self.blocks` - /// - /// **Validation steps:** - /// 1. Parse CAR blocks from `self.blocks` into temporary storage - /// 2. Load commit object and verify signature - /// 3. Apply operations to previous MST (using temporary storage for new blocks) - /// 4. Verify result matches commit.data (new MST root) - /// - /// Returns the new MST root CID on success. - pub async fn validate_v1_0( - &self, - prev_mst_root: Option, - prev_storage: Arc, - pubkey: &PublicKey<'_>, - ) -> Result { - // 1. Parse CAR blocks from the firehose message into temporary storage - let parsed = parse_car_bytes(&self.blocks).await?; - let temp_storage = MemoryBlockStore::new_from_blocks(parsed.blocks); - - // 2. Create layered storage: reads from temp first, then prev; writes to temp only - // This avoids copying all previous MST blocks - let layered_storage = LayeredBlockStore::new(temp_storage.clone(), prev_storage); - - // 3. Extract and verify commit object from temporary storage - let commit_cid: IpldCid = self - .commit - .to_ipld() - .map_err(|e| RepoError::invalid_cid_conversion(e, "commit CID"))?; - let commit_bytes = temp_storage - .get(&commit_cid) - .await? - .ok_or_else(|| RepoError::not_found("commit block", &commit_cid))?; - - let commit = super::Commit::from_cbor(&commit_bytes)?; - - // Verify DID matches - if commit.did().as_ref() != self.repo.as_ref() { - return Err(RepoError::invalid_commit(format!( +/// Validate a sync v1.0 commit +/// +/// **Requirements:** +/// - Must have previous MST state (potentially full repository) +/// - All blocks needed for validation must be in `self.blocks` +/// +/// **Validation steps:** +/// 1. Parse CAR blocks from `self.blocks` into temporary storage +/// 2. Load commit object and verify signature +/// 3. Apply operations to previous MST (using temporary storage for new blocks) +/// 4. Verify result matches commit.data (new MST root) +/// +/// Returns the new MST root CID on success. +pub async fn validate_v1_0( + fh_commit: &Commit<'_>, + prev_mst_root: Option, + prev_storage: Arc, + pubkey: &PublicKey<'_>, +) -> Result { + // 1. Parse CAR blocks from the firehose message into temporary storage + let parsed = parse_car_bytes(&fh_commit.blocks).await?; + let temp_storage = MemoryBlockStore::new_from_blocks(parsed.blocks); + + // 2. Create layered storage: reads from temp first, then prev; writes to temp only + // This avoids copying all previous MST blocks + let layered_storage = LayeredBlockStore::new(temp_storage.clone(), prev_storage); + + // 3. Extract and verify commit object from temporary storage + let commit_cid: IpldCid = fh_commit + .commit + .to_ipld() + .map_err(|e| RepoError::invalid_cid_conversion(e, "commit CID"))?; + let commit_bytes = temp_storage + .get(&commit_cid) + .await? + .ok_or_else(|| RepoError::not_found("commit block", &commit_cid))?; + + let commit = super::Commit::from_cbor(&commit_bytes)?; + + // Verify DID matches + if commit.did().as_ref() != fh_commit.repo.as_ref() { + return Err(RepoError::invalid_commit(format!( "DID mismatch: commit has {}, message has {}", commit.did(), - self.repo + fh_commit.repo )) .with_help("DID mismatch indicates the commit was signed by a different identity - verify the commit is from the expected repository")); - } - - // Verify signature - commit.verify(pubkey)?; + } - let layered_arc = Arc::new(layered_storage); + // Verify signature + commit.verify(pubkey)?; - // 4. Load previous MST state from layered storage (or start empty) - let prev_mst = if let Some(prev_root) = prev_mst_root { - Mst::load(layered_arc.clone(), prev_root, None) - } else { - Mst::new(layered_arc.clone()) - }; + let layered_arc = Arc::new(layered_storage); - // 5. Load new MST from commit.data (claimed result) - let expected_root = *commit.data(); - let new_mst = Mst::load(layered_arc, expected_root, None); + // 4. Load previous MST state from layered storage (or start empty) + let prev_mst = if let Some(prev_root) = prev_mst_root { + Mst::load(layered_arc.clone(), prev_root, None) + } else { + Mst::new(layered_arc.clone()) + }; - // 6. Compute diff to get verified write ops (with actual prev values from tree state) - let diff = prev_mst.diff(&new_mst).await?; - let verified_ops = diff.to_verified_ops(); + // 5. Load new MST from commit.data (claimed result) + let expected_root = *commit.data(); + let new_mst = Mst::load(layered_arc, expected_root, None); - // 7. Apply verified ops to prev MST - let computed_mst = prev_mst.batch(&verified_ops).await?; + // 6. Compute diff to get verified write ops (with actual prev values from tree state) + let diff = prev_mst.diff(&new_mst).await?; + let verified_ops = diff.to_verified_ops(); - // 8. Verify computed result matches claimed result - let computed_root = computed_mst.get_pointer().await?; + // 7. Apply verified ops to prev MST + let computed_mst = prev_mst.batch(&verified_ops).await?; - if computed_root != expected_root { - return Err(RepoError::cid_mismatch(format!( - "MST root mismatch: expected {}, got {}", - expected_root, computed_root - ))); - } + // 8. Verify computed result matches claimed result + let computed_root = computed_mst.get_pointer().await?; - Ok(expected_root) + if computed_root != expected_root { + return Err(RepoError::cid_mismatch(format!( + "MST root mismatch: expected {}, got {}", + expected_root, computed_root + ))); } - /// Validate a sync v1.1 commit (inductive validation) - /// - /// **Requirements:** - /// - `self.prev_data` must be Some (contains previous MST root) - /// - All blocks needed for validation must be in `self.blocks` - /// - /// **Validation steps:** - /// 1. Parse CAR blocks from `self.blocks` into temporary storage - /// 2. Load commit object and verify signature - /// 3. Start from `prev_data` MST root (loaded from temp storage) - /// 4. Apply operations (with prev CID validation for updates/deletes) - /// 5. Verify result matches commit.data (new MST root) - /// - /// Returns the new MST root CID on success. - /// - /// **Inductive property:** Can validate without any external state besides the blocks - /// in this message. The `prev_data` field provides the starting MST root, and operations - /// include `prev` CIDs for validation. All necessary blocks must be in the CAR bytes. - /// - /// Note: Because this uses the same merkle search tree struct as the repository itself, - /// this is far from the most efficient possible validation function possible. The repo - /// tree struct carries extra information. However, - /// it has the virtue of making everything self-validating. - pub async fn validate_v1_1(&self, pubkey: &PublicKey<'_>) -> Result { - // 1. Require prev_data for v1.1 - let prev_data_cid: IpldCid = self - .prev_data - .as_ref() - .ok_or_else(|| { - RepoError::invalid_commit("Sync v1.1 validation requires prev_data field") - })? - .to_ipld() - .map_err(|e| RepoError::invalid_cid_conversion(e, "prev_data CID"))?; - - // 2. Parse CAR blocks from the firehose message into temporary storage - let parsed = parse_car_bytes(&self.blocks).await?; - - let temp_storage = Arc::new(MemoryBlockStore::new_from_blocks(parsed.blocks)); - - // 3. Extract and verify commit object from temporary storage - let commit_cid: IpldCid = self - .commit - .to_ipld() - .map_err(|e| RepoError::invalid_cid_conversion(e, "commit CID"))?; - let commit_bytes = temp_storage - .get(&commit_cid) - .await? - .ok_or_else(|| RepoError::not_found("commit block", &commit_cid))?; - - let commit = super::Commit::from_cbor(&commit_bytes)?; - - // Verify DID matches - if commit.did().as_ref() != self.repo.as_ref() { - return Err(RepoError::invalid_commit(format!( + Ok(expected_root) +} + +/// Validate a sync v1.1 commit (inductive validation) +/// +/// **Requirements:** +/// - `self.prev_data` must be Some (contains previous MST root) +/// - All blocks needed for validation must be in `self.blocks` +/// +/// **Validation steps:** +/// 1. Parse CAR blocks from `self.blocks` into temporary storage +/// 2. Load commit object and verify signature +/// 3. Start from `prev_data` MST root (loaded from temp storage) +/// 4. Apply operations (with prev CID validation for updates/deletes) +/// 5. Verify result matches commit.data (new MST root) +/// +/// Returns the new MST root CID on success. +/// +/// **Inductive property:** Can validate without any external state besides the blocks +/// in this message. The `prev_data` field provides the starting MST root, and operations +/// include `prev` CIDs for validation. All necessary blocks must be in the CAR bytes. +/// +/// Note: Because this uses the same merkle search tree struct as the repository itself, +/// this is far from the most efficient possible validation function possible. The repo +/// tree struct carries extra information. However, +/// it has the virtue of making everything self-validating. +pub async fn validate_v1_1(fh_commit: &Commit<'_>, pubkey: &PublicKey<'_>) -> Result { + // 1. Require prev_data for v1.1 + let prev_data_cid: IpldCid = fh_commit + .prev_data + .as_ref() + .ok_or_else(|| RepoError::invalid_commit("Sync v1.1 validation requires prev_data field"))? + .to_ipld() + .map_err(|e| RepoError::invalid_cid_conversion(e, "prev_data CID"))?; + + // 2. Parse CAR blocks from the firehose message into temporary storage + let parsed = parse_car_bytes(&fh_commit.blocks).await?; + + let temp_storage = Arc::new(MemoryBlockStore::new_from_blocks(parsed.blocks)); + + // 3. Extract and verify commit object from temporary storage + let commit_cid: IpldCid = fh_commit + .commit + .to_ipld() + .map_err(|e| RepoError::invalid_cid_conversion(e, "commit CID"))?; + let commit_bytes = temp_storage + .get(&commit_cid) + .await? + .ok_or_else(|| RepoError::not_found("commit block", &commit_cid))?; + + let commit = super::Commit::from_cbor(&commit_bytes)?; + + // Verify DID matches + if commit.did().as_ref() != fh_commit.repo.as_ref() { + return Err(RepoError::invalid_commit(format!( "DID mismatch: commit has {}, message has {}", commit.did(), - self.repo + fh_commit.repo )) .with_help("DID mismatch indicates the commit was signed by a different identity - verify the commit is from the expected repository")); - } + } - // Verify signature - commit.verify(pubkey)?; - - // 5. Load new MST from commit.data (claimed result) - let expected_root = *commit.data(); - - let mut new_mst = Mst::load(temp_storage, expected_root, None); - - let verified_ops = self - .ops - .iter() - .filter_map(|op| op.to_invertible_op().ok()) - .collect::>(); - if verified_ops.len() != self.ops.len() { - return Err(RepoError::invalid_commit(format!( - "Invalid commit: expected {} ops, got {}", - self.ops.len(), - verified_ops.len() - ))); - } + // Verify signature + commit.verify(pubkey)?; + + // 5. Load new MST from commit.data (claimed result) + let expected_root = *commit.data(); + + let mut new_mst = Mst::load(temp_storage, expected_root, None); + + let verified_ops = fh_commit + .ops + .iter() + .filter_map(|op| to_invertible_op(op).ok()) + .collect::>(); + if verified_ops.len() != fh_commit.ops.len() { + return Err(RepoError::invalid_commit(format!( + "Invalid commit: expected {} ops, got {}", + fh_commit.ops.len(), + verified_ops.len() + ))); + } - for op in verified_ops { - if let Ok(inverted) = new_mst.invert_op(op.clone()).await { - if !inverted { - return Err(RepoError::invalid_commit(format!( - "Invalid commit: op {:?} is not invertible", - op - ))); - } + for op in verified_ops { + if let Ok(inverted) = new_mst.invert_op(op.clone()).await { + if !inverted { + return Err(RepoError::invalid_commit(format!( + "Invalid commit: op {:?} is not invertible", + op + ))); } } - // 8. Verify computed previous state matches claimed previous state - let computed_root = new_mst.get_pointer().await?; - - if computed_root != prev_data_cid { - return Err(RepoError::cid_mismatch(format!( - "MST root mismatch: expected {}, got {}", - prev_data_cid, computed_root - ))); - } - - Ok(expected_root) } + // 8. Verify computed previous state matches claimed previous state + let computed_root = new_mst.get_pointer().await?; + + if computed_root != prev_data_cid { + return Err(RepoError::cid_mismatch(format!( + "MST root mismatch: expected {}, got {}", + prev_data_cid, computed_root + ))); + } + + Ok(expected_root) } #[cfg(test)] @@ -419,9 +273,11 @@ mod tests { use crate::commit::Commit; use crate::mst::{Mst, RecordWriteOp}; use crate::storage::MemoryBlockStore; + use jacquard_common::IntoStatic; use jacquard_common::types::crypto::{KeyCodec, PublicKey}; + use jacquard_common::types::did::Did; use jacquard_common::types::recordkey::Rkey; - use jacquard_common::types::string::{Nsid, RecordKey}; + use jacquard_common::types::string::{Datetime, Nsid, RecordKey}; use jacquard_common::types::tid::Ticker; use jacquard_common::types::value::RawData; use smol_str::SmolStr; @@ -507,7 +363,7 @@ mod tests { .unwrap(); // Validate using v1.1 validation - let result = firehose_commit.validate_v1_1(&pubkey).await; + let result = validate_v1_1(&firehose_commit, &pubkey).await; if let Err(ref e) = result { eprintln!("Validation error: {}", e); } @@ -560,9 +416,8 @@ mod tests { firehose_commit.prev_data = None; // Validate using v1.0 validation with previous storage - let result = firehose_commit - .validate_v1_0(Some(prev_root), storage.clone(), &pubkey) - .await; + let result = + validate_v1_0(&firehose_commit, Some(prev_root), storage.clone(), &pubkey).await; assert!(result.is_ok(), "Valid v1.0 commit should pass validation"); @@ -612,7 +467,7 @@ mod tests { .await .unwrap(); - let result = firehose_commit.validate_v1_1(&pubkey).await; + let result = validate_v1_1(&firehose_commit, &pubkey).await; assert!(result.is_ok(), "Multiple creates should validate"); } @@ -685,7 +540,7 @@ mod tests { .await .unwrap(); - let result = firehose_commit.validate_v1_1(&pubkey).await; + let result = validate_v1_1(&firehose_commit, &pubkey).await; assert!( result.is_ok(), "Update and delete operations should validate" @@ -740,7 +595,7 @@ mod tests { firehose_commit.blocks = bad_car.into(); - let result = firehose_commit.validate_v1_1(&pubkey).await; + let result = validate_v1_1(&firehose_commit, &pubkey).await; assert!( result.is_err(), "Validation should fail when commit block is missing" @@ -802,7 +657,7 @@ mod tests { firehose_commit.blocks = bad_car.into(); - let result = firehose_commit.validate_v1_1(&pubkey).await; + let result = validate_v1_1(&firehose_commit, &pubkey).await; assert!( result.is_err(), "Validation should fail when MST blocks are missing" @@ -863,7 +718,7 @@ mod tests { .await .unwrap(); - let result = firehose_commit.validate_v1_1(&pubkey).await; + let result = validate_v1_1(&firehose_commit, &pubkey).await; assert!( result.is_err(), "Validation should fail when commit has wrong MST root" @@ -905,7 +760,7 @@ mod tests { firehose_commit.repo = wrong_did; - let result = firehose_commit.validate_v1_1(&pubkey).await; + let result = validate_v1_1(&firehose_commit, &pubkey).await; assert!( result.is_err(), "Validation should fail with mismatched DID" @@ -952,7 +807,7 @@ mod tests { .await .unwrap(); - let result = firehose_commit.validate_v1_1(&wrong_pubkey).await; + let result = validate_v1_1(&firehose_commit, &wrong_pubkey).await; assert!( result.is_err(), "Validation should fail with wrong public key" @@ -993,7 +848,7 @@ mod tests { // Strip prev_data to make it invalid for v1.1 firehose_commit.prev_data = None; - let result = firehose_commit.validate_v1_1(&pubkey).await; + let result = validate_v1_1(&firehose_commit, &pubkey).await; assert!( result.is_err(), "v1.1 validation should fail without prev_data" @@ -1040,7 +895,7 @@ mod tests { // Use wrong prev_data CID (point to commit instead of MST root) firehose_commit.prev_data = Some(firehose_commit.commit.clone()); - let result = firehose_commit.validate_v1_1(&pubkey).await; + let result = validate_v1_1(&firehose_commit, &pubkey).await; assert!( result.is_err(), "Validation should fail with wrong prev_data CID" diff --git a/crates/jacquard-repo/src/mst/diff.rs b/crates/jacquard-repo/src/mst/diff.rs index 49556cd1..a7bf5927 100644 --- a/crates/jacquard-repo/src/mst/diff.rs +++ b/crates/jacquard-repo/src/mst/diff.rs @@ -170,6 +170,7 @@ impl MstDiff { path: key.as_str().into(), cid: Some(CidLink::from(*cid)), prev: None, + extra_data: None, }); } @@ -180,6 +181,7 @@ impl MstDiff { path: key.as_str().into(), cid: Some(CidLink::from(*new_cid)), prev: Some(CidLink::from(*old_cid)), + extra_data: None, }); } @@ -190,6 +192,7 @@ impl MstDiff { path: key.as_str().into(), cid: None, // null for deletes prev: Some(CidLink::from(*old_cid)), + extra_data: None, }); } @@ -220,12 +223,17 @@ impl Mst { // Remove duplicate blocks: nodes that appear in both new_mst_blocks and removed_mst_blocks // are unchanged nodes that were traversed during the diff but shouldn't be counted as created/deleted. // This happens when we step into subtrees with different parent CIDs but encounter identical child nodes. - let created_set: std::collections::HashSet<_> = diff.new_mst_blocks.keys().copied().collect(); - let removed_set: std::collections::HashSet<_> = diff.removed_mst_blocks.iter().copied().collect(); - let duplicates: std::collections::HashSet<_> = created_set.intersection(&removed_set).copied().collect(); - - diff.new_mst_blocks.retain(|cid, _| !duplicates.contains(cid)); - diff.removed_mst_blocks.retain(|cid| !duplicates.contains(cid)); + let created_set: std::collections::HashSet<_> = + diff.new_mst_blocks.keys().copied().collect(); + let removed_set: std::collections::HashSet<_> = + diff.removed_mst_blocks.iter().copied().collect(); + let duplicates: std::collections::HashSet<_> = + created_set.intersection(&removed_set).copied().collect(); + + diff.new_mst_blocks + .retain(|cid, _| !duplicates.contains(cid)); + diff.removed_mst_blocks + .retain(|cid| !duplicates.contains(cid)); Ok(diff) } @@ -420,8 +428,12 @@ async fn serialize_and_track_mst( // Serialize the MST node let entries = tree.get_entries().await?; let node_data = serialize_node_data(&entries).await?; - let cbor = serde_ipld_dagcbor::to_vec(&node_data) - .map_err(|e| RepoError::serialization(e).with_context(format!("serializing MST node for diff tracking: {}", tree_cid)))?; + let cbor = serde_ipld_dagcbor::to_vec(&node_data).map_err(|e| { + RepoError::serialization(e).with_context(format!( + "serializing MST node for diff tracking: {}", + tree_cid + )) + })?; // Track the serialized block diff.new_mst_blocks.insert(tree_cid, Bytes::from(cbor)); diff --git a/crates/jacquard-repo/src/repo.rs b/crates/jacquard-repo/src/repo.rs index e63628ae..6fc20c8a 100644 --- a/crates/jacquard-repo/src/repo.rs +++ b/crates/jacquard-repo/src/repo.rs @@ -82,7 +82,7 @@ impl CommitData { repo: repo.clone().into_static(), rev: self.rev.clone(), seq, - since: self.since.clone().unwrap_or_else(|| self.rev.clone()), + since: Some(self.since.clone().unwrap_or_else(|| self.rev.clone())), time, commit: CidLink::from(self.cid), blocks: blocks_car.into(), @@ -91,6 +91,7 @@ impl CommitData { blobs, too_big: false, rebase: false, + extra_data: None, }) } } diff --git a/crates/jacquard-repo/tests/large_proof_tests.rs b/crates/jacquard-repo/tests/large_proof_tests.rs index e0b9a261..0e7ec6ab 100644 --- a/crates/jacquard-repo/tests/large_proof_tests.rs +++ b/crates/jacquard-repo/tests/large_proof_tests.rs @@ -10,6 +10,7 @@ use jacquard_common::types::tid::Ticker; use jacquard_common::types::value::RawData; use jacquard_repo::Repository; use jacquard_repo::car::read_car_header; +use jacquard_repo::commit::firehose::validate_v1_1; use jacquard_repo::mst::RecordWriteOp; use jacquard_repo::storage::{BlockStore, MemoryBlockStore}; use rand::Rng; @@ -224,8 +225,7 @@ async fn test_stress_random_operations() { .await .unwrap(); - firehose_commit - .validate_v1_1(&pubkey) + validate_v1_1(&firehose_commit, &pubkey) .await .expect("Initial batch should validate"); @@ -266,8 +266,7 @@ async fn test_stress_random_operations() { .await .unwrap(); - firehose_commit - .validate_v1_1(&pubkey) + validate_v1_1(&firehose_commit, &pubkey) .await .unwrap_or_else(|e| { eprintln!( @@ -336,7 +335,7 @@ async fn test_stress_large_batches() { .await .unwrap(); - firehose_commit.validate_v1_1(&pubkey).await.unwrap(); + validate_v1_1(&firehose_commit, &pubkey).await.unwrap(); for batch_num in 1..=5000 { let batch_size = rng.gen_range(1..=20); @@ -355,8 +354,7 @@ async fn test_stress_large_batches() { .await .unwrap(); - firehose_commit - .validate_v1_1(&pubkey) + validate_v1_1(&firehose_commit, &pubkey) .await .unwrap_or_else(|e| { panic!( @@ -441,8 +439,7 @@ async fn test_stress_with_fixture() { .await .unwrap(); - firehose_commit - .validate_v1_1(&pubkey) + validate_v1_1(&firehose_commit, &pubkey) .await .unwrap_or_else(|e| panic!("Fixture validation failed at batch {}: {}", batch_num, e)); }