diff --git a/crates/tranquil-store/src/metastore/commit_ops.rs b/crates/tranquil-store/src/metastore/commit_ops.rs index 4392a53..0e9003a 100644 --- a/crates/tranquil-store/src/metastore/commit_ops.rs +++ b/crates/tranquil-store/src/metastore/commit_ops.rs @@ -212,8 +212,8 @@ impl CommitOps { 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 CommitOps { .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 CommitOps { .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 CommitOps { 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 CommitOps { 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 CommitOps { 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 CommitOps { 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) -> Nsid { + Nsid::new(name).expect("test collection is a valid NSID") + } + + fn test_rkey(key: impl Into) -> 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(); diff --git a/crates/tranquil-store/src/metastore/recovery.rs b/crates/tranquil-store/src/metastore/recovery.rs index 66c82b3..faa3905 100644 --- a/crates/tranquil-store/src/metastore/recovery.rs +++ b/crates/tranquil-store/src/metastore/recovery.rs @@ -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>, pub block_deletes: Vec>, pub backlink_adds: Vec, - pub backlink_remove_uris: Vec, + pub backlink_remove_uris: Vec, } #[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, } #[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, + new_rev: String, + record_upserts: Vec, + record_deletes: Vec, + block_inserts: Vec>, + block_deletes: Vec>, + backlink_adds: Vec, + backlink_remove_uris: Vec, +} + +#[derive(Deserialize)] +struct UnvalidatedRecordMutationUpsert { + collection: String, + rkey: String, + cid_bytes: Vec, +} + +#[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 { + 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::>>()?; + 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::>>()?; + 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::>>()?; + 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::>>()?; + 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::(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> = 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::>()?, + }; + 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, 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 { + #[derive(Serialize)] + struct V1 { + new_root_cid: Vec, + new_rev: String, + record_upserts: Vec, + record_deletes: Vec<(String, String)>, + block_inserts: Vec>, + block_deletes: Vec>, + backlink_adds: Vec<(String, u8, String)>, + backlink_remove_uris: Vec, + } + #[derive(Serialize)] + struct V1Upsert { + collection: String, + rkey: String, + cid_bytes: Vec, + } + + 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); + } } diff --git a/crates/tranquil-store/tests/metastore_crash.rs b/crates/tranquil-store/tests/metastore_crash.rs index b77ebc1..66a62b7 100644 --- a/crates/tranquil-store/tests/metastore_crash.rs +++ b/crates/tranquil-store/tests/metastore_crash.rs @@ -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 { let arb_upsert = ( - "[a-z\\.]{5,20}", - "[a-z0-9]{3,10}", + NSID_PATTERN, + RKEY_PATTERN, prop::collection::vec(any::(), 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::(), 0..64), @@ -78,7 +86,10 @@ fn arb_mutation_set() -> impl Strategy { prop::collection::vec(prop::collection::vec(any::(), 4..36), 0..20), prop::collection::vec(prop::collection::vec(any::(), 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( |(