Our Personal Data Server from scratch!
0

Configure Feed

Select the types of activity you want to include in your feed.

store: revalidate stored mutation sets on replay

Lewis: May this revision serve well! <lu5a@proton.me>

author
Lewis
date (Jul 24, 2026, 8:27 PM +0300) commit 8dd06ab3 parent cd4e606d change-id kotpknor
+435 -133
+102 -78
crates/tranquil-store/src/metastore/commit_ops.rs
··· 212 212 let cid_bytes = cid_link_to_bytes(&u.cid) 213 213 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 214 214 Ok(RecordMutationUpsert { 215 - collection: u.collection.as_str().to_owned(), 216 - rkey: u.rkey.as_str().to_owned(), 215 + collection: u.collection.clone(), 216 + rkey: u.rkey.clone(), 217 217 cid_bytes, 218 218 }) 219 219 }) ··· 222 222 .record_deletes 223 223 .iter() 224 224 .map(|d| RecordMutationDelete { 225 - collection: d.collection.as_str().to_owned(), 226 - rkey: d.rkey.as_str().to_owned(), 225 + collection: d.collection.clone(), 226 + rkey: d.rkey.clone(), 227 227 }) 228 228 .collect(), 229 229 block_inserts: input.new_block_cids.clone(), ··· 232 232 .backlinks_to_add 233 233 .iter() 234 234 .map(|bl| BacklinkMutation { 235 - uri: bl.uri.as_str().to_owned(), 235 + uri: bl.uri.clone(), 236 236 path: path_to_discriminant(bl.path), 237 237 link_to: bl.link_to.clone(), 238 238 }) 239 239 .collect(), 240 - backlink_remove_uris: input 241 - .backlinks_to_remove 242 - .iter() 243 - .map(|uri| uri.as_str().to_owned()) 244 - .collect(), 240 + backlink_remove_uris: input.backlinks_to_remove.clone(), 245 241 }; 246 242 let mutation_set_bytes = mutation_set 247 243 .serialize() ··· 372 368 self.scan_users_missing_prefix( 373 369 record_blobs_user_prefix, 374 370 |meta, user_id| { 375 - let did = meta 376 - .did 377 - .map(Did::from) 378 - .ok_or(MetastoreError::CorruptData("repo_meta missing did field"))?; 371 + let did = match meta.did { 372 + None => Err(MetastoreError::CorruptData("repo_meta missing DID field")), 373 + Some(d) => Did::new(d) 374 + .map_err(|_| MetastoreError::CorruptData("corrupt repo_meta did")), 375 + }?; 379 376 Ok(UserNeedingRecordBlobsBackfill { user_id, did }) 380 377 }, 381 378 limit_usize, ··· 394 391 repo_root_cid: root_cid, 395 392 repo_rev: match meta.repo_rev.is_empty() { 396 393 true => None, 397 - false => Some(Tid::from(meta.repo_rev)), 394 + false => Some(Tid::new(meta.repo_rev).map_err(|_| { 395 + MetastoreError::CorruptData("corrupt repo_meta repo_rev") 396 + })?), 398 397 }, 399 398 }) 400 399 }, ··· 424 423 425 424 let user_hash = match parse_user_hash_from_key(&key_bytes) { 426 425 Some(h) => h, 427 - None => return Some(Err(MetastoreError::CorruptData("invalid repo_meta key"))), 426 + None => { 427 + tracing::warn!("skipping a repo_meta row whose key doesn't parse"); 428 + return None; 429 + } 428 430 }; 429 431 430 432 let check_prefix = make_prefix(user_hash); ··· 442 444 let meta = match RepoMetaValue::deserialize(&val_bytes) { 443 445 Some(v) => v, 444 446 None => { 445 - return Some(Err(MetastoreError::CorruptData( 446 - "invalid repo_meta value", 447 - ))); 447 + tracing::warn!( 448 + user_hash = user_hash.raw(), 449 + "skipping a repo_meta row whose value doesn't decode" 450 + ); 451 + return None; 448 452 } 449 453 }; 450 454 let user_id = match self.user_hashes.get_uuid(&user_hash) { 451 455 Some(id) => id, 452 456 None => { 453 - return Some(Err(MetastoreError::CorruptData( 454 - "user_hash has no reverse mapping", 455 - ))); 457 + tracing::warn!( 458 + user_hash = user_hash.raw(), 459 + "skipping a repo_meta row whose user_hash has no reverse mapping" 460 + ); 461 + return None; 456 462 } 457 463 }; 458 - Some(build_result(meta, user_id)) 464 + match build_result(meta, user_id) { 465 + Ok(v) => Some(Ok(v)), 466 + Err(e) => { 467 + tracing::warn!( 468 + user_id = %user_id, 469 + error = %e, 470 + "skipping a repo_meta row the current validators reject" 471 + ); 472 + None 473 + } 474 + } 459 475 } 460 476 } 461 477 }) ··· 528 544 } 529 545 530 546 fn test_did(name: &str) -> Did { 531 - Did::from(format!("did:plc:{name}")) 547 + Did::new(format!("did:plc:{name}")).expect("test DID is well-formed") 532 548 } 533 549 534 550 fn test_handle(name: &str) -> Handle { 535 - Handle::from(format!("{name}.test.invalid")) 551 + Handle::new(format!("{name}.oyster.cafe")).expect("test handle is valid") 552 + } 553 + 554 + fn test_nsid(name: impl Into<String>) -> Nsid { 555 + Nsid::new(name).expect("test collection is a valid NSID") 556 + } 557 + 558 + fn test_rkey(key: impl Into<String>) -> Rkey { 559 + Rkey::new(key).expect("test record key is a valid rkey") 536 560 } 537 561 538 562 fn test_rev(seq: u64) -> Tid { ··· 582 606 583 607 let new_root = test_cid_link(2); 584 608 let record_cid = test_cid_link(3); 585 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 586 - let rkey = Rkey::from("3k2abc".to_string()); 609 + let collection = test_nsid("app.bsky.feed.post"); 610 + let rkey = test_rkey("3k2abc"); 587 611 588 612 let input = ApplyCommitInput { 589 613 user_id, 590 614 did: did.clone(), 591 615 expected_root_cid: Some(root_cid.clone()), 592 616 new_root_cid: new_root.clone(), 593 - new_rev: Tid::from("rev1".to_string()), 617 + new_rev: Tid::new("3k2aaaaaaaaab").unwrap(), 594 618 new_block_cids: vec![vec![0x01, 0x02]], 595 619 obsolete_block_cids: vec![], 596 620 record_upserts: vec![tranquil_db_traits::RecordUpsert { ··· 610 634 blobs: None, 611 635 blocks: None, 612 636 prev_data_cid: None, 613 - rev: Some(Tid::from("rev1".to_string())), 637 + rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()), 614 638 }, 615 639 }; 616 640 ··· 619 643 620 644 let repo = h.metastore.repo_ops().get_repo(user_id).unwrap().unwrap(); 621 645 assert_eq!(repo.repo_root_cid, new_root); 622 - assert_eq!(repo.repo_rev.as_deref(), Some("rev1")); 646 + assert_eq!(repo.repo_rev.as_deref(), Some("3k2aaaaaaaaab")); 623 647 624 648 let found_cid = h 625 649 .metastore ··· 644 668 did, 645 669 expected_root_cid: Some(stale_root), 646 670 new_root_cid: new_root, 647 - new_rev: Tid::from("rev1".to_string()), 671 + new_rev: Tid::new("3k2aaaaaaaaab").unwrap(), 648 672 new_block_cids: vec![], 649 673 obsolete_block_cids: vec![], 650 674 record_upserts: vec![], ··· 660 684 blobs: None, 661 685 blocks: None, 662 686 prev_data_cid: None, 663 - rev: Some(Tid::from("rev1".to_string())), 687 + rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()), 664 688 }, 665 689 }; 666 690 ··· 681 705 did: test_did("nonexistent"), 682 706 expected_root_cid: None, 683 707 new_root_cid: test_cid_link(1), 684 - new_rev: Tid::from("rev1".to_string()), 708 + new_rev: Tid::new("3k2aaaaaaaaab").unwrap(), 685 709 new_block_cids: vec![], 686 710 obsolete_block_cids: vec![], 687 711 record_upserts: vec![], ··· 715 739 716 740 let mid_root = test_cid_link(21); 717 741 let record_cid = test_cid_link(22); 718 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 719 - let rkey = Rkey::from("3k2del".to_string()); 742 + let collection = test_nsid("app.bsky.feed.post"); 743 + let rkey = test_rkey("3k2del"); 720 744 721 745 let insert_input = ApplyCommitInput { 722 746 user_id, 723 747 did: did.clone(), 724 748 expected_root_cid: Some(root_cid.clone()), 725 749 new_root_cid: mid_root.clone(), 726 - new_rev: Tid::from("rev1".to_string()), 750 + new_rev: Tid::new("3k2aaaaaaaaab").unwrap(), 727 751 new_block_cids: vec![], 728 752 obsolete_block_cids: vec![], 729 753 record_upserts: vec![tranquil_db_traits::RecordUpsert { ··· 743 767 blobs: None, 744 768 blocks: None, 745 769 prev_data_cid: None, 746 - rev: Some(Tid::from("rev1".to_string())), 770 + rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()), 747 771 }, 748 772 }; 749 773 ops.apply_commit(insert_input).unwrap(); ··· 762 786 did: did.clone(), 763 787 expected_root_cid: Some(mid_root.clone()), 764 788 new_root_cid: final_root.clone(), 765 - new_rev: Tid::from("rev2".to_string()), 789 + new_rev: Tid::new("3k2aaaaaaaaac").unwrap(), 766 790 new_block_cids: vec![], 767 791 obsolete_block_cids: vec![], 768 792 record_upserts: vec![], ··· 781 805 blobs: None, 782 806 blocks: None, 783 807 prev_data_cid: None, 784 - rev: Some(Tid::from("rev2".to_string())), 808 + rev: Some(Tid::new("3k2aaaaaaaaac").unwrap()), 785 809 }, 786 810 }; 787 811 ops.apply_commit(delete_input).unwrap(); ··· 807 831 did: did.clone(), 808 832 expected_root_cid: Some(root_cid.clone()), 809 833 new_root_cid: new_root.clone(), 810 - new_rev: Tid::from("rev1".to_string()), 834 + new_rev: Tid::new("3k2aaaaaaaaab").unwrap(), 811 835 new_block_cids: vec![], 812 836 obsolete_block_cids: vec![], 813 837 record_upserts: vec![], ··· 823 847 blobs: None, 824 848 blocks: None, 825 849 prev_data_cid: None, 826 - rev: Some(Tid::from("rev1".to_string())), 850 + rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()), 827 851 }, 828 852 }; 829 853 ··· 833 857 let event = ops.event_ops.get_event_by_seq(seq).unwrap().unwrap(); 834 858 assert_eq!(event.did, did); 835 859 assert_eq!(event.event_type, RepoEventType::Commit); 836 - assert_eq!(event.rev.as_deref(), Some("rev1")); 860 + assert_eq!(event.rev.as_deref(), Some("3k2aaaaaaaaab")); 837 861 } 838 862 839 863 #[test] ··· 842 866 let ops = make_commit_ops(&h); 843 867 let (user_id, _did, root_cid) = create_test_repo(&h, "bailey", 40); 844 868 845 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 846 - let rkey = Rkey::from("3k2import".to_string()); 869 + let collection = test_nsid("app.bsky.feed.post"); 870 + let rkey = test_rkey("3k2import"); 847 871 let record_cid = test_cid_link(41); 848 872 849 873 ops.import_repo_data( ··· 915 939 did: did_a.clone(), 916 940 expected_root_cid: Some(root_a), 917 941 new_root_cid: new_root.clone(), 918 - new_rev: Tid::from("rev1".to_string()), 942 + new_rev: Tid::new("3k2aaaaaaaaab").unwrap(), 919 943 new_block_cids: vec![vec![0x01, 0x02, 0x03]], 920 944 obsolete_block_cids: vec![], 921 945 record_upserts: vec![], ··· 931 955 blobs: None, 932 956 blocks: None, 933 957 prev_data_cid: None, 934 - rev: Some(Tid::from("rev1".to_string())), 958 + rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()), 935 959 }, 936 960 }; 937 961 ops.apply_commit(input).unwrap(); ··· 953 977 did: did.clone(), 954 978 expected_root_cid: None, 955 979 new_root_cid: new_root.clone(), 956 - new_rev: Tid::from("rev_force".to_string()), 980 + new_rev: Tid::new("3k2aaaaaaaaaz").unwrap(), 957 981 new_block_cids: vec![], 958 982 obsolete_block_cids: vec![], 959 983 record_upserts: vec![], ··· 969 993 blobs: None, 970 994 blocks: None, 971 995 prev_data_cid: None, 972 - rev: Some(Tid::from("rev_force".to_string())), 996 + rev: Some(Tid::new("3k2aaaaaaaaaz").unwrap()), 973 997 }, 974 998 }; 975 999 ··· 983 1007 984 1008 let h = setup(); 985 1009 let ops = make_commit_ops(&h); 986 - let (user_id, did, root_cid) = create_test_repo(&h, "backlink_upd", 90); 1010 + let (user_id, did, root_cid) = create_test_repo(&h, "backlink-upd", 90); 987 1011 988 - let collection = Nsid::from("app.bsky.feed.like".to_string()); 989 - let rkey = Rkey::from("3k2like1".to_string()); 1012 + let collection = test_nsid("app.bsky.feed.like"); 1013 + let rkey = test_rkey("3k2like1"); 990 1014 let record_cid = test_cid_link(91); 991 1015 let record_uri = AtUri::from_parts(&did, &collection, &rkey); 992 1016 ··· 996 1020 did: did.clone(), 997 1021 expected_root_cid: Some(root_cid.clone()), 998 1022 new_root_cid: mid_root.clone(), 999 - new_rev: Tid::from("rev1".to_string()), 1023 + new_rev: Tid::new("3k2aaaaaaaaab").unwrap(), 1000 1024 new_block_cids: vec![], 1001 1025 obsolete_block_cids: vec![], 1002 1026 record_upserts: vec![tranquil_db_traits::RecordUpsert { ··· 1020 1044 blobs: None, 1021 1045 blocks: None, 1022 1046 prev_data_cid: None, 1023 - rev: Some(Tid::from("rev1".to_string())), 1047 + rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()), 1024 1048 }, 1025 1049 }; 1026 1050 ops.apply_commit(create_input).unwrap(); ··· 1040 1064 did: did.clone(), 1041 1065 expected_root_cid: Some(mid_root.clone()), 1042 1066 new_root_cid: final_root.clone(), 1043 - new_rev: Tid::from("rev2".to_string()), 1067 + new_rev: Tid::new("3k2aaaaaaaaac").unwrap(), 1044 1068 new_block_cids: vec![], 1045 1069 obsolete_block_cids: vec![], 1046 1070 record_upserts: vec![tranquil_db_traits::RecordUpsert { ··· 1064 1088 blobs: None, 1065 1089 blocks: None, 1066 1090 prev_data_cid: None, 1067 - rev: Some(Tid::from("rev2".to_string())), 1091 + rev: Some(Tid::new("3k2aaaaaaaaac").unwrap()), 1068 1092 }, 1069 1093 }; 1070 1094 ops.apply_commit(update_input).unwrap(); ··· 1091 1115 std::fs::create_dir_all(&segments_dir).unwrap(); 1092 1116 1093 1117 let user_id = Uuid::new_v4(); 1094 - let did = test_did("crash_alice"); 1095 - let handle = test_handle("crash_alice"); 1118 + let did = test_did("crash-nel"); 1119 + let handle = test_handle("crash-nel"); 1096 1120 let initial_root = test_cid_link(200); 1097 1121 let new_root = test_cid_link(201); 1098 1122 let record_cid = test_cid_link(202); 1099 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1100 - let rkey = Rkey::from("3k2crash".to_string()); 1123 + let collection = test_nsid("app.bsky.feed.post"); 1124 + let rkey = test_rkey("3k2crash"); 1101 1125 1102 1126 let event_log = EventLog::open( 1103 1127 EventLogConfig { ··· 1138 1162 did: did.clone(), 1139 1163 expected_root_cid: Some(initial_root.clone()), 1140 1164 new_root_cid: new_root.clone(), 1141 - new_rev: Tid::from("rev1".to_string()), 1165 + new_rev: Tid::new("3k2aaaaaaaaab").unwrap(), 1142 1166 new_block_cids: vec![vec![0xAA, 0xBB]], 1143 1167 obsolete_block_cids: vec![], 1144 1168 record_upserts: vec![tranquil_db_traits::RecordUpsert { ··· 1158 1182 blobs: None, 1159 1183 blocks: None, 1160 1184 prev_data_cid: None, 1161 - rev: Some(Tid::from("rev1".to_string())), 1185 + rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()), 1162 1186 }, 1163 1187 }; 1164 1188 ··· 1216 1240 std::fs::create_dir_all(&segments_dir).unwrap(); 1217 1241 1218 1242 let user_id = Uuid::new_v4(); 1219 - let did = test_did("crash_bob"); 1220 - let handle = test_handle("crash_bob"); 1243 + let did = test_did("crash-teq"); 1244 + let handle = test_handle("crash-teq"); 1221 1245 let initial_root = test_cid_link(210); 1222 1246 let new_root = test_cid_link(211); 1223 1247 let record_cid = test_cid_link(212); 1224 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1225 - let rkey = Rkey::from("3k2bob".to_string()); 1248 + let collection = test_nsid("app.bsky.feed.post"); 1249 + let rkey = test_rkey("3k2bob"); 1226 1250 1227 1251 let event_log = EventLog::open( 1228 1252 EventLogConfig { ··· 1260 1284 let event_ops = metastore.event_ops(Arc::clone(&bridge)); 1261 1285 let mutation_set = super::CommitMutationSet { 1262 1286 new_root_cid: super::cid_link_to_bytes(&new_root).unwrap(), 1263 - new_rev: test_rev(1), 1287 + new_rev: Tid::new("3k2abcdefghij").unwrap(), 1264 1288 record_upserts: vec![super::RecordMutationUpsert { 1265 - collection: collection.as_str().to_owned(), 1266 - rkey: rkey.as_str().to_owned(), 1289 + collection: collection.clone(), 1290 + rkey: rkey.clone(), 1267 1291 cid_bytes: super::cid_link_to_bytes(&record_cid).unwrap(), 1268 1292 }], 1269 1293 record_deletes: vec![], ··· 1283 1307 blobs: None, 1284 1308 blocks: None, 1285 1309 prev_data_cid: None, 1286 - rev: Some(Tid::from("rev1".to_string())), 1310 + rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()), 1287 1311 }; 1288 1312 1289 1313 let mut batch = metastore.database().batch(); ··· 1325 1349 1326 1350 let repo_after = metastore.repo_ops().get_repo(user_id).unwrap().unwrap(); 1327 1351 assert_eq!(repo_after.repo_root_cid, new_root); 1328 - assert_eq!(repo_after.repo_rev.as_deref(), Some("rev1")); 1352 + assert_eq!(repo_after.repo_rev.as_deref(), Some("3k2abcdefghij")); 1329 1353 1330 1354 let record_after = metastore 1331 1355 .record_ops() ··· 1356 1380 1357 1381 let h = setup(); 1358 1382 let ops = make_commit_ops(&h); 1359 - let (user_id, did, root_cid) = create_test_repo(&h, "col_iso", 95); 1383 + let (user_id, did, root_cid) = create_test_repo(&h, "col-iso", 95); 1360 1384 1361 - let col_like = Nsid::from("app.bsky.feed.like".to_string()); 1362 - let col_repost = Nsid::from("app.bsky.feed.repost".to_string()); 1363 - let rkey = Rkey::from("same_rkey".to_string()); 1385 + let col_like = test_nsid("app.bsky.feed.like"); 1386 + let col_repost = test_nsid("app.bsky.feed.repost"); 1387 + let rkey = test_rkey("same_rkey"); 1364 1388 let target = "at://did:plc:someone/app.bsky.feed.post/p1"; 1365 1389 1366 1390 let mid_root = test_cid_link(96); ··· 1372 1396 did: did.clone(), 1373 1397 expected_root_cid: Some(root_cid.clone()), 1374 1398 new_root_cid: mid_root.clone(), 1375 - new_rev: Tid::from("rev1".to_string()), 1399 + new_rev: Tid::new("3k2aaaaaaaaab").unwrap(), 1376 1400 new_block_cids: vec![], 1377 1401 obsolete_block_cids: vec![], 1378 1402 record_upserts: vec![ ··· 1410 1434 blobs: None, 1411 1435 blocks: None, 1412 1436 prev_data_cid: None, 1413 - rev: Some(Tid::from("rev1".to_string())), 1437 + rev: Some(Tid::new("3k2aaaaaaaaab").unwrap()), 1414 1438 }, 1415 1439 }; 1416 1440 ops.apply_commit(input).unwrap(); ··· 1441 1465 did: did.clone(), 1442 1466 expected_root_cid: Some(mid_root.clone()), 1443 1467 new_root_cid: final_root.clone(), 1444 - new_rev: Tid::from("rev2".to_string()), 1468 + new_rev: Tid::new("3k2aaaaaaaaac").unwrap(), 1445 1469 new_block_cids: vec![], 1446 1470 obsolete_block_cids: vec![], 1447 1471 record_upserts: vec![], ··· 1457 1481 blobs: None, 1458 1482 blocks: None, 1459 1483 prev_data_cid: None, 1460 - rev: Some(Tid::from("rev2".to_string())), 1484 + rev: Some(Tid::new("3k2aaaaaaaaac").unwrap()), 1461 1485 }, 1462 1486 }; 1463 1487 ops.apply_commit(remove_like).unwrap();
+307 -40
crates/tranquil-store/src/metastore/recovery.rs
··· 1 1 use std::collections::HashSet; 2 2 3 3 use serde::{Deserialize, Serialize}; 4 - use tranquil_types::{Nsid, Rkey, Tid}; 4 + use tranquil_types::{AtUri, Nsid, Rkey, Tid}; 5 5 6 6 use super::backlink_ops::remove_backlinks_for_record; 7 7 use super::backlinks::{BacklinkValue, backlink_by_user_key, backlink_key, discriminant_to_path}; 8 8 use super::encoding::KeyReader; 9 9 use super::keys::{KeyTag, UserHash}; 10 - use super::records::{RecordValue, record_key}; 10 + use super::records::{RecordValue, record_by_cid_key, record_key}; 11 11 use super::repo_meta::{RepoMetaValue, repo_meta_key}; 12 12 use super::user_blocks::{user_block_key, user_block_user_prefix}; 13 13 use crate::metastore::MetastoreError; 14 14 15 - const MUTATION_SET_VERSION: u8 = 1; 15 + const MUTATION_SET_VERSION: u8 = 2; 16 + const MUTATION_SET_VERSION_UNVALIDATED: u8 = 1; 16 17 17 18 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] 18 19 pub struct CommitMutationSet { ··· 23 24 pub block_inserts: Vec<Vec<u8>>, 24 25 pub block_deletes: Vec<Vec<u8>>, 25 26 pub backlink_adds: Vec<BacklinkMutation>, 26 - pub backlink_remove_uris: Vec<String>, 27 + pub backlink_remove_uris: Vec<AtUri>, 27 28 } 28 29 29 30 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] 30 31 pub struct RecordMutationUpsert { 31 - pub collection: String, 32 - pub rkey: String, 32 + pub collection: Nsid, 33 + pub rkey: Rkey, 33 34 pub cid_bytes: Vec<u8>, 34 35 } 35 36 36 37 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] 37 38 pub struct RecordMutationDelete { 38 - pub collection: String, 39 - pub rkey: String, 39 + pub collection: Nsid, 40 + pub rkey: Rkey, 40 41 } 41 42 42 43 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] 43 44 pub struct BacklinkMutation { 44 - pub uri: String, 45 + pub uri: AtUri, 45 46 pub path: u8, 46 47 pub link_to: String, 47 48 } 48 49 50 + #[derive(Deserialize)] 51 + struct UnvalidatedCommitMutationSet { 52 + new_root_cid: Vec<u8>, 53 + new_rev: String, 54 + record_upserts: Vec<UnvalidatedRecordMutationUpsert>, 55 + record_deletes: Vec<UnvalidatedRecordMutationDelete>, 56 + block_inserts: Vec<Vec<u8>>, 57 + block_deletes: Vec<Vec<u8>>, 58 + backlink_adds: Vec<UnvalidatedBacklinkMutation>, 59 + backlink_remove_uris: Vec<String>, 60 + } 61 + 62 + #[derive(Deserialize)] 63 + struct UnvalidatedRecordMutationUpsert { 64 + collection: String, 65 + rkey: String, 66 + cid_bytes: Vec<u8>, 67 + } 68 + 69 + #[derive(Deserialize)] 70 + struct UnvalidatedRecordMutationDelete { 71 + collection: String, 72 + rkey: String, 73 + } 74 + 75 + #[derive(Deserialize)] 76 + struct UnvalidatedBacklinkMutation { 77 + uri: String, 78 + path: u8, 79 + link_to: String, 80 + } 81 + 82 + impl UnvalidatedCommitMutationSet { 83 + fn into_validated(self) -> Option<CommitMutationSet> { 84 + let warn = |field: &str, value: &str| { 85 + tracing::warn!( 86 + field, 87 + value, 88 + "version 1 CommitMutationSet has a value the current validators reject. \ 89 + Skipping the whole set rahter than replaying it in part." 90 + ); 91 + }; 92 + let new_rev = Tid::new(self.new_rev.clone()) 93 + .inspect_err(|_| warn("new_rev", &self.new_rev)) 94 + .ok()?; 95 + let record_upserts = self 96 + .record_upserts 97 + .into_iter() 98 + .map(|u| { 99 + Some(RecordMutationUpsert { 100 + collection: Nsid::new(u.collection.clone()) 101 + .inspect_err(|_| warn("record_upserts.collection", &u.collection)) 102 + .ok()?, 103 + rkey: Rkey::new(u.rkey.clone()) 104 + .inspect_err(|_| warn("record_upserts.rkey", &u.rkey)) 105 + .ok()?, 106 + cid_bytes: u.cid_bytes, 107 + }) 108 + }) 109 + .collect::<Option<Vec<_>>>()?; 110 + let record_deletes = self 111 + .record_deletes 112 + .into_iter() 113 + .map(|d| { 114 + Some(RecordMutationDelete { 115 + collection: Nsid::new(d.collection.clone()) 116 + .inspect_err(|_| warn("record_deletes.collection", &d.collection)) 117 + .ok()?, 118 + rkey: Rkey::new(d.rkey.clone()) 119 + .inspect_err(|_| warn("record_deletes.rkey", &d.rkey)) 120 + .ok()?, 121 + }) 122 + }) 123 + .collect::<Option<Vec<_>>>()?; 124 + let backlink_adds = self 125 + .backlink_adds 126 + .into_iter() 127 + .map(|b| { 128 + Some(BacklinkMutation { 129 + uri: AtUri::new(b.uri.clone()) 130 + .inspect_err(|_| warn("backlink_adds.uri", &b.uri)) 131 + .ok()?, 132 + path: b.path, 133 + link_to: b.link_to, 134 + }) 135 + }) 136 + .collect::<Option<Vec<_>>>()?; 137 + let backlink_remove_uris = self 138 + .backlink_remove_uris 139 + .into_iter() 140 + .map(|uri| { 141 + AtUri::new(uri.clone()) 142 + .inspect_err(|_| warn("backlink_remove_uris", &uri)) 143 + .ok() 144 + }) 145 + .collect::<Option<Vec<_>>>()?; 146 + Some(CommitMutationSet { 147 + new_root_cid: self.new_root_cid, 148 + new_rev, 149 + record_upserts, 150 + record_deletes, 151 + block_inserts: self.block_inserts, 152 + block_deletes: self.block_deletes, 153 + backlink_adds, 154 + backlink_remove_uris, 155 + }) 156 + } 157 + } 158 + 49 159 const MAX_MUTATION_SET_ENTRIES: usize = 50_000; 50 160 51 161 impl CommitMutationSet { ··· 91 201 None 92 202 } 93 203 }, 204 + MUTATION_SET_VERSION_UNVALIDATED => { 205 + match postcard::from_bytes::<UnvalidatedCommitMutationSet>(payload) { 206 + Ok(v) => v.into_validated(), 207 + Err(e) => { 208 + tracing::warn!(%e, "failed to deserialize version 1 CommitMutationSet payload"); 209 + None 210 + } 211 + } 212 + } 94 213 _ => { 95 214 tracing::warn!(version, "unknown CommitMutationSet version"); 96 215 None ··· 117 236 let meta_key = repo_meta_key(user_hash); 118 237 batch.insert(repo_data, meta_key.as_slice(), updated_meta.serialize()); 119 238 120 - mutation_set.record_upserts.iter().for_each(|u| { 121 - let key = record_key( 122 - user_hash, 123 - &Nsid::from(u.collection.clone()), 124 - &Rkey::from(u.rkey.clone()), 125 - ); 239 + mutation_set.record_upserts.iter().try_for_each(|u| { 240 + let key = record_key(user_hash, &u.collection, &u.rkey); 241 + let previous = stored_record(repo_data, key.as_slice())?; 242 + if let Some(prev) = previous 243 + .as_ref() 244 + .map(|p| &p.record_cid) 245 + .filter(|prev| *prev != &u.cid_bytes) 246 + { 247 + let stale = record_by_cid_key(user_hash, prev, &u.collection, &u.rkey); 248 + batch.remove(repo_data, stale.as_slice()); 249 + } 126 250 let value = RecordValue { 127 251 record_cid: u.cid_bytes.clone(), 128 - takedown_ref: None, 252 + takedown_ref: previous.and_then(|p| p.takedown_ref), 129 253 }; 254 + let reverse = record_by_cid_key(user_hash, &u.cid_bytes, &u.collection, &u.rkey); 130 255 batch.insert(repo_data, key.as_slice(), value.serialize()); 131 - }); 256 + batch.insert(repo_data, reverse.as_slice(), []); 257 + Ok::<(), MetastoreError>(()) 258 + })?; 132 259 133 - mutation_set.record_deletes.iter().for_each(|d| { 134 - let key = record_key( 135 - user_hash, 136 - &Nsid::from(d.collection.clone()), 137 - &Rkey::from(d.rkey.clone()), 138 - ); 260 + mutation_set.record_deletes.iter().try_for_each(|d| { 261 + let key = record_key(user_hash, &d.collection, &d.rkey); 262 + if let Some(prev) = stored_record(repo_data, key.as_slice())? { 263 + let reverse = record_by_cid_key(user_hash, &prev.record_cid, &d.collection, &d.rkey); 264 + batch.remove(repo_data, reverse.as_slice()); 265 + } 139 266 batch.remove(repo_data, key.as_slice()); 140 - }); 267 + Ok::<(), MetastoreError>(()) 268 + })?; 141 269 142 - mutation_set.block_inserts.iter().for_each(|cid_bytes| { 143 - let key = user_block_key(user_hash, &mutation_set.new_rev, cid_bytes); 144 - batch.insert(repo_data, key.as_slice(), []); 145 - }); 270 + let already_recorded: HashSet<Vec<u8>> = match mutation_set.block_inserts.is_empty() { 271 + true => HashSet::new(), 272 + false => repo_data 273 + .prefix(user_block_user_prefix(user_hash).as_slice()) 274 + .map(|guard| { 275 + let (key_bytes, _) = guard.into_inner().map_err(MetastoreError::Fjall)?; 276 + Ok(extract_cid_from_user_block_key(key_bytes.as_ref()).map(|c| c.to_vec())) 277 + }) 278 + .filter_map(Result::transpose) 279 + .collect::<Result<_, MetastoreError>>()?, 280 + }; 281 + mutation_set 282 + .block_inserts 283 + .iter() 284 + .filter(|cid_bytes| !cid_bytes.is_empty()) 285 + .filter(|cid_bytes| !already_recorded.contains(cid_bytes.as_slice())) 286 + .for_each(|cid_bytes| { 287 + let key = user_block_key(user_hash, &mutation_set.new_rev, cid_bytes); 288 + batch.insert(repo_data, key.as_slice(), []); 289 + }); 146 290 147 291 delete_user_blocks_by_cid_scan(batch, repo_data, user_hash, &mutation_set.block_deletes)?; 148 292 149 293 mutation_set 150 294 .backlink_remove_uris 151 295 .iter() 152 - .try_for_each(|uri_str| { 153 - let uri = tranquil_types::AtUri::from(uri_str.clone()); 296 + .try_for_each(|uri| { 154 297 let collection = uri.collection().ok_or(MetastoreError::CorruptData( 155 298 "backlink URI missing collection", 156 299 ))?; ··· 162 305 })?; 163 306 164 307 mutation_set.backlink_adds.iter().try_for_each(|bl| { 165 - let uri = tranquil_types::AtUri::from(bl.uri.clone()); 166 - let collection = uri.collection().ok_or(MetastoreError::CorruptData( 308 + let collection = bl.uri.collection().ok_or(MetastoreError::CorruptData( 167 309 "backlink URI missing collection", 168 310 ))?; 169 - let rkey = uri 311 + let rkey = bl 312 + .uri 170 313 .rkey() 171 314 .ok_or(MetastoreError::CorruptData("backlink URI missing rkey"))?; 172 315 ··· 181 324 Some(_) => { 182 325 let primary = backlink_key(&bl.link_to, user_hash, collection, rkey); 183 326 let value = BacklinkValue { 184 - source_uri: bl.uri.clone(), 327 + source_uri: bl.uri.as_str().to_owned(), 185 328 path: bl.path, 186 329 }; 187 330 batch.insert(indexes, primary.as_slice(), value.serialize()); ··· 194 337 }) 195 338 } 196 339 340 + fn stored_record( 341 + repo_data: &fjall::Keyspace, 342 + key: &[u8], 343 + ) -> Result<Option<RecordValue>, MetastoreError> { 344 + Ok(repo_data 345 + .get(key) 346 + .map_err(MetastoreError::Fjall)? 347 + .and_then(|raw| RecordValue::deserialize(&raw))) 348 + } 349 + 197 350 fn delete_user_blocks_by_cid_scan( 198 351 batch: &mut fjall::OwnedWriteBatch, 199 352 repo_data: &fjall::Keyspace, ··· 256 409 new_root_cid: vec![0x01, 0x71, 0x12, 0x20], 257 410 new_rev: Tid::new("3k2abcdefghij").unwrap(), 258 411 record_upserts: vec![RecordMutationUpsert { 259 - collection: "app.bsky.feed.post".to_owned(), 260 - rkey: "3k2abc".to_owned(), 412 + collection: Nsid::new("app.bsky.feed.post").unwrap(), 413 + rkey: Rkey::new("3k2abc").unwrap(), 261 414 cid_bytes: vec![0xDE, 0xAD], 262 415 }], 263 416 record_deletes: vec![RecordMutationDelete { 264 - collection: "app.bsky.feed.like".to_owned(), 265 - rkey: "3k2del".to_owned(), 417 + collection: Nsid::new("app.bsky.feed.like").unwrap(), 418 + rkey: Rkey::new("3k2del").unwrap(), 266 419 }], 267 420 block_inserts: vec![vec![0x01, 0x02]], 268 421 block_deletes: vec![vec![0x03, 0x04]], 269 422 backlink_adds: vec![BacklinkMutation { 270 - uri: "at://did:plc:olaren/app.bsky.feed.like/3k2abc".to_owned(), 423 + uri: AtUri::new("at://did:plc:olaren/app.bsky.feed.like/3k2abc").unwrap(), 271 424 path: 1, 272 425 link_to: "at://did:plc:teq/app.bsky.feed.post/3k2xyz".to_owned(), 273 426 }], 274 - backlink_remove_uris: vec!["at://did:plc:olaren/app.bsky.feed.like/3k2old".to_owned()], 427 + backlink_remove_uris: vec![ 428 + AtUri::new("at://did:plc:olaren/app.bsky.feed.like/3k2old").unwrap(), 429 + ], 275 430 }; 276 431 277 432 let bytes = ms.serialize().unwrap(); ··· 281 436 } 282 437 283 438 #[test] 439 + fn mutation_set_rejects_a_field_corrupted_into_a_structurally_valid_decode() { 440 + let ms = CommitMutationSet { 441 + new_root_cid: vec![0x01], 442 + new_rev: Tid::new("3k2abcdefghij").unwrap(), 443 + record_upserts: vec![RecordMutationUpsert { 444 + collection: Nsid::new("app.bsky.feed.post").unwrap(), 445 + rkey: Rkey::new("3k2abc").unwrap(), 446 + cid_bytes: vec![0x02], 447 + }], 448 + record_deletes: vec![], 449 + block_inserts: vec![], 450 + block_deletes: vec![], 451 + backlink_adds: vec![], 452 + backlink_remove_uris: vec![], 453 + }; 454 + 455 + let mut bytes = ms.serialize().unwrap(); 456 + let nsid_start = bytes 457 + .windows(b"app.bsky.feed.post".len()) 458 + .position(|w| w == b"app.bsky.feed.post") 459 + .expect("serialized form contains the collection"); 460 + bytes[nsid_start] = b'!'; 461 + 462 + assert!( 463 + CommitMutationSet::deserialize(&bytes).is_none(), 464 + "a corrupt field that still decodes as a string must be rejected" 465 + ); 466 + } 467 + 468 + #[test] 284 469 fn mutation_set_empty_roundtrip() { 285 470 let ms = CommitMutationSet { 286 471 new_root_cid: vec![], ··· 312 497 let mut bytes = ms.serialize().unwrap(); 313 498 bytes[0] = 99; 314 499 assert!(CommitMutationSet::deserialize(&bytes).is_none()); 500 + } 501 + 502 + fn v1_payload(rev: &str, collection: &str, rkey: &str) -> Vec<u8> { 503 + #[derive(Serialize)] 504 + struct V1 { 505 + new_root_cid: Vec<u8>, 506 + new_rev: String, 507 + record_upserts: Vec<V1Upsert>, 508 + record_deletes: Vec<(String, String)>, 509 + block_inserts: Vec<Vec<u8>>, 510 + block_deletes: Vec<Vec<u8>>, 511 + backlink_adds: Vec<(String, u8, String)>, 512 + backlink_remove_uris: Vec<String>, 513 + } 514 + #[derive(Serialize)] 515 + struct V1Upsert { 516 + collection: String, 517 + rkey: String, 518 + cid_bytes: Vec<u8>, 519 + } 520 + 521 + let payload = postcard::to_allocvec(&V1 { 522 + new_root_cid: vec![0x01, 0x71], 523 + new_rev: rev.to_owned(), 524 + record_upserts: vec![V1Upsert { 525 + collection: collection.to_owned(), 526 + rkey: rkey.to_owned(), 527 + cid_bytes: vec![0xDE, 0xAD], 528 + }], 529 + record_deletes: vec![], 530 + block_inserts: vec![vec![0x01, 0x02]], 531 + block_deletes: vec![], 532 + backlink_adds: vec![], 533 + backlink_remove_uris: vec![], 534 + }) 535 + .unwrap(); 536 + std::iter::once(MUTATION_SET_VERSION_UNVALIDATED) 537 + .chain(payload) 538 + .collect() 539 + } 540 + 541 + #[test] 542 + fn a_version_1_payload_written_by_an_older_binary_still_replays() { 543 + let decoded = CommitMutationSet::deserialize(&v1_payload( 544 + "3k2abcdefghij", 545 + "app.bsky.feed.post", 546 + "3k2abc", 547 + )) 548 + .expect("a version 1 payload with valid values decodes"); 549 + 550 + assert_eq!(decoded.new_rev.as_str(), "3k2abcdefghij"); 551 + assert_eq!(decoded.record_upserts[0].rkey.as_str(), "3k2abc"); 552 + assert_eq!(decoded.block_inserts, vec![vec![0x01, 0x02]]); 553 + } 554 + 555 + #[test] 556 + fn a_version_1_payload_the_current_validators_reject_is_skipped_not_fatal() { 557 + assert!( 558 + CommitMutationSet::deserialize(&v1_payload("3k2abcdefghij", "not/an/nsid", "3k2abc")) 559 + .is_none(), 560 + "a collection the validators reject skips the whole set instead of replaying it in part" 561 + ); 562 + assert!( 563 + CommitMutationSet::deserialize(&v1_payload("0", "app.bsky.feed.post", "3k2abc")) 564 + .is_none(), 565 + "a rev that isn't a TID skips the set" 566 + ); 567 + } 568 + 569 + #[test] 570 + fn a_current_payload_is_written_at_the_validated_version() { 571 + let ms = CommitMutationSet { 572 + new_root_cid: vec![0x01], 573 + new_rev: Tid::new("3k2abcdefghij").unwrap(), 574 + record_upserts: vec![], 575 + record_deletes: vec![], 576 + block_inserts: vec![], 577 + block_deletes: vec![], 578 + backlink_adds: vec![], 579 + backlink_remove_uris: vec![], 580 + }; 581 + assert_eq!(ms.serialize().unwrap()[0], MUTATION_SET_VERSION); 315 582 } 316 583 }
+26 -15
crates/tranquil-store/tests/metastore_crash.rs
··· 10 10 }; 11 11 use tranquil_store::metastore::{Metastore, MetastoreConfig}; 12 12 use tranquil_store::{sim_proptest_cases, sim_seed_range}; 13 - use tranquil_types::{CidLink, Did, Handle, Tid}; 13 + use tranquil_types::{AtUri, CidLink, Did, Handle, Nsid, Rkey, Tid}; 14 14 use uuid::Uuid; 15 15 16 16 const NAMES: &[&str] = &["olaren", "teq", "nel", "lyna", "bailey"]; ··· 27 27 28 28 fn test_did(seed: u64) -> Did { 29 29 let name = NAMES[(seed as usize) % NAMES.len()]; 30 - Did::from(format!("did:plc:{name}{seed}")) 30 + Did::new(format!("did:plc:{name}{seed}")).expect("generated DID is valid") 31 31 } 32 32 33 33 fn test_handle(seed: u64) -> Handle { ··· 46 46 Uuid::from_u128(seed as u128 | 0x4000_0000_0000_0000_8000_0000_0000_0000) 47 47 } 48 48 49 + const NSID_PATTERN: &str = "[a-z]{3,8}\\.[a-z]{3,8}\\.[a-z]{3,8}"; 50 + const RKEY_PATTERN: &str = "[a-z0-9]{3,10}"; 49 51 const TID_PATTERN: &str = "[234567abcdefghij][234567abcdefghijklmnopqrstuvwxyz]{12}"; 52 + const AT_URI_PATTERN: &str = 53 + "at://did:plc:[a-z]{3,8}/[a-z]{3,8}\\.[a-z]{3,8}\\.[a-z]{3,8}/[a-z0-9]{3,8}"; 50 54 51 55 fn arb_mutation_set() -> impl Strategy<Value = CommitMutationSet> { 52 56 let arb_upsert = ( 53 - "[a-z\\.]{5,20}", 54 - "[a-z0-9]{3,10}", 57 + NSID_PATTERN, 58 + RKEY_PATTERN, 55 59 prop::collection::vec(any::<u8>(), 4..36), 56 60 ) 57 61 .prop_map(|(collection, rkey, cid_bytes)| RecordMutationUpsert { 58 - collection, 59 - rkey, 62 + collection: Nsid::new(collection).expect("generated NSID is valid"), 63 + rkey: Rkey::new(rkey).expect("generated rkey is valid"), 60 64 cid_bytes, 61 65 }); 62 66 63 - let arb_delete = ("[a-z\\.]{5,20}", "[a-z0-9]{3,10}") 64 - .prop_map(|(collection, rkey)| RecordMutationDelete { collection, rkey }); 67 + let arb_delete = 68 + (NSID_PATTERN, RKEY_PATTERN).prop_map(|(collection, rkey)| RecordMutationDelete { 69 + collection: Nsid::new(collection).expect("generated NSID is valid"), 70 + rkey: Rkey::new(rkey).expect("generated rkey is valid"), 71 + }); 65 72 66 - let arb_backlink = ( 67 - "at://did:plc:[a-z]{3,8}/[a-z\\.]{5,20}/[a-z0-9]{3,8}", 68 - 0u8..4, 69 - "at://did:plc:[a-z]{3,8}/[a-z\\.]{5,20}/[a-z0-9]{3,8}", 70 - ) 71 - .prop_map(|(uri, path, link_to)| BacklinkMutation { uri, path, link_to }); 73 + let arb_backlink = (AT_URI_PATTERN, 0u8..4, AT_URI_PATTERN).prop_map(|(uri, path, link_to)| { 74 + BacklinkMutation { 75 + uri: AtUri::new(uri).expect("generated AT URI is valid"), 76 + path, 77 + link_to, 78 + } 79 + }); 72 80 73 81 ( 74 82 prop::collection::vec(any::<u8>(), 0..64), ··· 78 86 prop::collection::vec(prop::collection::vec(any::<u8>(), 4..36), 0..20), 79 87 prop::collection::vec(prop::collection::vec(any::<u8>(), 4..36), 0..20), 80 88 prop::collection::vec(arb_backlink, 0..5), 81 - prop::collection::vec("at://did:plc:[a-z]{3,8}/[a-z\\.]{5,20}/[a-z0-9]{3,8}", 0..5), 89 + prop::collection::vec( 90 + AT_URI_PATTERN.prop_map(|uri| AtUri::new(uri).expect("generated AT URI is valid")), 91 + 0..5, 92 + ), 82 93 ) 83 94 .prop_map( 84 95 |(