diff --git a/Cargo.lock b/Cargo.lock index 7ac7952..55f82ec 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7455,7 +7455,7 @@ dependencies = [ [[package]] name = "tranquil-api" -version = "0.5.5" +version = "0.5.6" dependencies = [ "anyhow", "axum", @@ -7506,7 +7506,7 @@ dependencies = [ [[package]] name = "tranquil-auth" -version = "0.5.5" +version = "0.5.6" dependencies = [ "anyhow", "base32", @@ -7529,7 +7529,7 @@ dependencies = [ [[package]] name = "tranquil-cache" -version = "0.5.5" +version = "0.5.6" dependencies = [ "async-trait", "base64 0.22.1", @@ -7543,7 +7543,7 @@ dependencies = [ [[package]] name = "tranquil-comms" -version = "0.5.5" +version = "0.5.6" dependencies = [ "async-trait", "base64 0.22.1", @@ -7561,7 +7561,7 @@ dependencies = [ [[package]] name = "tranquil-config" -version = "0.5.5" +version = "0.5.6" dependencies = [ "confique", "serde", @@ -7569,7 +7569,7 @@ dependencies = [ [[package]] name = "tranquil-crypto" -version = "0.5.5" +version = "0.5.6" dependencies = [ "aes-gcm", "base64 0.22.1", @@ -7585,7 +7585,7 @@ dependencies = [ [[package]] name = "tranquil-db" -version = "0.5.5" +version = "0.5.6" dependencies = [ "async-trait", "chrono", @@ -7602,7 +7602,7 @@ dependencies = [ [[package]] name = "tranquil-db-traits" -version = "0.5.5" +version = "0.5.6" dependencies = [ "async-trait", "base64 0.22.1", @@ -7618,7 +7618,7 @@ dependencies = [ [[package]] name = "tranquil-infra" -version = "0.5.5" +version = "0.5.6" dependencies = [ "async-trait", "bytes", @@ -7629,7 +7629,7 @@ dependencies = [ [[package]] name = "tranquil-lexicon" -version = "0.5.5" +version = "0.5.6" dependencies = [ "chrono", "futures", @@ -7648,7 +7648,7 @@ dependencies = [ [[package]] name = "tranquil-oauth" -version = "0.5.5" +version = "0.5.6" dependencies = [ "anyhow", "axum", @@ -7671,7 +7671,7 @@ dependencies = [ [[package]] name = "tranquil-oauth-server" -version = "0.5.5" +version = "0.5.6" dependencies = [ "axum", "base64 0.22.1", @@ -7704,7 +7704,7 @@ dependencies = [ [[package]] name = "tranquil-pds" -version = "0.5.5" +version = "0.5.6" dependencies = [ "aes-gcm", "anyhow", @@ -7796,7 +7796,7 @@ dependencies = [ [[package]] name = "tranquil-repo" -version = "0.5.5" +version = "0.5.6" dependencies = [ "bytes", "cid", @@ -7808,7 +7808,7 @@ dependencies = [ [[package]] name = "tranquil-ripple" -version = "0.5.5" +version = "0.5.6" dependencies = [ "async-trait", "backon", @@ -7833,7 +7833,7 @@ dependencies = [ [[package]] name = "tranquil-scopes" -version = "0.5.5" +version = "0.5.6" dependencies = [ "axum", "futures", @@ -7849,7 +7849,7 @@ dependencies = [ [[package]] name = "tranquil-server" -version = "0.5.5" +version = "0.5.6" dependencies = [ "axum", "clap", @@ -7870,7 +7870,7 @@ dependencies = [ [[package]] name = "tranquil-signal" -version = "0.5.5" +version = "0.5.6" dependencies = [ "async-trait", "chrono", @@ -7893,7 +7893,7 @@ dependencies = [ [[package]] name = "tranquil-storage" -version = "0.5.5" +version = "0.5.6" dependencies = [ "async-trait", "aws-config", @@ -7910,7 +7910,7 @@ dependencies = [ [[package]] name = "tranquil-store" -version = "0.5.5" +version = "0.5.6" dependencies = [ "async-trait", "bytes", @@ -7959,7 +7959,7 @@ dependencies = [ [[package]] name = "tranquil-sync" -version = "0.5.5" +version = "0.5.6" dependencies = [ "anyhow", "axum", @@ -7981,7 +7981,7 @@ dependencies = [ [[package]] name = "tranquil-types" -version = "0.5.5" +version = "0.5.6" dependencies = [ "chrono", "cid", diff --git a/Cargo.toml b/Cargo.toml index 9117f0b..fcb7bcc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -26,7 +26,7 @@ members = [ ] [workspace.package] -version = "0.5.5" +version = "0.5.6" edition = "2024" license = "AGPL-3.0-or-later" diff --git a/crates/tranquil-config/src/lib.rs b/crates/tranquil-config/src/lib.rs index 8583872..93e7a34 100644 --- a/crates/tranquil-config/src/lib.rs +++ b/crates/tranquil-config/src/lib.rs @@ -65,6 +65,9 @@ pub fn ensure_test_defaults() { if env::var("ENABLE_PDS_HOSTED_DID_WEB").is_err() { env::set_var("ENABLE_PDS_HOSTED_DID_WEB", "true"); } + if env::var("TRANQUIL_LEXICON_OFFLINE").is_err() { + env::set_var("TRANQUIL_LEXICON_OFFLINE", "1"); + } } TranquilConfig::builder() .env() diff --git a/crates/tranquil-pds/src/repo_ops.rs b/crates/tranquil-pds/src/repo_ops.rs index 2171890..5845788 100644 --- a/crates/tranquil-pds/src/repo_ops.rs +++ b/crates/tranquil-pds/src/repo_ops.rs @@ -6,17 +6,19 @@ use crate::types::{Did, Handle, Nsid, Rkey}; use backon::{ExponentialBuilder, Retryable}; use bytes::Bytes; use cid::Cid; +use jacquard_common::smol_str::SmolStr; use jacquard_common::types::{integer::LimitedU32, string::Tid}; use jacquard_repo::commit::Commit; -use jacquard_repo::mst::Mst; use jacquard_repo::mst::util::compute_cid; +use jacquard_repo::mst::{Mst, VerifiedWriteOp}; use jacquard_repo::storage::BlockStore; use k256::ecdsa::SigningKey; use serde_json::{Value, json}; +use std::collections::{BTreeMap, HashSet}; use std::str::FromStr; use std::sync::Arc; use tokio::sync::OwnedMutexGuard; -use tracing::error; +use tracing::{error, warn}; use tranquil_db_traits::SequenceNumber; use uuid::Uuid; @@ -236,13 +238,89 @@ pub async fn finalize_repo_write( ApiError::InternalError(None) })?; - let block_bytes = ctx.tracking_store.take_written_blocks(); + let written_bytes = ctx.tracking_store.take_written_blocks(); + let new_tree_cids: Vec = written_bytes.keys().copied().collect(); - let storage_for_diff = Arc::new(ctx.tracking_store.clone()); - let original_settled = Mst::load(storage_for_diff.clone(), ctx.prev_data_cid, None); - let new_settled = Mst::load(storage_for_diff, new_mst_root, None); + let storage_for_proof = Arc::new(ctx.tracking_store.clone()); + let original_settled = Mst::load(storage_for_proof.clone(), ctx.prev_data_cid, None); + let new_settled = Mst::load(storage_for_proof.clone(), new_mst_root, None); - let new_tree_cids: Vec = block_bytes.keys().copied().collect(); + let mut inverse_trace = new_settled.clone(); + let mut non_invertible: Vec = Vec::new(); + let mut invert_errors: Vec = Vec::new(); + for op in params.ops.iter() { + let (collection, rkey) = match op { + RecordOp::Create { + collection, rkey, .. + } + | RecordOp::Update { + collection, rkey, .. + } + | RecordOp::Delete { + collection, rkey, .. + } => (collection, rkey), + }; + let key = SmolStr::new(format!("{}/{}", collection, rkey)); + let verified = match op { + RecordOp::Create { cid, .. } => VerifiedWriteOp::Create { + key, + cid: *cid.as_cid(), + }, + RecordOp::Update { cid, prev, .. } => VerifiedWriteOp::Update { + key, + cid: *cid.as_cid(), + prev: *prev.as_cid(), + }, + RecordOp::Delete { prev, .. } => VerifiedWriteOp::Delete { + key, + prev: *prev.as_cid(), + }, + }; + match inverse_trace.invert_op(verified.clone()).await { + Ok(true) => {} + Ok(false) => non_invertible.push(format!("{:?}", verified)), + Err(e) => invert_errors.push(format!("{:?} -> {:?}", verified, e)), + } + } + if !non_invertible.is_empty() { + warn!( + user_id = %params.user_id, + count = non_invertible.len(), + ops = ?non_invertible, + "firehose proof walk: ops not invertible on new MST, consumer will reject frame" + ); + } + if !invert_errors.is_empty() { + warn!( + user_id = %params.user_id, + count = invert_errors.len(), + failures = ?invert_errors, + "firehose proof walk: invert_op errored, cover blocks may be incomplete" + ); + } + + let read_cid_set: HashSet = ctx.tracking_store.get_read_cids().into_iter().collect(); + let missing_read_cids: Vec = read_cid_set + .iter() + .copied() + .filter(|cid| !written_bytes.contains_key(cid)) + .collect(); + let mut relevant: BTreeMap = BTreeMap::new(); + if !missing_read_cids.is_empty() { + let fetched = ctx + .tracking_store + .get_many(&missing_read_cids) + .await + .map_err(|e| { + error!("fetch cover read bytes: {e}"); + ApiError::InternalError(None) + })?; + for (cid, maybe) in missing_read_cids.into_iter().zip(fetched) { + if let Some(bytes) = maybe { + relevant.insert(cid, bytes); + } + } + } let obsolete_cids = match original_settled.diff(&new_settled).await { Ok(diff) => { @@ -263,6 +341,9 @@ pub async fn finalize_repo_write( } }; + let mut block_bytes = written_bytes; + block_bytes.extend(relevant); + let result = commit_and_log( state, CommitParams { diff --git a/crates/tranquil-pds/tests/mst_inductive_firehose.rs b/crates/tranquil-pds/tests/mst_inductive_firehose.rs new file mode 100644 index 0000000..713ce16 --- /dev/null +++ b/crates/tranquil-pds/tests/mst_inductive_firehose.rs @@ -0,0 +1,640 @@ +mod common; +mod mst_verify; + +use std::collections::BTreeMap; +use std::str::FromStr; +use std::sync::Arc; + +use cid::Cid; +use common::*; +use jacquard_common::smol_str::SmolStr; +use jacquard_repo::commit::Commit; +use jacquard_repo::mst::{Mst, VerifiedWriteOp}; +use jacquard_repo::storage::{BlockStore, MemoryBlockStore}; +use mst_verify::{extract_event_blocks, inline_to_store}; +use reqwest::StatusCode; +use serde_json::{Value, json}; +use tranquil_db_traits::{RepoEventType, SequenceNumber, SequencedEvent}; +use tranquil_types::Did; + +async fn new_commit_data_cid( + storage: &Arc, + commit_cid: &Cid, +) -> Result { + let commit_bytes = storage + .get(commit_cid) + .await + .map_err(|e| format!("get commit: {e:?}"))? + .ok_or_else(|| format!("CAR missing commit block {commit_cid}"))?; + let commit = Commit::from_cbor(&commit_bytes).map_err(|e| format!("parse commit: {e:?}"))?; + Ok(*commit.data()) +} + +fn ops_json(event: &SequencedEvent) -> Result<&Vec, String> { + event + .ops + .as_ref() + .and_then(|v| v.as_array()) + .ok_or_else(|| "event.ops not an array".into()) +} + +fn parse_op_to_verified(op: &Value) -> Result { + let action = op["action"].as_str().ok_or("op.action missing")?; + let path = op["path"].as_str().ok_or("op.path missing")?; + let key = SmolStr::new(path); + match action { + "create" => { + let cid_str = op["cid"].as_str().ok_or("create missing cid")?; + let cid = Cid::from_str(cid_str).map_err(|e| format!("parse cid: {e:?}"))?; + Ok(VerifiedWriteOp::Create { key, cid }) + } + "update" => { + let cid_str = op["cid"].as_str().ok_or("update missing cid")?; + let cid = Cid::from_str(cid_str).map_err(|e| format!("parse cid: {e:?}"))?; + let prev_str = op["prev"].as_str().ok_or("update missing prev")?; + let prev = Cid::from_str(prev_str).map_err(|e| format!("parse prev: {e:?}"))?; + Ok(VerifiedWriteOp::Update { key, cid, prev }) + } + "delete" => { + let prev_str = op["prev"].as_str().ok_or("delete missing prev")?; + let prev = Cid::from_str(prev_str).map_err(|e| format!("parse prev: {e:?}"))?; + Ok(VerifiedWriteOp::Delete { key, prev }) + } + other => Err(format!("unknown op action: {other}")), + } +} + +async fn verify_inductive_forward(event: &SequencedEvent) -> Result<(Cid, Cid), String> { + let prev_data_cid = event + .prev_data_cid + .as_ref() + .and_then(|c| c.to_cid()) + .ok_or_else(|| "event missing prev_data_cid".to_string())?; + let commit_cid = event + .commit_cid + .as_ref() + .and_then(|c| c.to_cid()) + .ok_or_else(|| "event missing commit_cid".to_string())?; + + let storage = inline_to_store(extract_event_blocks(event)?); + let expected_new_data = new_commit_data_cid(&storage, &commit_cid).await?; + + let mut mst = Mst::load(storage.clone(), prev_data_cid, None); + for op_value in ops_json(event)? { + let action = op_value["action"].as_str().ok_or("op.action missing")?; + let path = op_value["path"].as_str().ok_or("op.path missing")?; + match action { + "create" | "update" => { + let cid = Cid::from_str(op_value["cid"].as_str().ok_or("op.cid missing")?) + .map_err(|e| format!("parse op.cid: {e:?}"))?; + mst = mst + .add(path, cid) + .await + .map_err(|e| format!("mst.add({path}): {e:?}"))?; + } + "delete" => { + mst = mst + .delete(path) + .await + .map_err(|e| format!("mst.delete({path}): {e:?}"))?; + } + other => return Err(format!("unknown op action: {other}")), + } + } + let computed = mst + .persist() + .await + .map_err(|e| format!("mst.persist: {e:?}"))?; + Ok((expected_new_data, computed)) +} + +async fn verify_inductive_inverse(event: &SequencedEvent) -> Result<(Cid, Cid), String> { + let prev_data_cid = event + .prev_data_cid + .as_ref() + .and_then(|c| c.to_cid()) + .ok_or_else(|| "event missing prev_data_cid".to_string())?; + let commit_cid = event + .commit_cid + .as_ref() + .and_then(|c| c.to_cid()) + .ok_or_else(|| "event missing commit_cid".to_string())?; + + let storage = inline_to_store(extract_event_blocks(event)?); + let new_data_cid = new_commit_data_cid(&storage, &commit_cid).await?; + + let mut mst = Mst::load(storage.clone(), new_data_cid, None); + for op_value in ops_json(event)? { + let verified = parse_op_to_verified(op_value)?; + let inverted = mst + .invert_op(verified.clone()) + .await + .map_err(|e| format!("invert_op({verified:?}): {e:?}"))?; + if !inverted { + return Err(format!("op not invertible: {verified:?}")); + } + } + let computed_prev = mst + .get_pointer() + .await + .map_err(|e| format!("get_pointer: {e:?}"))?; + Ok((prev_data_cid, computed_prev)) +} + +fn report_failures(total: usize, failures: &[String], mode: &str) { + assert!( + failures.is_empty(), + "{} of {total} {mode} commit events failed inductive verification:\n - {}", + failures.len(), + failures.join("\n - "), + ); +} + +async fn apply_writes_batch(client: &reqwest::Client, token: &str, did: &str, writes: Vec) { + let payload = json!({ "repo": did, "writes": writes }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(token) + .json(&payload) + .send() + .await + .expect("applyWrites request failed"); + assert_eq!( + res.status(), + StatusCode::OK, + "applyWrites failed: {:?}", + res.text().await + ); +} + +async fn create_record(client: &reqwest::Client, token: &str, did: &str, col: &str, rkey: &str) { + let now = chrono::Utc::now().to_rfc3339(); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.createRecord", + base_url().await + )) + .bearer_auth(token) + .json(&json!({ + "repo": did, + "collection": col, + "rkey": rkey, + "record": { + "$type": col, + "text": format!("post {rkey}"), + "createdAt": now, + } + })) + .send() + .await + .expect("createRecord request failed"); + assert_eq!(res.status(), StatusCode::OK, "createRecord failed"); +} + +async fn put_record( + client: &reqwest::Client, + token: &str, + did: &str, + col: &str, + rkey: &str, + text: &str, +) { + let now = chrono::Utc::now().to_rfc3339(); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.putRecord", + base_url().await + )) + .bearer_auth(token) + .json(&json!({ + "repo": did, + "collection": col, + "rkey": rkey, + "record": { + "$type": col, + "text": text, + "createdAt": now, + } + })) + .send() + .await + .expect("putRecord request failed"); + assert_eq!(res.status(), StatusCode::OK, "putRecord failed"); +} + +async fn delete_record(client: &reqwest::Client, token: &str, did: &str, col: &str, rkey: &str) { + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.deleteRecord", + base_url().await + )) + .bearer_auth(token) + .json(&json!({ "repo": did, "collection": col, "rkey": rkey })) + .send() + .await + .expect("deleteRecord request failed"); + assert_eq!(res.status(), StatusCode::OK, "deleteRecord failed"); +} + +const COLLECTION: &str = "app.bsky.feed.post"; +fn rkey_for(prefix: &str, i: usize) -> String { + format!("3k{prefix}{:08}", i) +} + +async fn our_commit_events(did: &str) -> Vec { + let repos = get_test_repos().await; + let typed_did = Did::new(did.to_string()).unwrap(); + let events = repos + .repo + .get_events_since_seq(SequenceNumber::ZERO, None) + .await + .expect("get_events_since_seq"); + events + .into_iter() + .filter(|e| e.did == typed_did && e.event_type == RepoEventType::Commit) + .collect() +} + +#[tokio::test] +async fn inductive_forward_verifies_delete_commits() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + + let now = chrono::Utc::now().to_rfc3339(); + const N_CREATE: usize = 200; + + let all_writes: Vec = (0..N_CREATE) + .map(|i| { + json!({ + "$type": "com.atproto.repo.applyWrites#create", + "collection": COLLECTION, + "rkey": rkey_for("del", i), + "value": { + "$type": COLLECTION, + "text": format!("record {i}"), + "createdAt": now, + } + }) + }) + .collect(); + for chunk in all_writes.chunks(50) { + apply_writes_batch(&client, &token, &did, chunk.to_vec()).await; + } + + let delete_indices: Vec = (10..N_CREATE).step_by(7).collect(); + for i in &delete_indices { + delete_record(&client, &token, &did, COLLECTION, &rkey_for("del", *i)).await; + } + + let our = our_commit_events(&did).await; + let delete_events: Vec<&SequencedEvent> = our + .iter() + .filter(|e| { + ops_json(e) + .map(|arr| arr.iter().any(|op| op["action"].as_str() == Some("delete"))) + .unwrap_or(false) + }) + .collect(); + assert_eq!(delete_events.len(), delete_indices.len()); + + let mut failures = Vec::new(); + for e in &delete_events { + match verify_inductive_forward(e).await { + Ok((exp, got)) if exp == got => {} + Ok((exp, got)) => failures.push(format!( + "seq={}: root mismatch exp={exp} got={got}", + e.seq.as_i64() + )), + Err(msg) => failures.push(format!("seq={}: {msg}", e.seq.as_i64())), + } + } + report_failures(delete_events.len(), &failures, "delete forward"); +} + +#[tokio::test] +async fn inductive_forward_verifies_create_commits() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + for i in 0..60usize { + create_record(&client, &token, &did, COLLECTION, &rkey_for("cre", i)).await; + } + + let our = our_commit_events(&did).await; + let create_events: Vec<&SequencedEvent> = our + .iter() + .filter(|e| { + ops_json(e) + .map(|arr| arr.iter().all(|op| op["action"].as_str() == Some("create"))) + .unwrap_or(false) + && e.prev_data_cid.is_some() + }) + .collect(); + assert!(!create_events.is_empty()); + + let mut failures = Vec::new(); + for e in &create_events { + match verify_inductive_forward(e).await { + Ok((exp, got)) if exp == got => {} + Ok((exp, got)) => failures.push(format!( + "seq={}: root mismatch exp={exp} got={got}", + e.seq.as_i64() + )), + Err(msg) => failures.push(format!("seq={}: {msg}", e.seq.as_i64())), + } + } + report_failures(create_events.len(), &failures, "create forward"); +} + +#[tokio::test] +async fn inductive_forward_verifies_update_commits() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + + let now = chrono::Utc::now().to_rfc3339(); + let creates: Vec = (0..80) + .map(|i| { + json!({ + "$type": "com.atproto.repo.applyWrites#create", + "collection": COLLECTION, + "rkey": rkey_for("upd", i), + "value": { + "$type": COLLECTION, + "text": format!("original {i}"), + "createdAt": now, + } + }) + }) + .collect(); + for chunk in creates.chunks(40) { + apply_writes_batch(&client, &token, &did, chunk.to_vec()).await; + } + + for i in (0..80).step_by(3) { + put_record( + &client, + &token, + &did, + COLLECTION, + &rkey_for("upd", i), + &format!("updated {i}"), + ) + .await; + } + + let our = our_commit_events(&did).await; + let update_events: Vec<&SequencedEvent> = our + .iter() + .filter(|e| { + ops_json(e) + .map(|arr| arr.iter().any(|op| op["action"].as_str() == Some("update"))) + .unwrap_or(false) + }) + .collect(); + assert!(!update_events.is_empty()); + + let mut failures = Vec::new(); + for e in &update_events { + match verify_inductive_forward(e).await { + Ok((exp, got)) if exp == got => {} + Ok((exp, got)) => failures.push(format!( + "seq={}: root mismatch exp={exp} got={got}", + e.seq.as_i64() + )), + Err(msg) => failures.push(format!("seq={}: {msg}", e.seq.as_i64())), + } + } + report_failures(update_events.len(), &failures, "update forward"); +} + +#[tokio::test] +async fn inductive_forward_verifies_mixed_applywrites() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + + let now = chrono::Utc::now().to_rfc3339(); + let seed: Vec = (0..120) + .map(|i| { + json!({ + "$type": "com.atproto.repo.applyWrites#create", + "collection": COLLECTION, + "rkey": rkey_for("mix", i), + "value": { + "$type": COLLECTION, + "text": format!("seed {i}"), + "createdAt": now, + } + }) + }) + .collect(); + for chunk in seed.chunks(40) { + apply_writes_batch(&client, &token, &did, chunk.to_vec()).await; + } + + let mixed: Vec = (0..40) + .flat_map(|i| { + vec![ + json!({ + "$type": "com.atproto.repo.applyWrites#create", + "collection": COLLECTION, + "rkey": rkey_for("mxc", i), + "value": { + "$type": COLLECTION, + "text": format!("new {i}"), + "createdAt": now, + } + }), + json!({ + "$type": "com.atproto.repo.applyWrites#update", + "collection": COLLECTION, + "rkey": rkey_for("mix", i), + "value": { + "$type": COLLECTION, + "text": format!("updated-mix {i}"), + "createdAt": now, + } + }), + json!({ + "$type": "com.atproto.repo.applyWrites#delete", + "collection": COLLECTION, + "rkey": rkey_for("mix", i + 60), + }), + ] + }) + .collect(); + apply_writes_batch(&client, &token, &did, mixed).await; + + let our = our_commit_events(&did).await; + let last = our + .iter() + .rfind(|e| e.prev_data_cid.is_some()) + .expect("at least one non-genesis commit"); + + let actions: Vec<&str> = ops_json(last) + .unwrap() + .iter() + .filter_map(|op| op["action"].as_str()) + .collect(); + assert!(actions.contains(&"create")); + assert!(actions.contains(&"update")); + assert!(actions.contains(&"delete")); + + let (exp, got) = verify_inductive_forward(last) + .await + .expect("mixed applyWrites forward verify"); + assert_eq!(exp, got, "mixed applyWrites commit forward-verify mismatch"); +} + +#[tokio::test] +async fn inductive_inverse_verifies_every_commit() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + + let now = chrono::Utc::now().to_rfc3339(); + let seed: Vec = (0..100) + .map(|i| { + json!({ + "$type": "com.atproto.repo.applyWrites#create", + "collection": COLLECTION, + "rkey": rkey_for("inv", i), + "value": { + "$type": COLLECTION, + "text": format!("seed {i}"), + "createdAt": now, + } + }) + }) + .collect(); + for chunk in seed.chunks(50) { + apply_writes_batch(&client, &token, &did, chunk.to_vec()).await; + } + for i in (0..100).step_by(5) { + put_record( + &client, + &token, + &did, + COLLECTION, + &rkey_for("inv", i), + &format!("upd {i}"), + ) + .await; + } + for i in (2..100).step_by(11) { + delete_record(&client, &token, &did, COLLECTION, &rkey_for("inv", i)).await; + } + + let our = our_commit_events(&did).await; + let non_genesis: Vec<&SequencedEvent> = our + .iter() + .filter(|e| e.prev_data_cid.is_some() && ops_json(e).is_ok()) + .collect(); + assert!(!non_genesis.is_empty()); + + let mut failures = Vec::new(); + for e in &non_genesis { + match verify_inductive_inverse(e).await { + Ok((exp, got)) if exp == got => {} + Ok((exp, got)) => failures.push(format!( + "seq={}: inverse root mismatch exp={exp} got={got}", + e.seq.as_i64() + )), + Err(msg) => failures.push(format!("seq={}: {msg}", e.seq.as_i64())), + } + } + report_failures(non_genesis.len(), &failures, "any inverse"); +} + +#[tokio::test] +async fn prev_cid_chain_walks_to_genesis() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + for i in 0..8 { + create_record(&client, &token, &did, COLLECTION, &rkey_for("cha", i)).await; + } + + let our = our_commit_events(&did).await; + assert!(our.len() >= 2); + + let last = our.last().unwrap(); + let mut current_prev: Option = last.prev_cid.as_ref().and_then(|c| c.to_cid()); + let head_commit_cid = last + .commit_cid + .as_ref() + .and_then(|c| c.to_cid()) + .expect("head commit_cid"); + + let by_commit: BTreeMap = our + .iter() + .filter_map(|e| { + e.commit_cid + .as_ref() + .and_then(|c| c.to_cid()) + .map(|c| (c, e)) + }) + .collect(); + + let mut visited = 1; + while let Some(prev) = current_prev { + let e = by_commit + .get(&prev) + .unwrap_or_else(|| panic!("prev commit {prev} missing from event list")); + visited += 1; + current_prev = e.prev_cid.as_ref().and_then(|c| c.to_cid()); + } + assert!( + visited >= 2, + "chain too short: visited={visited}, head_commit={head_commit_cid}" + ); + assert_eq!( + visited, + our.len(), + "chain did not reach genesis: walked {visited}, have {}", + our.len() + ); +} + +#[tokio::test] +async fn record_bytes_present_in_car_for_creates() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + + let now = chrono::Utc::now().to_rfc3339(); + let writes: Vec = (0..5) + .map(|i| { + json!({ + "$type": "com.atproto.repo.applyWrites#create", + "collection": COLLECTION, + "rkey": rkey_for("rec", i), + "value": { + "$type": COLLECTION, + "text": format!("rec {i}"), + "createdAt": now, + } + }) + }) + .collect(); + apply_writes_batch(&client, &token, &did, writes).await; + + let our = our_commit_events(&did).await; + let latest = our.iter().rfind(|e| e.prev_data_cid.is_some()).unwrap(); + + let inline = extract_event_blocks(latest).unwrap(); + let have_cids: std::collections::HashSet = inline + .iter() + .map(|b| Cid::read_bytes(b.cid_bytes.as_slice()).unwrap()) + .collect(); + + for op in ops_json(latest).unwrap() { + if op["action"].as_str() == Some("create") + && let Some(cid_str) = op["cid"].as_str() + { + let cid = Cid::from_str(cid_str).unwrap(); + assert!( + have_cids.contains(&cid), + "create op record CID {cid} not present in CAR inline blocks" + ); + } + } +} diff --git a/crates/tranquil-pds/tests/mst_verify/mod.rs b/crates/tranquil-pds/tests/mst_verify/mod.rs new file mode 100644 index 0000000..9938c9b --- /dev/null +++ b/crates/tranquil-pds/tests/mst_verify/mod.rs @@ -0,0 +1,26 @@ +use std::collections::BTreeMap; +use std::sync::Arc; + +use bytes::Bytes; +use cid::Cid; +use jacquard_repo::storage::MemoryBlockStore; +use tranquil_db_traits::{EventBlockInline, EventBlocks, SequencedEvent}; + +pub fn extract_event_blocks(event: &SequencedEvent) -> Result<&[EventBlockInline], String> { + match event.blocks.as_ref() { + Some(EventBlocks::Inline(v)) => Ok(v.as_slice()), + Some(EventBlocks::LegacyCids(_)) => Err("legacy cids, not inline".into()), + None => Err("event missing blocks".into()), + } +} + +pub fn inline_to_store(inline: &[EventBlockInline]) -> Arc { + let map: BTreeMap = inline + .iter() + .map(|b| { + let cid = Cid::read_bytes(b.cid_bytes.as_slice()).expect("valid cid bytes"); + (cid, Bytes::from(b.data.clone())) + }) + .collect(); + Arc::new(MemoryBlockStore::new_from_blocks(map)) +} diff --git a/crates/tranquil-signal/src/client.rs b/crates/tranquil-signal/src/client.rs index d3c5dbf..7a7f167 100644 --- a/crates/tranquil-signal/src/client.rs +++ b/crates/tranquil-signal/src/client.rs @@ -122,11 +122,7 @@ pub struct MessageTooLong { impl fmt::Display for MessageTooLong { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - write!( - f, - "message body is {} bytes, max {}", - self.len, self.max - ) + write!(f, "message body is {} bytes, max {}", self.len, self.max) } } diff --git a/crates/tranquil-store/src/gauntlet/runner.rs b/crates/tranquil-store/src/gauntlet/runner.rs index 5c65901..9eb5100 100644 --- a/crates/tranquil-store/src/gauntlet/runner.rs +++ b/crates/tranquil-store/src/gauntlet/runner.rs @@ -577,7 +577,8 @@ where oracle.record_crash(); } shutdown_harness(&mut harness); - match reopen_with_recovery(&mut open, &mut crash, tolerate_op_errors, reopen_backoff).await { + match reopen_with_recovery(&mut open, &mut crash, tolerate_op_errors, reopen_backoff).await + { Ok(reopened) => { harness = Some(reopened); let n = restarts_counter.fetch_add(1, Ordering::Relaxed) + 1; @@ -613,7 +614,8 @@ where crash(); oracle.record_crash(); shutdown_harness(&mut harness); - match reopen_with_recovery(&mut open, &mut crash, tolerate_op_errors, reopen_backoff).await { + match reopen_with_recovery(&mut open, &mut crash, tolerate_op_errors, reopen_backoff).await + { Ok(reopened) => harness = Some(reopened), Err(detail) => { violations.push(InvariantViolation { @@ -654,7 +656,8 @@ where { let pre_snapshot = snapshot_block_index(&live.store); shutdown_harness(&mut harness); - match reopen_with_recovery(&mut open, &mut crash, tolerate_op_errors, reopen_backoff).await { + match reopen_with_recovery(&mut open, &mut crash, tolerate_op_errors, reopen_backoff).await + { Ok(reopened) => { let post_snapshot = snapshot_block_index(&reopened.store); if let Some(detail) = diff_snapshots(&pre_snapshot, &post_snapshot) { @@ -1676,7 +1679,14 @@ where oracle.record_crash(); } shutdown_harness(&mut harness); - match reopen_with_recovery(&mut open, &mut crash, tolerate_op_errors, reopen_backoff).await { + match reopen_with_recovery( + &mut open, + &mut crash, + tolerate_op_errors, + reopen_backoff, + ) + .await + { Ok(reopened) => { harness = Some(reopened); let n = restarts_counter.fetch_add(1, Ordering::Relaxed) + 1; @@ -1713,7 +1723,8 @@ where crash(); oracle.record_crash(); shutdown_harness(&mut harness); - match reopen_with_recovery(&mut open, &mut crash, tolerate_op_errors, reopen_backoff).await { + match reopen_with_recovery(&mut open, &mut crash, tolerate_op_errors, reopen_backoff).await + { Ok(reopened) => harness = Some(reopened), Err(detail) => { violations.push(InvariantViolation { @@ -1755,7 +1766,8 @@ where { let pre_snapshot = snapshot_block_index(&live.store); shutdown_harness(&mut harness); - match reopen_with_recovery(&mut open, &mut crash, tolerate_op_errors, reopen_backoff).await { + match reopen_with_recovery(&mut open, &mut crash, tolerate_op_errors, reopen_backoff).await + { Ok(reopened) => { let post_snapshot = snapshot_block_index(&reopened.store); if let Some(detail) = diff_snapshots(&pre_snapshot, &post_snapshot) {