store: revalidate stored mutation sets on replay

Lewis: May this revision serve well! <lu5a@proton.me>
This commit is contained in:
Lewis
2026-07-25 08:27:40 +03:00
committed by Tangled
parent 6ed568dbfb
commit 17905115d8
3 changed files with 435 additions and 133 deletions
+102 -78
View File
@@ -212,8 +212,8 @@ impl<S: StorageIO + 'static> CommitOps<S> {
let cid_bytes = cid_link_to_bytes(&u.cid)
.map_err(|e| ApplyCommitError::Database(e.to_string()))?;
Ok(RecordMutationUpsert {
collection: u.collection.as_str().to_owned(),
rkey: u.rkey.as_str().to_owned(),
collection: u.collection.clone(),
rkey: u.rkey.clone(),
cid_bytes,
})
})
@@ -222,8 +222,8 @@ impl<S: StorageIO + 'static> CommitOps<S> {
.record_deletes
.iter()
.map(|d| RecordMutationDelete {
collection: d.collection.as_str().to_owned(),
rkey: d.rkey.as_str().to_owned(),
collection: d.collection.clone(),
rkey: d.rkey.clone(),
})
.collect(),
block_inserts: input.new_block_cids.clone(),
@@ -232,16 +232,12 @@ impl<S: StorageIO + 'static> CommitOps<S> {
.backlinks_to_add
.iter()
.map(|bl| BacklinkMutation {
uri: bl.uri.as_str().to_owned(),
uri: bl.uri.clone(),
path: path_to_discriminant(bl.path),
link_to: bl.link_to.clone(),
})
.collect(),
backlink_remove_uris: input
.backlinks_to_remove
.iter()
.map(|uri| uri.as_str().to_owned())
.collect(),
backlink_remove_uris: input.backlinks_to_remove.clone(),
};
let mutation_set_bytes = mutation_set
.serialize()
@@ -372,10 +368,11 @@ impl<S: StorageIO + 'static> CommitOps<S> {
self.scan_users_missing_prefix(
record_blobs_user_prefix,
|meta, user_id| {
let did = meta
.did
.map(Did::from)
.ok_or(MetastoreError::CorruptData("repo_meta missing did field"))?;
let did = match meta.did {
None => Err(MetastoreError::CorruptData("repo_meta missing DID field")),
Some(d) => Did::new(d)
.map_err(|_| MetastoreError::CorruptData("corrupt repo_meta did")),
}?;
Ok(UserNeedingRecordBlobsBackfill { user_id, did })
},
limit_usize,
@@ -394,7 +391,9 @@ impl<S: StorageIO + 'static> CommitOps<S> {
repo_root_cid: root_cid,
repo_rev: match meta.repo_rev.is_empty() {
true => None,
false => Some(Tid::from(meta.repo_rev)),
false => Some(Tid::new(meta.repo_rev).map_err(|_| {
MetastoreError::CorruptData("corrupt repo_meta repo_rev")
})?),
},
})
},
@@ -424,7 +423,10 @@ impl<S: StorageIO + 'static> CommitOps<S> {
let user_hash = match parse_user_hash_from_key(&key_bytes) {
Some(h) => h,
None => return Some(Err(MetastoreError::CorruptData("invalid repo_meta key"))),
None => {
tracing::warn!("skipping a repo_meta row whose key doesn't parse");
return None;
}
};
let check_prefix = make_prefix(user_hash);
@@ -442,20 +444,34 @@ impl<S: StorageIO + 'static> CommitOps<S> {
let meta = match RepoMetaValue::deserialize(&val_bytes) {
Some(v) => v,
None => {
return Some(Err(MetastoreError::CorruptData(
"invalid repo_meta value",
)));
tracing::warn!(
user_hash = user_hash.raw(),
"skipping a repo_meta row whose value doesn't decode"
);
return None;
}
};
let user_id = match self.user_hashes.get_uuid(&user_hash) {
Some(id) => id,
None => {
return Some(Err(MetastoreError::CorruptData(
"user_hash has no reverse mapping",
)));
tracing::warn!(
user_hash = user_hash.raw(),
"skipping a repo_meta row whose user_hash has no reverse mapping"
);
return None;
}
};
Some(build_result(meta, user_id))
match build_result(meta, user_id) {
Ok(v) => Some(Ok(v)),
Err(e) => {
tracing::warn!(
user_id = %user_id,
error = %e,
"skipping a repo_meta row the current validators reject"
);
None
}
}
}
}
})
@@ -528,11 +544,19 @@ mod tests {
}
fn test_did(name: &str) -> Did {
Did::from(format!("did:plc:{name}"))
Did::new(format!("did:plc:{name}")).expect("test DID is well-formed")
}
fn test_handle(name: &str) -> Handle {
Handle::from(format!("{name}.test.invalid"))
Handle::new(format!("{name}.oyster.cafe")).expect("test handle is valid")
}
fn test_nsid(name: impl Into<String>) -> Nsid {
Nsid::new(name).expect("test collection is a valid NSID")
}
fn test_rkey(key: impl Into<String>) -> Rkey {
Rkey::new(key).expect("test record key is a valid rkey")
}
fn test_rev(seq: u64) -> Tid {
@@ -582,15 +606,15 @@ mod tests {
let new_root = test_cid_link(2);
let record_cid = test_cid_link(3);
let collection = Nsid::from("app.bsky.feed.post".to_string());
let rkey = Rkey::from("3k2abc".to_string());
let collection = test_nsid("app.bsky.feed.post");
let rkey = test_rkey("3k2abc");
let input = ApplyCommitInput {
user_id,
did: did.clone(),
expected_root_cid: Some(root_cid.clone()),
new_root_cid: new_root.clone(),
new_rev: Tid::from("rev1".to_string()),
new_rev: Tid::new("3k2aaaaaaaaab").unwrap(),
new_block_cids: vec![vec![0x01, 0x02]],
obsolete_block_cids: vec![],
record_upserts: vec![tranquil_db_traits::RecordUpsert {
@@ -610,7 +634,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev1".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()),
},
};
@@ -619,7 +643,7 @@ mod tests {
let repo = h.metastore.repo_ops().get_repo(user_id).unwrap().unwrap();
assert_eq!(repo.repo_root_cid, new_root);
assert_eq!(repo.repo_rev.as_deref(), Some("rev1"));
assert_eq!(repo.repo_rev.as_deref(), Some("3k2aaaaaaaaab"));
let found_cid = h
.metastore
@@ -644,7 +668,7 @@ mod tests {
did,
expected_root_cid: Some(stale_root),
new_root_cid: new_root,
new_rev: Tid::from("rev1".to_string()),
new_rev: Tid::new("3k2aaaaaaaaab").unwrap(),
new_block_cids: vec![],
obsolete_block_cids: vec![],
record_upserts: vec![],
@@ -660,7 +684,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev1".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()),
},
};
@@ -681,7 +705,7 @@ mod tests {
did: test_did("nonexistent"),
expected_root_cid: None,
new_root_cid: test_cid_link(1),
new_rev: Tid::from("rev1".to_string()),
new_rev: Tid::new("3k2aaaaaaaaab").unwrap(),
new_block_cids: vec![],
obsolete_block_cids: vec![],
record_upserts: vec![],
@@ -715,15 +739,15 @@ mod tests {
let mid_root = test_cid_link(21);
let record_cid = test_cid_link(22);
let collection = Nsid::from("app.bsky.feed.post".to_string());
let rkey = Rkey::from("3k2del".to_string());
let collection = test_nsid("app.bsky.feed.post");
let rkey = test_rkey("3k2del");
let insert_input = ApplyCommitInput {
user_id,
did: did.clone(),
expected_root_cid: Some(root_cid.clone()),
new_root_cid: mid_root.clone(),
new_rev: Tid::from("rev1".to_string()),
new_rev: Tid::new("3k2aaaaaaaaab").unwrap(),
new_block_cids: vec![],
obsolete_block_cids: vec![],
record_upserts: vec![tranquil_db_traits::RecordUpsert {
@@ -743,7 +767,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev1".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()),
},
};
ops.apply_commit(insert_input).unwrap();
@@ -762,7 +786,7 @@ mod tests {
did: did.clone(),
expected_root_cid: Some(mid_root.clone()),
new_root_cid: final_root.clone(),
new_rev: Tid::from("rev2".to_string()),
new_rev: Tid::new("3k2aaaaaaaaac").unwrap(),
new_block_cids: vec![],
obsolete_block_cids: vec![],
record_upserts: vec![],
@@ -781,7 +805,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev2".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaac").unwrap()),
},
};
ops.apply_commit(delete_input).unwrap();
@@ -807,7 +831,7 @@ mod tests {
did: did.clone(),
expected_root_cid: Some(root_cid.clone()),
new_root_cid: new_root.clone(),
new_rev: Tid::from("rev1".to_string()),
new_rev: Tid::new("3k2aaaaaaaaab").unwrap(),
new_block_cids: vec![],
obsolete_block_cids: vec![],
record_upserts: vec![],
@@ -823,7 +847,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev1".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()),
},
};
@@ -833,7 +857,7 @@ mod tests {
let event = ops.event_ops.get_event_by_seq(seq).unwrap().unwrap();
assert_eq!(event.did, did);
assert_eq!(event.event_type, RepoEventType::Commit);
assert_eq!(event.rev.as_deref(), Some("rev1"));
assert_eq!(event.rev.as_deref(), Some("3k2aaaaaaaaab"));
}
#[test]
@@ -842,8 +866,8 @@ mod tests {
let ops = make_commit_ops(&h);
let (user_id, _did, root_cid) = create_test_repo(&h, "bailey", 40);
let collection = Nsid::from("app.bsky.feed.post".to_string());
let rkey = Rkey::from("3k2import".to_string());
let collection = test_nsid("app.bsky.feed.post");
let rkey = test_rkey("3k2import");
let record_cid = test_cid_link(41);
ops.import_repo_data(
@@ -915,7 +939,7 @@ mod tests {
did: did_a.clone(),
expected_root_cid: Some(root_a),
new_root_cid: new_root.clone(),
new_rev: Tid::from("rev1".to_string()),
new_rev: Tid::new("3k2aaaaaaaaab").unwrap(),
new_block_cids: vec![vec![0x01, 0x02, 0x03]],
obsolete_block_cids: vec![],
record_upserts: vec![],
@@ -931,7 +955,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev1".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()),
},
};
ops.apply_commit(input).unwrap();
@@ -953,7 +977,7 @@ mod tests {
did: did.clone(),
expected_root_cid: None,
new_root_cid: new_root.clone(),
new_rev: Tid::from("rev_force".to_string()),
new_rev: Tid::new("3k2aaaaaaaaaz").unwrap(),
new_block_cids: vec![],
obsolete_block_cids: vec![],
record_upserts: vec![],
@@ -969,7 +993,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev_force".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaaz").unwrap()),
},
};
@@ -983,10 +1007,10 @@ mod tests {
let h = setup();
let ops = make_commit_ops(&h);
let (user_id, did, root_cid) = create_test_repo(&h, "backlink_upd", 90);
let (user_id, did, root_cid) = create_test_repo(&h, "backlink-upd", 90);
let collection = Nsid::from("app.bsky.feed.like".to_string());
let rkey = Rkey::from("3k2like1".to_string());
let collection = test_nsid("app.bsky.feed.like");
let rkey = test_rkey("3k2like1");
let record_cid = test_cid_link(91);
let record_uri = AtUri::from_parts(&did, &collection, &rkey);
@@ -996,7 +1020,7 @@ mod tests {
did: did.clone(),
expected_root_cid: Some(root_cid.clone()),
new_root_cid: mid_root.clone(),
new_rev: Tid::from("rev1".to_string()),
new_rev: Tid::new("3k2aaaaaaaaab").unwrap(),
new_block_cids: vec![],
obsolete_block_cids: vec![],
record_upserts: vec![tranquil_db_traits::RecordUpsert {
@@ -1020,7 +1044,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev1".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()),
},
};
ops.apply_commit(create_input).unwrap();
@@ -1040,7 +1064,7 @@ mod tests {
did: did.clone(),
expected_root_cid: Some(mid_root.clone()),
new_root_cid: final_root.clone(),
new_rev: Tid::from("rev2".to_string()),
new_rev: Tid::new("3k2aaaaaaaaac").unwrap(),
new_block_cids: vec![],
obsolete_block_cids: vec![],
record_upserts: vec![tranquil_db_traits::RecordUpsert {
@@ -1064,7 +1088,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev2".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaac").unwrap()),
},
};
ops.apply_commit(update_input).unwrap();
@@ -1091,13 +1115,13 @@ mod tests {
std::fs::create_dir_all(&segments_dir).unwrap();
let user_id = Uuid::new_v4();
let did = test_did("crash_alice");
let handle = test_handle("crash_alice");
let did = test_did("crash-nel");
let handle = test_handle("crash-nel");
let initial_root = test_cid_link(200);
let new_root = test_cid_link(201);
let record_cid = test_cid_link(202);
let collection = Nsid::from("app.bsky.feed.post".to_string());
let rkey = Rkey::from("3k2crash".to_string());
let collection = test_nsid("app.bsky.feed.post");
let rkey = test_rkey("3k2crash");
let event_log = EventLog::open(
EventLogConfig {
@@ -1138,7 +1162,7 @@ mod tests {
did: did.clone(),
expected_root_cid: Some(initial_root.clone()),
new_root_cid: new_root.clone(),
new_rev: Tid::from("rev1".to_string()),
new_rev: Tid::new("3k2aaaaaaaaab").unwrap(),
new_block_cids: vec![vec![0xAA, 0xBB]],
obsolete_block_cids: vec![],
record_upserts: vec![tranquil_db_traits::RecordUpsert {
@@ -1158,7 +1182,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev1".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()),
},
};
@@ -1216,13 +1240,13 @@ mod tests {
std::fs::create_dir_all(&segments_dir).unwrap();
let user_id = Uuid::new_v4();
let did = test_did("crash_bob");
let handle = test_handle("crash_bob");
let did = test_did("crash-teq");
let handle = test_handle("crash-teq");
let initial_root = test_cid_link(210);
let new_root = test_cid_link(211);
let record_cid = test_cid_link(212);
let collection = Nsid::from("app.bsky.feed.post".to_string());
let rkey = Rkey::from("3k2bob".to_string());
let collection = test_nsid("app.bsky.feed.post");
let rkey = test_rkey("3k2bob");
let event_log = EventLog::open(
EventLogConfig {
@@ -1260,10 +1284,10 @@ mod tests {
let event_ops = metastore.event_ops(Arc::clone(&bridge));
let mutation_set = super::CommitMutationSet {
new_root_cid: super::cid_link_to_bytes(&new_root).unwrap(),
new_rev: test_rev(1),
new_rev: Tid::new("3k2abcdefghij").unwrap(),
record_upserts: vec![super::RecordMutationUpsert {
collection: collection.as_str().to_owned(),
rkey: rkey.as_str().to_owned(),
collection: collection.clone(),
rkey: rkey.clone(),
cid_bytes: super::cid_link_to_bytes(&record_cid).unwrap(),
}],
record_deletes: vec![],
@@ -1283,7 +1307,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev1".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()),
};
let mut batch = metastore.database().batch();
@@ -1325,7 +1349,7 @@ mod tests {
let repo_after = metastore.repo_ops().get_repo(user_id).unwrap().unwrap();
assert_eq!(repo_after.repo_root_cid, new_root);
assert_eq!(repo_after.repo_rev.as_deref(), Some("rev1"));
assert_eq!(repo_after.repo_rev.as_deref(), Some("3k2abcdefghij"));
let record_after = metastore
.record_ops()
@@ -1356,11 +1380,11 @@ mod tests {
let h = setup();
let ops = make_commit_ops(&h);
let (user_id, did, root_cid) = create_test_repo(&h, "col_iso", 95);
let (user_id, did, root_cid) = create_test_repo(&h, "col-iso", 95);
let col_like = Nsid::from("app.bsky.feed.like".to_string());
let col_repost = Nsid::from("app.bsky.feed.repost".to_string());
let rkey = Rkey::from("same_rkey".to_string());
let col_like = test_nsid("app.bsky.feed.like");
let col_repost = test_nsid("app.bsky.feed.repost");
let rkey = test_rkey("same_rkey");
let target = "at://did:plc:someone/app.bsky.feed.post/p1";
let mid_root = test_cid_link(96);
@@ -1372,7 +1396,7 @@ mod tests {
did: did.clone(),
expected_root_cid: Some(root_cid.clone()),
new_root_cid: mid_root.clone(),
new_rev: Tid::from("rev1".to_string()),
new_rev: Tid::new("3k2aaaaaaaaab").unwrap(),
new_block_cids: vec![],
obsolete_block_cids: vec![],
record_upserts: vec![
@@ -1410,7 +1434,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev1".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()),
},
};
ops.apply_commit(input).unwrap();
@@ -1441,7 +1465,7 @@ mod tests {
did: did.clone(),
expected_root_cid: Some(mid_root.clone()),
new_root_cid: final_root.clone(),
new_rev: Tid::from("rev2".to_string()),
new_rev: Tid::new("3k2aaaaaaaaac").unwrap(),
new_block_cids: vec![],
obsolete_block_cids: vec![],
record_upserts: vec![],
@@ -1457,7 +1481,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev2".to_string())),
rev: Some(Tid::new("3k2aaaaaaaaac").unwrap()),
},
};
ops.apply_commit(remove_like).unwrap();
+307 -40
View File
@@ -1,18 +1,19 @@
use std::collections::HashSet;
use serde::{Deserialize, Serialize};
use tranquil_types::{Nsid, Rkey, Tid};
use tranquil_types::{AtUri, Nsid, Rkey, Tid};
use super::backlink_ops::remove_backlinks_for_record;
use super::backlinks::{BacklinkValue, backlink_by_user_key, backlink_key, discriminant_to_path};
use super::encoding::KeyReader;
use super::keys::{KeyTag, UserHash};
use super::records::{RecordValue, record_key};
use super::records::{RecordValue, record_by_cid_key, record_key};
use super::repo_meta::{RepoMetaValue, repo_meta_key};
use super::user_blocks::{user_block_key, user_block_user_prefix};
use crate::metastore::MetastoreError;
const MUTATION_SET_VERSION: u8 = 1;
const MUTATION_SET_VERSION: u8 = 2;
const MUTATION_SET_VERSION_UNVALIDATED: u8 = 1;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CommitMutationSet {
@@ -23,29 +24,138 @@ pub struct CommitMutationSet {
pub block_inserts: Vec<Vec<u8>>,
pub block_deletes: Vec<Vec<u8>>,
pub backlink_adds: Vec<BacklinkMutation>,
pub backlink_remove_uris: Vec<String>,
pub backlink_remove_uris: Vec<AtUri>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RecordMutationUpsert {
pub collection: String,
pub rkey: String,
pub collection: Nsid,
pub rkey: Rkey,
pub cid_bytes: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RecordMutationDelete {
pub collection: String,
pub rkey: String,
pub collection: Nsid,
pub rkey: Rkey,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BacklinkMutation {
pub uri: String,
pub uri: AtUri,
pub path: u8,
pub link_to: String,
}
#[derive(Deserialize)]
struct UnvalidatedCommitMutationSet {
new_root_cid: Vec<u8>,
new_rev: String,
record_upserts: Vec<UnvalidatedRecordMutationUpsert>,
record_deletes: Vec<UnvalidatedRecordMutationDelete>,
block_inserts: Vec<Vec<u8>>,
block_deletes: Vec<Vec<u8>>,
backlink_adds: Vec<UnvalidatedBacklinkMutation>,
backlink_remove_uris: Vec<String>,
}
#[derive(Deserialize)]
struct UnvalidatedRecordMutationUpsert {
collection: String,
rkey: String,
cid_bytes: Vec<u8>,
}
#[derive(Deserialize)]
struct UnvalidatedRecordMutationDelete {
collection: String,
rkey: String,
}
#[derive(Deserialize)]
struct UnvalidatedBacklinkMutation {
uri: String,
path: u8,
link_to: String,
}
impl UnvalidatedCommitMutationSet {
fn into_validated(self) -> Option<CommitMutationSet> {
let warn = |field: &str, value: &str| {
tracing::warn!(
field,
value,
"version 1 CommitMutationSet has a value the current validators reject. \
Skipping the whole set rahter than replaying it in part."
);
};
let new_rev = Tid::new(self.new_rev.clone())
.inspect_err(|_| warn("new_rev", &self.new_rev))
.ok()?;
let record_upserts = self
.record_upserts
.into_iter()
.map(|u| {
Some(RecordMutationUpsert {
collection: Nsid::new(u.collection.clone())
.inspect_err(|_| warn("record_upserts.collection", &u.collection))
.ok()?,
rkey: Rkey::new(u.rkey.clone())
.inspect_err(|_| warn("record_upserts.rkey", &u.rkey))
.ok()?,
cid_bytes: u.cid_bytes,
})
})
.collect::<Option<Vec<_>>>()?;
let record_deletes = self
.record_deletes
.into_iter()
.map(|d| {
Some(RecordMutationDelete {
collection: Nsid::new(d.collection.clone())
.inspect_err(|_| warn("record_deletes.collection", &d.collection))
.ok()?,
rkey: Rkey::new(d.rkey.clone())
.inspect_err(|_| warn("record_deletes.rkey", &d.rkey))
.ok()?,
})
})
.collect::<Option<Vec<_>>>()?;
let backlink_adds = self
.backlink_adds
.into_iter()
.map(|b| {
Some(BacklinkMutation {
uri: AtUri::new(b.uri.clone())
.inspect_err(|_| warn("backlink_adds.uri", &b.uri))
.ok()?,
path: b.path,
link_to: b.link_to,
})
})
.collect::<Option<Vec<_>>>()?;
let backlink_remove_uris = self
.backlink_remove_uris
.into_iter()
.map(|uri| {
AtUri::new(uri.clone())
.inspect_err(|_| warn("backlink_remove_uris", &uri))
.ok()
})
.collect::<Option<Vec<_>>>()?;
Some(CommitMutationSet {
new_root_cid: self.new_root_cid,
new_rev,
record_upserts,
record_deletes,
block_inserts: self.block_inserts,
block_deletes: self.block_deletes,
backlink_adds,
backlink_remove_uris,
})
}
}
const MAX_MUTATION_SET_ENTRIES: usize = 50_000;
impl CommitMutationSet {
@@ -91,6 +201,15 @@ impl CommitMutationSet {
None
}
},
MUTATION_SET_VERSION_UNVALIDATED => {
match postcard::from_bytes::<UnvalidatedCommitMutationSet>(payload) {
Ok(v) => v.into_validated(),
Err(e) => {
tracing::warn!(%e, "failed to deserialize version 1 CommitMutationSet payload");
None
}
}
}
_ => {
tracing::warn!(version, "unknown CommitMutationSet version");
None
@@ -117,40 +236,64 @@ pub fn replay_mutation_set(
let meta_key = repo_meta_key(user_hash);
batch.insert(repo_data, meta_key.as_slice(), updated_meta.serialize());
mutation_set.record_upserts.iter().for_each(|u| {
let key = record_key(
user_hash,
&Nsid::from(u.collection.clone()),
&Rkey::from(u.rkey.clone()),
);
mutation_set.record_upserts.iter().try_for_each(|u| {
let key = record_key(user_hash, &u.collection, &u.rkey);
let previous = stored_record(repo_data, key.as_slice())?;
if let Some(prev) = previous
.as_ref()
.map(|p| &p.record_cid)
.filter(|prev| *prev != &u.cid_bytes)
{
let stale = record_by_cid_key(user_hash, prev, &u.collection, &u.rkey);
batch.remove(repo_data, stale.as_slice());
}
let value = RecordValue {
record_cid: u.cid_bytes.clone(),
takedown_ref: None,
takedown_ref: previous.and_then(|p| p.takedown_ref),
};
let reverse = record_by_cid_key(user_hash, &u.cid_bytes, &u.collection, &u.rkey);
batch.insert(repo_data, key.as_slice(), value.serialize());
});
batch.insert(repo_data, reverse.as_slice(), []);
Ok::<(), MetastoreError>(())
})?;
mutation_set.record_deletes.iter().for_each(|d| {
let key = record_key(
user_hash,
&Nsid::from(d.collection.clone()),
&Rkey::from(d.rkey.clone()),
);
mutation_set.record_deletes.iter().try_for_each(|d| {
let key = record_key(user_hash, &d.collection, &d.rkey);
if let Some(prev) = stored_record(repo_data, key.as_slice())? {
let reverse = record_by_cid_key(user_hash, &prev.record_cid, &d.collection, &d.rkey);
batch.remove(repo_data, reverse.as_slice());
}
batch.remove(repo_data, key.as_slice());
});
Ok::<(), MetastoreError>(())
})?;
mutation_set.block_inserts.iter().for_each(|cid_bytes| {
let key = user_block_key(user_hash, &mutation_set.new_rev, cid_bytes);
batch.insert(repo_data, key.as_slice(), []);
});
let already_recorded: HashSet<Vec<u8>> = match mutation_set.block_inserts.is_empty() {
true => HashSet::new(),
false => repo_data
.prefix(user_block_user_prefix(user_hash).as_slice())
.map(|guard| {
let (key_bytes, _) = guard.into_inner().map_err(MetastoreError::Fjall)?;
Ok(extract_cid_from_user_block_key(key_bytes.as_ref()).map(|c| c.to_vec()))
})
.filter_map(Result::transpose)
.collect::<Result<_, MetastoreError>>()?,
};
mutation_set
.block_inserts
.iter()
.filter(|cid_bytes| !cid_bytes.is_empty())
.filter(|cid_bytes| !already_recorded.contains(cid_bytes.as_slice()))
.for_each(|cid_bytes| {
let key = user_block_key(user_hash, &mutation_set.new_rev, cid_bytes);
batch.insert(repo_data, key.as_slice(), []);
});
delete_user_blocks_by_cid_scan(batch, repo_data, user_hash, &mutation_set.block_deletes)?;
mutation_set
.backlink_remove_uris
.iter()
.try_for_each(|uri_str| {
let uri = tranquil_types::AtUri::from(uri_str.clone());
.try_for_each(|uri| {
let collection = uri.collection().ok_or(MetastoreError::CorruptData(
"backlink URI missing collection",
))?;
@@ -162,11 +305,11 @@ pub fn replay_mutation_set(
})?;
mutation_set.backlink_adds.iter().try_for_each(|bl| {
let uri = tranquil_types::AtUri::from(bl.uri.clone());
let collection = uri.collection().ok_or(MetastoreError::CorruptData(
let collection = bl.uri.collection().ok_or(MetastoreError::CorruptData(
"backlink URI missing collection",
))?;
let rkey = uri
let rkey = bl
.uri
.rkey()
.ok_or(MetastoreError::CorruptData("backlink URI missing rkey"))?;
@@ -181,7 +324,7 @@ pub fn replay_mutation_set(
Some(_) => {
let primary = backlink_key(&bl.link_to, user_hash, collection, rkey);
let value = BacklinkValue {
source_uri: bl.uri.clone(),
source_uri: bl.uri.as_str().to_owned(),
path: bl.path,
};
batch.insert(indexes, primary.as_slice(), value.serialize());
@@ -194,6 +337,16 @@ pub fn replay_mutation_set(
})
}
fn stored_record(
repo_data: &fjall::Keyspace,
key: &[u8],
) -> Result<Option<RecordValue>, MetastoreError> {
Ok(repo_data
.get(key)
.map_err(MetastoreError::Fjall)?
.and_then(|raw| RecordValue::deserialize(&raw)))
}
fn delete_user_blocks_by_cid_scan(
batch: &mut fjall::OwnedWriteBatch,
repo_data: &fjall::Keyspace,
@@ -256,22 +409,24 @@ mod tests {
new_root_cid: vec![0x01, 0x71, 0x12, 0x20],
new_rev: Tid::new("3k2abcdefghij").unwrap(),
record_upserts: vec![RecordMutationUpsert {
collection: "app.bsky.feed.post".to_owned(),
rkey: "3k2abc".to_owned(),
collection: Nsid::new("app.bsky.feed.post").unwrap(),
rkey: Rkey::new("3k2abc").unwrap(),
cid_bytes: vec![0xDE, 0xAD],
}],
record_deletes: vec![RecordMutationDelete {
collection: "app.bsky.feed.like".to_owned(),
rkey: "3k2del".to_owned(),
collection: Nsid::new("app.bsky.feed.like").unwrap(),
rkey: Rkey::new("3k2del").unwrap(),
}],
block_inserts: vec![vec![0x01, 0x02]],
block_deletes: vec![vec![0x03, 0x04]],
backlink_adds: vec![BacklinkMutation {
uri: "at://did:plc:olaren/app.bsky.feed.like/3k2abc".to_owned(),
uri: AtUri::new("at://did:plc:olaren/app.bsky.feed.like/3k2abc").unwrap(),
path: 1,
link_to: "at://did:plc:teq/app.bsky.feed.post/3k2xyz".to_owned(),
}],
backlink_remove_uris: vec!["at://did:plc:olaren/app.bsky.feed.like/3k2old".to_owned()],
backlink_remove_uris: vec![
AtUri::new("at://did:plc:olaren/app.bsky.feed.like/3k2old").unwrap(),
],
};
let bytes = ms.serialize().unwrap();
@@ -280,6 +435,36 @@ mod tests {
assert_eq!(recovered, ms);
}
#[test]
fn mutation_set_rejects_a_field_corrupted_into_a_structurally_valid_decode() {
let ms = CommitMutationSet {
new_root_cid: vec![0x01],
new_rev: Tid::new("3k2abcdefghij").unwrap(),
record_upserts: vec![RecordMutationUpsert {
collection: Nsid::new("app.bsky.feed.post").unwrap(),
rkey: Rkey::new("3k2abc").unwrap(),
cid_bytes: vec![0x02],
}],
record_deletes: vec![],
block_inserts: vec![],
block_deletes: vec![],
backlink_adds: vec![],
backlink_remove_uris: vec![],
};
let mut bytes = ms.serialize().unwrap();
let nsid_start = bytes
.windows(b"app.bsky.feed.post".len())
.position(|w| w == b"app.bsky.feed.post")
.expect("serialized form contains the collection");
bytes[nsid_start] = b'!';
assert!(
CommitMutationSet::deserialize(&bytes).is_none(),
"a corrupt field that still decodes as a string must be rejected"
);
}
#[test]
fn mutation_set_empty_roundtrip() {
let ms = CommitMutationSet {
@@ -313,4 +498,86 @@ mod tests {
bytes[0] = 99;
assert!(CommitMutationSet::deserialize(&bytes).is_none());
}
fn v1_payload(rev: &str, collection: &str, rkey: &str) -> Vec<u8> {
#[derive(Serialize)]
struct V1 {
new_root_cid: Vec<u8>,
new_rev: String,
record_upserts: Vec<V1Upsert>,
record_deletes: Vec<(String, String)>,
block_inserts: Vec<Vec<u8>>,
block_deletes: Vec<Vec<u8>>,
backlink_adds: Vec<(String, u8, String)>,
backlink_remove_uris: Vec<String>,
}
#[derive(Serialize)]
struct V1Upsert {
collection: String,
rkey: String,
cid_bytes: Vec<u8>,
}
let payload = postcard::to_allocvec(&V1 {
new_root_cid: vec![0x01, 0x71],
new_rev: rev.to_owned(),
record_upserts: vec![V1Upsert {
collection: collection.to_owned(),
rkey: rkey.to_owned(),
cid_bytes: vec![0xDE, 0xAD],
}],
record_deletes: vec![],
block_inserts: vec![vec![0x01, 0x02]],
block_deletes: vec![],
backlink_adds: vec![],
backlink_remove_uris: vec![],
})
.unwrap();
std::iter::once(MUTATION_SET_VERSION_UNVALIDATED)
.chain(payload)
.collect()
}
#[test]
fn a_version_1_payload_written_by_an_older_binary_still_replays() {
let decoded = CommitMutationSet::deserialize(&v1_payload(
"3k2abcdefghij",
"app.bsky.feed.post",
"3k2abc",
))
.expect("a version 1 payload with valid values decodes");
assert_eq!(decoded.new_rev.as_str(), "3k2abcdefghij");
assert_eq!(decoded.record_upserts[0].rkey.as_str(), "3k2abc");
assert_eq!(decoded.block_inserts, vec![vec![0x01, 0x02]]);
}
#[test]
fn a_version_1_payload_the_current_validators_reject_is_skipped_not_fatal() {
assert!(
CommitMutationSet::deserialize(&v1_payload("3k2abcdefghij", "not/an/nsid", "3k2abc"))
.is_none(),
"a collection the validators reject skips the whole set instead of replaying it in part"
);
assert!(
CommitMutationSet::deserialize(&v1_payload("0", "app.bsky.feed.post", "3k2abc"))
.is_none(),
"a rev that isn't a TID skips the set"
);
}
#[test]
fn a_current_payload_is_written_at_the_validated_version() {
let ms = CommitMutationSet {
new_root_cid: vec![0x01],
new_rev: Tid::new("3k2abcdefghij").unwrap(),
record_upserts: vec![],
record_deletes: vec![],
block_inserts: vec![],
block_deletes: vec![],
backlink_adds: vec![],
backlink_remove_uris: vec![],
};
assert_eq!(ms.serialize().unwrap()[0], MUTATION_SET_VERSION);
}
}
+26 -15
View File
@@ -10,7 +10,7 @@ use tranquil_store::metastore::recovery::{
};
use tranquil_store::metastore::{Metastore, MetastoreConfig};
use tranquil_store::{sim_proptest_cases, sim_seed_range};
use tranquil_types::{CidLink, Did, Handle, Tid};
use tranquil_types::{AtUri, CidLink, Did, Handle, Nsid, Rkey, Tid};
use uuid::Uuid;
const NAMES: &[&str] = &["olaren", "teq", "nel", "lyna", "bailey"];
@@ -27,7 +27,7 @@ fn open_metastore(path: &Path) -> Metastore {
fn test_did(seed: u64) -> Did {
let name = NAMES[(seed as usize) % NAMES.len()];
Did::from(format!("did:plc:{name}{seed}"))
Did::new(format!("did:plc:{name}{seed}")).expect("generated DID is valid")
}
fn test_handle(seed: u64) -> Handle {
@@ -46,29 +46,37 @@ fn test_uuid(seed: u64) -> Uuid {
Uuid::from_u128(seed as u128 | 0x4000_0000_0000_0000_8000_0000_0000_0000)
}
const NSID_PATTERN: &str = "[a-z]{3,8}\\.[a-z]{3,8}\\.[a-z]{3,8}";
const RKEY_PATTERN: &str = "[a-z0-9]{3,10}";
const TID_PATTERN: &str = "[234567abcdefghij][234567abcdefghijklmnopqrstuvwxyz]{12}";
const AT_URI_PATTERN: &str =
"at://did:plc:[a-z]{3,8}/[a-z]{3,8}\\.[a-z]{3,8}\\.[a-z]{3,8}/[a-z0-9]{3,8}";
fn arb_mutation_set() -> impl Strategy<Value = CommitMutationSet> {
let arb_upsert = (
"[a-z\\.]{5,20}",
"[a-z0-9]{3,10}",
NSID_PATTERN,
RKEY_PATTERN,
prop::collection::vec(any::<u8>(), 4..36),
)
.prop_map(|(collection, rkey, cid_bytes)| RecordMutationUpsert {
collection,
rkey,
collection: Nsid::new(collection).expect("generated NSID is valid"),
rkey: Rkey::new(rkey).expect("generated rkey is valid"),
cid_bytes,
});
let arb_delete = ("[a-z\\.]{5,20}", "[a-z0-9]{3,10}")
.prop_map(|(collection, rkey)| RecordMutationDelete { collection, rkey });
let arb_delete =
(NSID_PATTERN, RKEY_PATTERN).prop_map(|(collection, rkey)| RecordMutationDelete {
collection: Nsid::new(collection).expect("generated NSID is valid"),
rkey: Rkey::new(rkey).expect("generated rkey is valid"),
});
let arb_backlink = (
"at://did:plc:[a-z]{3,8}/[a-z\\.]{5,20}/[a-z0-9]{3,8}",
0u8..4,
"at://did:plc:[a-z]{3,8}/[a-z\\.]{5,20}/[a-z0-9]{3,8}",
)
.prop_map(|(uri, path, link_to)| BacklinkMutation { uri, path, link_to });
let arb_backlink = (AT_URI_PATTERN, 0u8..4, AT_URI_PATTERN).prop_map(|(uri, path, link_to)| {
BacklinkMutation {
uri: AtUri::new(uri).expect("generated AT URI is valid"),
path,
link_to,
}
});
(
prop::collection::vec(any::<u8>(), 0..64),
@@ -78,7 +86,10 @@ fn arb_mutation_set() -> impl Strategy<Value = CommitMutationSet> {
prop::collection::vec(prop::collection::vec(any::<u8>(), 4..36), 0..20),
prop::collection::vec(prop::collection::vec(any::<u8>(), 4..36), 0..20),
prop::collection::vec(arb_backlink, 0..5),
prop::collection::vec("at://did:plc:[a-z]{3,8}/[a-z\\.]{5,20}/[a-z0-9]{3,8}", 0..5),
prop::collection::vec(
AT_URI_PATTERN.prop_map(|uri| AtUri::new(uri).expect("generated AT URI is valid")),
0..5,
),
)
.prop_map(
|(