fix(tranquil-pds): firehose car carries inductive proof

Lewis: May this revision serve well! <lu5a@proton.me>
This commit is contained in:
Lewis
2026-04-26 20:11:27 +03:00
committed by Tangled
parent 0455dc20bd
commit fc817cea4c
8 changed files with 799 additions and 41 deletions
Generated
+22 -22
View File
@@ -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",
+1 -1
View File
@@ -26,7 +26,7 @@ members = [
]
[workspace.package]
version = "0.5.5"
version = "0.5.6"
edition = "2024"
license = "AGPL-3.0-or-later"
+3
View File
@@ -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()
+88 -7
View File
@@ -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<Cid> = 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<Cid> = block_bytes.keys().copied().collect();
let mut inverse_trace = new_settled.clone();
let mut non_invertible: Vec<String> = Vec::new();
let mut invert_errors: Vec<String> = 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<Cid> = ctx.tracking_store.get_read_cids().into_iter().collect();
let missing_read_cids: Vec<Cid> = read_cid_set
.iter()
.copied()
.filter(|cid| !written_bytes.contains_key(cid))
.collect();
let mut relevant: BTreeMap<Cid, Bytes> = 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 {
@@ -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<MemoryBlockStore>,
commit_cid: &Cid,
) -> Result<Cid, String> {
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<Value>, 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<VerifiedWriteOp, String> {
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<Value>) {
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<SequencedEvent> {
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<Value> = (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<usize> = (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<Value> = (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<Value> = (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<Value> = (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<Value> = (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<Cid> = 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<Cid, &SequencedEvent> = 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<Value> = (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<Cid> = 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"
);
}
}
}
@@ -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<MemoryBlockStore> {
let map: BTreeMap<Cid, Bytes> = 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))
}
+1 -5
View File
@@ -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)
}
}
+18 -6
View File
@@ -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) {