Our Personal Data Server from scratch!
0

Configure Feed

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

store: record-by-cid reverse index

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

author
Lewis
date (Jul 24, 2026, 8:27 PM +0300) commit bc052830 parent 6d59fe0d change-id wxwzqnpt
+903 -133
+25
.sqlx/query-105807a41c7337e7aa46bace29ab613030fd4fbf6845baddab9c0b2009972c02.json
··· 1 + { 2 + "db_name": "PostgreSQL", 3 + "query": "\n SELECT DISTINCT r.record_cid AS \"record_cid!\"\n FROM records r\n WHERE r.repo_id = $1\n AND r.record_cid = ANY($2)\n AND NOT EXISTS (\n SELECT 1 FROM UNNEST($3::text[], $4::text[]) AS k(collection, rkey)\n WHERE k.collection = r.collection AND k.rkey = r.rkey\n )\n ", 4 + "describe": { 5 + "columns": [ 6 + { 7 + "ordinal": 0, 8 + "name": "record_cid!", 9 + "type_info": "Text" 10 + } 11 + ], 12 + "parameters": { 13 + "Left": [ 14 + "Uuid", 15 + "TextArray", 16 + "TextArray", 17 + "TextArray" 18 + ] 19 + }, 20 + "nullable": [ 21 + false 22 + ] 23 + }, 24 + "hash": "105807a41c7337e7aa46bace29ab613030fd4fbf6845baddab9c0b2009972c02" 25 + }
+7
crates/tranquil-db-traits/src/repo.rs
··· 410 410 async fn get_record_by_cid(&self, cid: &CidLink) 411 411 -> Result<Option<RecordWithTakedown>, DbError>; 412 412 413 + async fn referenced_record_cids( 414 + &self, 415 + repo_id: Uuid, 416 + cids: &[CidLink], 417 + excluded_keys: &[(&Nsid, &Rkey)], 418 + ) -> Result<Vec<CidLink>, DbError>; 419 + 413 420 async fn set_record_takedown( 414 421 &self, 415 422 cid: &CidLink,
+39
crates/tranquil-db/src/postgres/repo.rs
··· 721 721 })) 722 722 } 723 723 724 + async fn referenced_record_cids( 725 + &self, 726 + repo_id: Uuid, 727 + cids: &[CidLink], 728 + excluded_keys: &[(&Nsid, &Rkey)], 729 + ) -> Result<Vec<CidLink>, DbError> { 730 + if cids.is_empty() { 731 + return Ok(Vec::new()); 732 + } 733 + 734 + let cid_strs: Vec<String> = cids.iter().map(|c| c.as_str().to_owned()).collect(); 735 + let (excluded_collections, excluded_rkeys): (Vec<String>, Vec<String>) = excluded_keys 736 + .iter() 737 + .map(|(collection, rkey)| (collection.as_str().to_owned(), rkey.as_str().to_owned())) 738 + .unzip(); 739 + 740 + let rows = sqlx::query_scalar!( 741 + r#" 742 + SELECT DISTINCT r.record_cid AS "record_cid!" 743 + FROM records r 744 + WHERE r.repo_id = $1 745 + AND r.record_cid = ANY($2) 746 + AND NOT EXISTS ( 747 + SELECT 1 FROM UNNEST($3::text[], $4::text[]) AS k(collection, rkey) 748 + WHERE k.collection = r.collection AND k.rkey = r.rkey 749 + ) 750 + "#, 751 + repo_id, 752 + &cid_strs, 753 + &excluded_collections, 754 + &excluded_rkeys 755 + ) 756 + .fetch_all(&self.pool) 757 + .await 758 + .map_err(map_sqlx_error)?; 759 + 760 + column_vec(rows, col::RECORDS_RECORD_CID) 761 + } 762 + 724 763 async fn set_record_takedown( 725 764 &self, 726 765 cid: &CidLink,
+21
crates/tranquil-store/src/metastore/client.rs
··· 400 400 recv(rx).await 401 401 } 402 402 403 + async fn referenced_record_cids( 404 + &self, 405 + repo_id: Uuid, 406 + cids: &[CidLink], 407 + excluded_keys: &[(&Nsid, &Rkey)], 408 + ) -> Result<Vec<CidLink>, DbError> { 409 + let (tx, rx) = oneshot::channel(); 410 + self.pool.send(MetastoreRequest::Record( 411 + RecordRequest::ReferencedRecordCids { 412 + repo_id, 413 + cids: cids.to_vec(), 414 + excluded_keys: excluded_keys 415 + .iter() 416 + .map(|(collection, rkey)| ((*collection).clone(), (*rkey).clone())) 417 + .collect(), 418 + tx, 419 + }, 420 + ))?; 421 + recv(rx).await 422 + } 423 + 403 424 async fn set_record_takedown( 404 425 &self, 405 426 cid: &CidLink,
+2 -1
crates/tranquil-store/src/metastore/commit_ops.rs
··· 181 181 .collect(); 182 182 183 183 self.record_ops 184 - .delete_records(&mut batch, user_hash, &deletes); 184 + .delete_records(&mut batch, user_hash, &deletes) 185 + .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 185 186 186 187 self.user_block_ops 187 188 .insert_user_blocks(&mut batch, user_hash, &input.new_block_cids, &input.new_rev)
+27 -2
crates/tranquil-store/src/metastore/handler.rs
··· 324 324 cid: CidLink, 325 325 tx: Tx<Option<tranquil_db_traits::RecordWithTakedown>>, 326 326 }, 327 + ReferencedRecordCids { 328 + repo_id: Uuid, 329 + cids: Vec<CidLink>, 330 + excluded_keys: Vec<(Nsid, Rkey)>, 331 + tx: Tx<Vec<CidLink>>, 332 + }, 327 333 SetRecordTakedown { 328 334 cid: CidLink, 329 335 takedown_ref: Option<String>, ··· 342 348 | Self::ListRecords { repo_id, .. } 343 349 | Self::GetAllRecords { repo_id, .. } 344 350 | Self::ListCollections { repo_id, .. } 345 - | Self::CountRecords { repo_id, .. } => uuid_to_routing(user_hashes, repo_id), 351 + | Self::CountRecords { repo_id, .. } 352 + | Self::ReferencedRecordCids { repo_id, .. } => uuid_to_routing(user_hashes, repo_id), 346 353 Self::CountAllRecords { .. } | Self::GetRecordByCid { .. } => Routing::Global, 347 354 Self::SetRecordTakedown { 348 355 scope_user: Some(user_id), ··· 2757 2764 state 2758 2765 .metastore 2759 2766 .record_ops() 2760 - .delete_records(&mut batch, user_hash, &deletes); 2767 + .delete_records(&mut batch, user_hash, &deletes) 2768 + .map_err(metastore_to_db)?; 2761 2769 batch.commit().map_err(|e| DbError::Query(e.to_string())) 2762 2770 })(); 2763 2771 let _ = tx.send(result); ··· 2858 2866 .record_ops() 2859 2867 .get_record_by_cid(&cid, None) 2860 2868 .map(|opt| opt.map(convert_record_with_takedown)) 2869 + .map_err(metastore_to_db); 2870 + let _ = tx.send(result); 2871 + } 2872 + RecordRequest::ReferencedRecordCids { 2873 + repo_id, 2874 + cids, 2875 + excluded_keys, 2876 + tx, 2877 + } => { 2878 + let excluded: Vec<(&Nsid, &Rkey)> = excluded_keys 2879 + .iter() 2880 + .map(|(collection, rkey)| (collection, rkey)) 2881 + .collect(); 2882 + let result = state 2883 + .metastore 2884 + .record_ops() 2885 + .referenced_record_cids(repo_id, &cids, &excluded) 2861 2886 .map_err(metastore_to_db); 2862 2887 let _ = tx.send(result); 2863 2888 }
+775 -129
crates/tranquil-store/src/metastore/record_ops.rs
··· 1 + use std::collections::{BTreeMap, HashSet}; 1 2 use std::sync::Arc; 2 3 3 4 use fjall::Keyspace; ··· 8 9 use super::encoding::{KeyReader, exclusive_upper_bound}; 9 10 use super::keys::UserHash; 10 11 use super::records::{ 11 - RecordValue, record_collection_prefix, record_key, record_user_prefix, records_prefix, 12 + RecordValue, record_by_cid_built_key, record_by_cid_index_prefix, record_by_cid_key, 13 + record_by_cid_key_from_suffix, record_by_cid_prefix, record_by_cid_suffix, 14 + record_by_cid_user_prefix, record_collection_prefix, record_key, 15 + record_key_user_hash_and_suffix, record_user_prefix, records_prefix, 12 16 }; 13 17 use super::repo_ops::{bytes_to_cid_link, cid_link_to_bytes}; 14 - use super::scan::{count_prefix, delete_all_by_prefix, point_lookup}; 18 + use super::scan::{ 19 + count_prefix, delete_all_by_prefix, delete_all_by_prefix_chunked, point_lookup, prefix_entries, 20 + stage_in_chunks, 21 + }; 15 22 use super::user_hash::UserHashMap; 16 23 17 24 use tranquil_types::{CidLink, Nsid, Rkey}; ··· 77 84 user_hash: UserHash, 78 85 records: &[RecordWrite<'_>], 79 86 ) -> Result<(), MetastoreError> { 80 - records.iter().try_for_each(|rec| { 81 - let key = record_key(user_hash, rec.collection, rec.rkey); 82 - let cid_bytes = cid_link_to_bytes(rec.cid)?; 83 - let existing_takedown = self 84 - .repo_data 85 - .get(key.as_slice()) 86 - .map_err(MetastoreError::Fjall)? 87 - .and_then(|raw| RecordValue::deserialize(&raw)) 88 - .and_then(|v| v.takedown_ref); 89 - let value = RecordValue { 90 - record_cid: cid_bytes, 91 - takedown_ref: existing_takedown, 92 - }; 93 - batch.insert(&self.repo_data, key.as_slice(), value.serialize()); 94 - Ok::<(), MetastoreError>(()) 95 - }) 87 + let last_write_per_key: BTreeMap<(&Nsid, &Rkey), &CidLink> = records 88 + .iter() 89 + .map(|rec| ((rec.collection, rec.rkey), rec.cid)) 90 + .collect(); 91 + 92 + last_write_per_key 93 + .into_iter() 94 + .try_for_each(|((collection, rkey), cid)| { 95 + let key = record_key(user_hash, collection, rkey); 96 + let cid_bytes = cid_link_to_bytes(cid)?; 97 + let existing = self 98 + .repo_data 99 + .get(key.as_slice()) 100 + .map_err(MetastoreError::Fjall)? 101 + .and_then(|raw| RecordValue::deserialize(&raw)); 102 + if let Some(prev) = &existing 103 + && prev.record_cid != cid_bytes 104 + { 105 + let stale = record_by_cid_key(user_hash, &prev.record_cid, collection, rkey); 106 + batch.remove(&self.repo_data, stale.as_slice()); 107 + } 108 + let value = RecordValue { 109 + record_cid: cid_bytes, 110 + takedown_ref: existing.and_then(|v| v.takedown_ref), 111 + }; 112 + let reverse = record_by_cid_key(user_hash, &value.record_cid, collection, rkey); 113 + batch.insert(&self.repo_data, key.as_slice(), value.serialize()); 114 + batch.insert(&self.repo_data, reverse.as_slice(), []); 115 + Ok::<(), MetastoreError>(()) 116 + }) 96 117 } 97 118 98 119 pub fn delete_records( ··· 100 121 batch: &mut fjall::OwnedWriteBatch, 101 122 user_hash: UserHash, 102 123 records: &[RecordDelete<'_>], 103 - ) { 104 - records.iter().for_each(|rec| { 124 + ) -> Result<(), MetastoreError> { 125 + records.iter().try_for_each(|rec| { 105 126 let key = record_key(user_hash, rec.collection, rec.rkey); 127 + let existing = self 128 + .repo_data 129 + .get(key.as_slice()) 130 + .map_err(MetastoreError::Fjall)? 131 + .and_then(|raw| RecordValue::deserialize(&raw)); 132 + if let Some(prev) = existing { 133 + let reverse = 134 + record_by_cid_key(user_hash, &prev.record_cid, rec.collection, rec.rkey); 135 + batch.remove(&self.repo_data, reverse.as_slice()); 136 + } 106 137 batch.remove(&self.repo_data, key.as_slice()); 107 - }); 138 + Ok(()) 139 + }) 108 140 } 109 141 110 142 pub fn delete_all_records( ··· 112 144 batch: &mut fjall::OwnedWriteBatch, 113 145 user_hash: UserHash, 114 146 ) -> Result<(), MetastoreError> { 115 - let prefix = record_user_prefix(user_hash); 116 - delete_all_by_prefix(&self.repo_data, batch, prefix.as_slice()) 147 + delete_all_by_prefix( 148 + &self.repo_data, 149 + batch, 150 + record_by_cid_user_prefix(user_hash).as_slice(), 151 + )?; 152 + delete_all_by_prefix( 153 + &self.repo_data, 154 + batch, 155 + record_user_prefix(user_hash).as_slice(), 156 + ) 157 + } 158 + 159 + pub fn referenced_record_cids( 160 + &self, 161 + repo_id: Uuid, 162 + cids: &[CidLink], 163 + excluded_keys: &[(&Nsid, &Rkey)], 164 + ) -> Result<Vec<CidLink>, MetastoreError> { 165 + let user_hash = self 166 + .user_hashes 167 + .get(&repo_id) 168 + .ok_or(MetastoreError::InvalidInput( 169 + "repo id has no user hash mapping", 170 + ))?; 171 + 172 + let excluded: HashSet<SmallVec<[u8; 128]>> = excluded_keys 173 + .iter() 174 + .map(|(collection, rkey)| record_by_cid_suffix(collection, rkey)) 175 + .collect(); 176 + 177 + cids.iter() 178 + .map(|cid| { 179 + let cid_bytes = cid_link_to_bytes(cid)?; 180 + let prefix = record_by_cid_prefix(user_hash, &cid_bytes); 181 + let referenced = self 182 + .repo_data 183 + .prefix(prefix.as_slice()) 184 + .find_map(|guard| match guard.into_inner() { 185 + Err(e) => Some(Err(MetastoreError::Fjall(e))), 186 + Ok((key_bytes, _)) => key_bytes 187 + .get(prefix.len()..) 188 + .is_none_or(|suffix| !excluded.contains(suffix)) 189 + .then_some(Ok::<(), MetastoreError>(())), 190 + }) 191 + .transpose()? 192 + .is_some(); 193 + Ok(referenced.then(|| cid.clone())) 194 + }) 195 + .filter_map(Result::transpose) 196 + .collect() 117 197 } 118 198 119 199 pub fn get_record_cid( ··· 286 366 }; 287 367 let (key_bytes, _) = guard.into_inner().map_err(MetastoreError::Fjall)?; 288 368 let collection = parse_record_key_collection(&key_bytes) 289 - .map(Nsid::from) 290 369 .ok_or(MetastoreError::CorruptData("invalid record key"))?; 291 370 let coll_prefix = record_collection_prefix(user_hash, &collection); 292 371 seek_from = exclusive_upper_bound(coll_prefix.as_slice()) ··· 339 418 340 419 match value.record_cid == target_bytes { 341 420 true => { 342 - let (coll_str, rkey_str) = match parse_record_key_fields(&key_bytes) { 421 + let (collection, rkey) = match parse_record_key_fields(&key_bytes) { 343 422 Some(pair) => pair, 344 423 None => { 345 424 return Some(Err(MetastoreError::CorruptData( ··· 365 444 }; 366 445 Some(Ok(RecordWithTakedown { 367 446 id: user_id, 368 - collection: Nsid::from(coll_str), 369 - rkey: Rkey::from(rkey_str), 447 + collection, 448 + rkey, 370 449 takedown_ref: value.takedown_ref, 371 450 })) 372 451 } ··· 430 509 .ok_or(MetastoreError::CorruptData("invalid record key"))?; 431 510 let cid = bytes_to_cid_link(&value.record_cid)?; 432 511 Ok(RecordInfo { 433 - rkey: Rkey::from(rkey), 512 + rkey, 434 513 record_cid: cid, 435 514 }) 436 515 } ··· 445 524 .ok_or(MetastoreError::CorruptData("invalid record key"))?; 446 525 let cid = bytes_to_cid_link(&value.record_cid)?; 447 526 Ok(FullRecordInfo { 448 - collection: Nsid::from(collection), 449 - rkey: Rkey::from(rkey), 527 + collection, 528 + rkey, 450 529 record_cid: cid, 451 530 }) 452 531 } 453 532 454 - fn parse_record_key_fields(key_bytes: &[u8]) -> Option<(String, String)> { 533 + fn parse_record_key_fields(key_bytes: &[u8]) -> Option<(Nsid, Rkey)> { 455 534 let mut reader = KeyReader::new(key_bytes); 456 535 let _tag = reader.tag()?; 457 536 let _user_hash = reader.u64()?; 458 - let collection = reader.string()?; 459 - let rkey = reader.string()?; 537 + let collection = Nsid::new(reader.string()?).ok()?; 538 + let rkey = Rkey::new(reader.string()?).ok()?; 460 539 Some((collection, rkey)) 461 540 } 462 541 463 - fn parse_record_key_collection(key_bytes: &[u8]) -> Option<String> { 542 + fn parse_record_key_collection(key_bytes: &[u8]) -> Option<Nsid> { 464 543 let mut reader = KeyReader::new(key_bytes); 465 544 let _tag = reader.tag()?; 466 545 let _user_hash = reader.u64()?; 467 - reader.string() 546 + Nsid::new(reader.string()?).ok() 547 + } 548 + 549 + // Committing in chunks like this means 550 + // the rebuild leaves entries under the index prefix 551 + // that only cover some of the records 552 + // if anything interrupts it part-way through 553 + // where the marker goes in last 554 + // and is the only evidence a rebuild ever finished 555 + // so its absence means wiping whatever is there and starting over 556 + // as opposed to trusting a partial index. 557 + pub fn rebuild_record_by_cid_index( 558 + db: &fjall::Database, 559 + repo_data: &Keyspace, 560 + ) -> Result<usize, MetastoreError> { 561 + let built_key = record_by_cid_built_key(); 562 + if repo_data 563 + .get(built_key.as_slice()) 564 + .map_err(MetastoreError::Fjall)? 565 + .is_some() 566 + { 567 + return Ok(0); 568 + } 569 + 570 + const COMMIT_EVERY: usize = 10_000; 571 + let discarded = delete_all_by_prefix_chunked( 572 + db, 573 + repo_data, 574 + record_by_cid_index_prefix().as_slice(), 575 + COMMIT_EVERY, 576 + )?; 577 + if discarded > 0 { 578 + tracing::warn!( 579 + entries = discarded, 580 + "discarded a record_by_cid index that an interrupted rebuild left wihtout its marker" 581 + ); 582 + } 583 + 584 + let mut skipped = 0usize; 585 + let written = stage_in_chunks( 586 + db, 587 + prefix_entries(repo_data, records_prefix().as_slice()), 588 + COMMIT_EVERY, 589 + |batch, (key_bytes, val_bytes)| { 590 + match ( 591 + record_key_user_hash_and_suffix(&key_bytes), 592 + RecordValue::deserialize(&val_bytes), 593 + ) { 594 + (Some((user_hash, suffix)), Some(value)) => { 595 + let reverse = 596 + record_by_cid_key_from_suffix(user_hash, &value.record_cid, suffix); 597 + batch.insert(repo_data, reverse.as_slice(), []); 598 + } 599 + _ => skipped += 1, 600 + } 601 + Ok(()) 602 + }, 603 + )?; 604 + if skipped > 0 { 605 + tracing::warn!( 606 + skipped, 607 + "record_by_cid rebuild skipped rows whose key or value doesn't decode. \ 608 + Those records are unreadable through every other path too." 609 + ); 610 + } 611 + 612 + let mut batch = db.batch(); 613 + batch.insert(repo_data, built_key.as_slice(), []); 614 + batch.commit()?; 615 + Ok(written.saturating_sub(skipped)) 468 616 } 469 617 470 618 fn parse_record_key_user_hash(key_bytes: &[u8]) -> Option<UserHash> { ··· 492 640 CidLink::from_cid(&c) 493 641 } 494 642 643 + fn test_rev(seq: u64) -> tranquil_types::Tid { 644 + const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz"; 645 + let s: String = (0..13) 646 + .rev() 647 + .map(|i| ALPHABET[((seq >> (i * 5)) & 0x1F) as usize] as char) 648 + .collect(); 649 + tranquil_types::Tid::new(s).expect("generated TID is valid") 650 + } 651 + 495 652 fn open_fresh() -> (tempfile::TempDir, Metastore) { 496 653 let dir = tempfile::TempDir::new().unwrap(); 497 654 let ms = Metastore::open(dir.path(), test_config()).unwrap(); ··· 499 656 } 500 657 501 658 fn test_did(name: &str) -> tranquil_types::Did { 502 - tranquil_types::Did::from(format!("did:plc:{name}")) 659 + tranquil_types::Did::new(format!("did:plc:{name}")).expect("test DID is well-formed") 503 660 } 504 661 505 662 fn test_handle(name: &str) -> tranquil_types::Handle { 506 - tranquil_types::Handle::from(format!("{name}.test.invalid")) 663 + tranquil_types::Handle::new(format!("{name}.oyster.cafe")).expect("test handle is valid") 664 + } 665 + 666 + fn test_nsid(name: impl Into<String>) -> Nsid { 667 + Nsid::new(name).expect("test collection is a valid NSID") 507 668 } 508 669 509 - fn test_rev(seq: u64) -> tranquil_types::Tid { 510 - const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz"; 511 - let s: String = (0..13) 512 - .rev() 513 - .map(|i| ALPHABET[((seq >> (i * 5)) & 0x1F) as usize] as char) 514 - .collect(); 515 - tranquil_types::Tid::new(s).expect("generated TID is valid") 670 + fn test_rkey(key: impl Into<String>) -> Rkey { 671 + Rkey::new(key).expect("test record key is a valid rkey") 516 672 } 517 673 518 674 fn setup_user(ms: &Metastore) -> (Uuid, super::super::keys::UserHash) { 519 675 let user_id = Uuid::new_v4(); 520 - let did = test_did("testuser"); 521 - let handle = test_handle("testuser"); 676 + let did = test_did("limpet"); 677 + let handle = test_handle("limpet"); 522 678 let cid = test_cid_link(0); 523 679 ms.repo_ops() 524 680 .create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(0)) ··· 565 721 let (user_id, user_hash) = setup_user(&ms); 566 722 let rec_ops = ms.record_ops(); 567 723 568 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 569 - let rkey = Rkey::from("3k2abcd".to_string()); 724 + let collection = test_nsid("app.bsky.feed.post"); 725 + let rkey = test_rkey("3k2abcd"); 570 726 let cid = test_cid_link(1); 571 727 572 728 let mut batch = ms.database().batch(); ··· 585 741 let (user_id, _) = setup_user(&ms); 586 742 let rec_ops = ms.record_ops(); 587 743 588 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 589 - let rkey = Rkey::from("nonexistent".to_string()); 744 + let collection = test_nsid("app.bsky.feed.post"); 745 + let rkey = test_rkey("nonexistent"); 590 746 assert!( 591 747 rec_ops 592 748 .get_record_cid(user_id, &collection, &rkey) ··· 601 757 let (user_id, user_hash) = setup_user(&ms); 602 758 let rec_ops = ms.record_ops(); 603 759 604 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 605 - let rkey = Rkey::from("3k2abcd".to_string()); 760 + let collection = test_nsid("app.bsky.feed.post"); 761 + let rkey = test_rkey("3k2abcd"); 606 762 let cid1 = test_cid_link(1); 607 763 let cid2 = test_cid_link(2); 608 764 ··· 628 784 let (user_id, user_hash) = setup_user(&ms); 629 785 let rec_ops = ms.record_ops(); 630 786 631 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 632 - let rkey = Rkey::from("3k2abcd".to_string()); 787 + let collection = test_nsid("app.bsky.feed.post"); 788 + let rkey = test_rkey("3k2abcd"); 633 789 let cid = test_cid_link(1); 634 790 635 791 let mut batch = ms.database().batch(); ··· 639 795 batch.commit().unwrap(); 640 796 641 797 let mut batch = ms.database().batch(); 642 - rec_ops.delete_records(&mut batch, user_hash, &[rd(&collection, &rkey)]); 798 + rec_ops 799 + .delete_records(&mut batch, user_hash, &[rd(&collection, &rkey)]) 800 + .unwrap(); 643 801 batch.commit().unwrap(); 644 802 645 803 assert!( ··· 656 814 let (user_id, user_hash) = setup_user(&ms); 657 815 let rec_ops = ms.record_ops(); 658 816 659 - let coll1 = Nsid::from("app.bsky.feed.post".to_string()); 660 - let coll2 = Nsid::from("app.bsky.feed.like".to_string()); 661 - let rkey_a = Rkey::from("a".to_string()); 662 - let rkey_b = Rkey::from("b".to_string()); 817 + let coll1 = test_nsid("app.bsky.feed.post"); 818 + let coll2 = test_nsid("app.bsky.feed.like"); 819 + let rkey_a = test_rkey("a"); 820 + let rkey_b = test_rkey("b"); 663 821 let cid1 = test_cid_link(1); 664 822 let cid2 = test_cid_link(2); 665 823 ··· 688 846 let (user_id, user_hash) = setup_user(&ms); 689 847 let rec_ops = ms.record_ops(); 690 848 691 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 692 - let rkeys: Vec<Rkey> = (0..5).map(|i| Rkey::from(format!("rkey{i:03}"))).collect(); 849 + let collection = test_nsid("app.bsky.feed.post"); 850 + let rkeys: Vec<Rkey> = (0..5).map(|i| test_rkey(format!("rkey{i:03}"))).collect(); 693 851 let cids: Vec<CidLink> = (0..5).map(|i| test_cid_link(i + 1)).collect(); 694 852 695 853 let writes: Vec<RecordWrite<'_>> = rkeys ··· 719 877 let (user_id, user_hash) = setup_user(&ms); 720 878 let rec_ops = ms.record_ops(); 721 879 722 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 723 - let rkeys: Vec<Rkey> = (0..5).map(|i| Rkey::from(format!("rkey{i:03}"))).collect(); 880 + let collection = test_nsid("app.bsky.feed.post"); 881 + let rkeys: Vec<Rkey> = (0..5).map(|i| test_rkey(format!("rkey{i:03}"))).collect(); 724 882 let cids: Vec<CidLink> = (0..5).map(|i| test_cid_link(i + 1)).collect(); 725 883 726 884 let writes: Vec<RecordWrite<'_>> = rkeys ··· 735 893 .unwrap(); 736 894 batch.commit().unwrap(); 737 895 738 - let cursor = Rkey::from("rkey003".to_string()); 896 + let cursor = test_rkey("rkey003"); 739 897 let results = rec_ops 740 898 .list_records(&lrq( 741 899 user_id, ··· 759 917 let (user_id, user_hash) = setup_user(&ms); 760 918 let rec_ops = ms.record_ops(); 761 919 762 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 763 - let rkeys: Vec<Rkey> = (0..5).map(|i| Rkey::from(format!("rkey{i:03}"))).collect(); 920 + let collection = test_nsid("app.bsky.feed.post"); 921 + let rkeys: Vec<Rkey> = (0..5).map(|i| test_rkey(format!("rkey{i:03}"))).collect(); 764 922 let cids: Vec<CidLink> = (0..5).map(|i| test_cid_link(i + 1)).collect(); 765 923 766 924 let writes: Vec<RecordWrite<'_>> = rkeys ··· 790 948 let (user_id, user_hash) = setup_user(&ms); 791 949 let rec_ops = ms.record_ops(); 792 950 793 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 794 - let rkeys: Vec<Rkey> = (0..5).map(|i| Rkey::from(format!("rkey{i:03}"))).collect(); 951 + let collection = test_nsid("app.bsky.feed.post"); 952 + let rkeys: Vec<Rkey> = (0..5).map(|i| test_rkey(format!("rkey{i:03}"))).collect(); 795 953 let cids: Vec<CidLink> = (0..5).map(|i| test_cid_link(i + 1)).collect(); 796 954 797 955 let writes: Vec<RecordWrite<'_>> = rkeys ··· 806 964 .unwrap(); 807 965 batch.commit().unwrap(); 808 966 809 - let cursor = Rkey::from("rkey001".to_string()); 967 + let cursor = test_rkey("rkey001"); 810 968 let results = rec_ops 811 969 .list_records(&lrq( 812 970 user_id, ··· 830 988 let (user_id, user_hash) = setup_user(&ms); 831 989 let rec_ops = ms.record_ops(); 832 990 833 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 834 - let rkeys: Vec<Rkey> = (0..10).map(|i| Rkey::from(format!("rkey{i:03}"))).collect(); 991 + let collection = test_nsid("app.bsky.feed.post"); 992 + let rkeys: Vec<Rkey> = (0..10).map(|i| test_rkey(format!("rkey{i:03}"))).collect(); 835 993 let cids: Vec<CidLink> = (0..10).map(|i| test_cid_link(i + 1)).collect(); 836 994 837 995 let writes: Vec<RecordWrite<'_>> = rkeys ··· 846 1004 .unwrap(); 847 1005 batch.commit().unwrap(); 848 1006 849 - let rkey_start = Rkey::from("rkey003".to_string()); 850 - let rkey_end = Rkey::from("rkey006".to_string()); 1007 + let rkey_start = test_rkey("rkey003"); 1008 + let rkey_end = test_rkey("rkey006"); 851 1009 let results = rec_ops 852 1010 .list_records(&lrq( 853 1011 user_id, ··· 870 1028 let (user_id, user_hash) = setup_user(&ms); 871 1029 let rec_ops = ms.record_ops(); 872 1030 873 - let coll1 = Nsid::from("app.bsky.feed.like".to_string()); 874 - let coll2 = Nsid::from("app.bsky.feed.post".to_string()); 875 - let rkey_a = Rkey::from("a".to_string()); 876 - let rkey_b = Rkey::from("b".to_string()); 877 - let rkey_c = Rkey::from("c".to_string()); 1031 + let coll1 = test_nsid("app.bsky.feed.like"); 1032 + let coll2 = test_nsid("app.bsky.feed.post"); 1033 + let rkey_a = test_rkey("a"); 1034 + let rkey_b = test_rkey("b"); 1035 + let rkey_c = test_rkey("c"); 878 1036 let cid1 = test_cid_link(1); 879 1037 let cid2 = test_cid_link(2); 880 1038 let cid3 = test_cid_link(3); ··· 906 1064 let (user_id, user_hash) = setup_user(&ms); 907 1065 let rec_ops = ms.record_ops(); 908 1066 909 - let coll1 = Nsid::from("app.bsky.feed.like".to_string()); 910 - let coll2 = Nsid::from("app.bsky.feed.post".to_string()); 911 - let coll3 = Nsid::from("app.bsky.graph.follow".to_string()); 1067 + let coll1 = test_nsid("app.bsky.feed.like"); 1068 + let coll2 = test_nsid("app.bsky.feed.post"); 1069 + let coll3 = test_nsid("app.bsky.graph.follow"); 912 1070 let rkeys: Vec<Rkey> = ["a", "b", "c", "d", "e"] 913 1071 .iter() 914 - .map(|s| Rkey::from(s.to_string())) 1072 + .map(|s| test_rkey(s.to_string())) 915 1073 .collect(); 916 1074 let cids: Vec<CidLink> = (1..=5).map(test_cid_link).collect(); 917 1075 ··· 947 1105 assert_eq!(rec_ops.count_records(user_id).unwrap(), 0); 948 1106 assert_eq!(rec_ops.count_all_records().unwrap(), 0); 949 1107 950 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 951 - let rkey_a = Rkey::from("a".to_string()); 952 - let rkey_b = Rkey::from("b".to_string()); 1108 + let collection = test_nsid("app.bsky.feed.post"); 1109 + let rkey_a = test_rkey("a"); 1110 + let rkey_b = test_rkey("b"); 953 1111 let cid1 = test_cid_link(1); 954 1112 let cid2 = test_cid_link(2); 955 1113 ··· 1005 1163 .unwrap(); 1006 1164 let hash2 = ms.user_hashes().get(&user2).unwrap(); 1007 1165 1008 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1009 - let rkey_a = Rkey::from("a".to_string()); 1010 - let rkey_b = Rkey::from("b".to_string()); 1166 + let collection = test_nsid("app.bsky.feed.post"); 1167 + let rkey_a = test_rkey("a"); 1168 + let rkey_b = test_rkey("b"); 1011 1169 let cid1 = test_cid_link(1); 1012 1170 let cid2 = test_cid_link(2); 1013 1171 ··· 1031 1189 let (user_id, user_hash) = setup_user(&ms); 1032 1190 let rec_ops = ms.record_ops(); 1033 1191 1034 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1035 - let rkey = Rkey::from("r1".to_string()); 1192 + let collection = test_nsid("app.bsky.feed.post"); 1193 + let rkey = test_rkey("r1"); 1036 1194 let cid = test_cid_link(42); 1037 1195 1038 1196 let mut batch = ms.database().batch(); ··· 1066 1224 let (user_id, user_hash) = setup_user(&ms); 1067 1225 let rec_ops = ms.record_ops(); 1068 1226 1069 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1070 - let rkey = Rkey::from("r1".to_string()); 1227 + let collection = test_nsid("app.bsky.feed.post"); 1228 + let rkey = test_rkey("r1"); 1071 1229 let cid = test_cid_link(42); 1072 1230 1073 1231 let mut batch = ms.database().batch(); ··· 1097 1255 let (user_id, user_hash) = setup_user(&ms); 1098 1256 let rec_ops = ms.record_ops(); 1099 1257 1100 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1101 - let rkey = Rkey::from("r1".to_string()); 1258 + let collection = test_nsid("app.bsky.feed.post"); 1259 + let rkey = test_rkey("r1"); 1102 1260 let cid = test_cid_link(42); 1103 1261 1104 1262 let mut batch = ms.database().batch(); ··· 1129 1287 let (_user_id, user_hash) = setup_user(&ms); 1130 1288 let rec_ops = ms.record_ops(); 1131 1289 1132 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1133 - let rkey = Rkey::from("r1".to_string()); 1290 + let collection = test_nsid("app.bsky.feed.post"); 1291 + let rkey = test_rkey("r1"); 1134 1292 let cid1 = test_cid_link(1); 1135 1293 let cid2 = test_cid_link(2); 1136 1294 ··· 1158 1316 fn records_survive_reopen() { 1159 1317 let dir = tempfile::TempDir::new().unwrap(); 1160 1318 let user_id; 1161 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1162 - let rkey = Rkey::from("durable".to_string()); 1319 + let collection = test_nsid("app.bsky.feed.post"); 1320 + let rkey = test_rkey("durable"); 1163 1321 let cid = test_cid_link(77); 1164 1322 1165 1323 { ··· 1202 1360 let (user_id, user_hash) = setup_user(&ms); 1203 1361 let rec_ops = ms.record_ops(); 1204 1362 1205 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1363 + let collection = test_nsid("app.bsky.feed.post"); 1206 1364 let rkeys_unordered = ["zebra", "apple", "mango", "banana"]; 1207 1365 1208 1366 let rkeys: Vec<Rkey> = rkeys_unordered 1209 1367 .iter() 1210 - .map(|rk| Rkey::from(rk.to_string())) 1368 + .map(|rk| test_rkey(rk.to_string())) 1211 1369 .collect(); 1212 1370 let cids: Vec<CidLink> = (0..rkeys.len()) 1213 1371 .map(|i| test_cid_link(i as u8 + 1)) ··· 1233 1391 } 1234 1392 1235 1393 #[test] 1236 - fn record_with_empty_rkey() { 1394 + fn record_with_shortest_rkey() { 1237 1395 let (_dir, ms) = open_fresh(); 1238 1396 let (user_id, user_hash) = setup_user(&ms); 1239 1397 let rec_ops = ms.record_ops(); 1240 1398 1241 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1242 - let rkey = Rkey::from(String::new()); 1399 + let collection = test_nsid("app.bsky.feed.post"); 1400 + let rkey = test_rkey("z"); 1243 1401 let cid = test_cid_link(1); 1244 1402 1245 1403 let mut batch = ms.database().batch(); ··· 1258 1416 } 1259 1417 1260 1418 #[test] 1261 - fn record_with_null_bytes_in_rkey() { 1419 + fn record_rkey_extending_another_stays_distinct() { 1262 1420 let (_dir, ms) = open_fresh(); 1263 1421 let (user_id, user_hash) = setup_user(&ms); 1264 1422 let rec_ops = ms.record_ops(); 1265 1423 1266 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1267 - let rkey_with_null = Rkey::from("abc\x00def".to_string()); 1268 - let rkey_plain = Rkey::from("abc".to_string()); 1424 + let collection = test_nsid("app.bsky.feed.post"); 1425 + let rkey_extended = test_rkey("abc.def"); 1426 + let rkey_plain = test_rkey("abc"); 1269 1427 let cid1 = test_cid_link(1); 1270 1428 let cid2 = test_cid_link(2); 1271 1429 ··· 1275 1433 &mut batch, 1276 1434 user_hash, 1277 1435 &[ 1278 - rw(&collection, &rkey_with_null, &cid1), 1436 + rw(&collection, &rkey_extended, &cid1), 1279 1437 rw(&collection, &rkey_plain, &cid2), 1280 1438 ], 1281 1439 ) ··· 1283 1441 batch.commit().unwrap(); 1284 1442 1285 1443 let found1 = rec_ops 1286 - .get_record_cid(user_id, &collection, &rkey_with_null) 1444 + .get_record_cid(user_id, &collection, &rkey_extended) 1287 1445 .unwrap(); 1288 1446 assert_eq!(found1, Some(cid1)); 1289 1447 ··· 1296 1454 .list_records(&lrq(user_id, &collection, None, 100, false, None, None)) 1297 1455 .unwrap(); 1298 1456 assert_eq!(results.len(), 2); 1299 - assert_eq!(results[0].rkey.as_str(), "abc\x00def"); 1457 + assert_eq!(results[0].rkey.as_str(), "abc.def"); 1300 1458 assert_eq!(results[1].rkey.as_str(), "abc"); 1301 1459 } 1302 1460 1303 1461 #[test] 1304 - fn record_with_null_bytes_in_collection() { 1462 + fn record_collection_extending_another_stays_distinct() { 1305 1463 let (_dir, ms) = open_fresh(); 1306 1464 let (user_id, user_hash) = setup_user(&ms); 1307 1465 let rec_ops = ms.record_ops(); 1308 1466 1309 - let coll_normal = Nsid::from("app.bsky.feed.post".to_string()); 1310 - let coll_with_null = Nsid::from("app.bsky.feed.post\x00extra".to_string()); 1311 - let rkey = Rkey::from("r1".to_string()); 1467 + let coll_normal = test_nsid("app.bsky.feed.post"); 1468 + let coll_extended = test_nsid("app.bsky.feed.post.extra"); 1469 + let rkey = test_rkey("r1"); 1312 1470 let cid1 = test_cid_link(1); 1313 1471 let cid2 = test_cid_link(2); 1314 1472 ··· 1319 1477 user_hash, 1320 1478 &[ 1321 1479 rw(&coll_normal, &rkey, &cid1), 1322 - rw(&coll_with_null, &rkey, &cid2), 1480 + rw(&coll_extended, &rkey, &cid2), 1323 1481 ], 1324 1482 ) 1325 1483 .unwrap(); ··· 1331 1489 assert_eq!(found1, Some(cid1)); 1332 1490 1333 1491 let found2 = rec_ops 1334 - .get_record_cid(user_id, &coll_with_null, &rkey) 1492 + .get_record_cid(user_id, &coll_extended, &rkey) 1335 1493 .unwrap(); 1336 1494 assert_eq!(found2, Some(cid2)); 1337 1495 ··· 1340 1498 .unwrap(); 1341 1499 assert_eq!(results_normal.len(), 1); 1342 1500 1343 - let results_null = rec_ops 1344 - .list_records(&lrq(user_id, &coll_with_null, None, 100, false, None, None)) 1501 + let results_extended = rec_ops 1502 + .list_records(&lrq(user_id, &coll_extended, None, 100, false, None, None)) 1345 1503 .unwrap(); 1346 - assert_eq!(results_null.len(), 1); 1504 + assert_eq!(results_extended.len(), 1); 1347 1505 1348 1506 let collections = rec_ops.list_collections(user_id).unwrap(); 1349 1507 assert_eq!(collections.len(), 2); ··· 1355 1513 let (user_id, user_hash) = setup_user(&ms); 1356 1514 let rec_ops = ms.record_ops(); 1357 1515 1358 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1359 - let rkeys: Vec<Rkey> = (0..5).map(|i| Rkey::from(format!("rkey{i:03}"))).collect(); 1516 + let collection = test_nsid("app.bsky.feed.post"); 1517 + let rkeys: Vec<Rkey> = (0..5).map(|i| test_rkey(format!("rkey{i:03}"))).collect(); 1360 1518 let cids: Vec<CidLink> = (0..5).map(|i| test_cid_link(i + 1)).collect(); 1361 1519 1362 1520 let writes: Vec<RecordWrite<'_>> = rkeys ··· 1371 1529 .unwrap(); 1372 1530 batch.commit().unwrap(); 1373 1531 1374 - let cursor = Rkey::from("a".to_string()); 1532 + let cursor = test_rkey("a"); 1375 1533 let results = rec_ops 1376 1534 .list_records(&lrq( 1377 1535 user_id, ··· 1392 1550 let (user_id, user_hash) = setup_user(&ms); 1393 1551 let rec_ops = ms.record_ops(); 1394 1552 1395 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1396 - let rkeys: Vec<Rkey> = (0..10).map(|i| Rkey::from(format!("rkey{i:03}"))).collect(); 1553 + let collection = test_nsid("app.bsky.feed.post"); 1554 + let rkeys: Vec<Rkey> = (0..10).map(|i| test_rkey(format!("rkey{i:03}"))).collect(); 1397 1555 let cids: Vec<CidLink> = (0..10).map(|i| test_cid_link(i + 1)).collect(); 1398 1556 1399 1557 let writes: Vec<RecordWrite<'_>> = rkeys ··· 1408 1566 .unwrap(); 1409 1567 batch.commit().unwrap(); 1410 1568 1411 - let rkey_start = Rkey::from("rkey002".to_string()); 1412 - let rkey_end = Rkey::from("rkey007".to_string()); 1569 + let rkey_start = test_rkey("rkey002"); 1570 + let rkey_end = test_rkey("rkey007"); 1413 1571 let results = rec_ops 1414 1572 .list_records(&lrq( 1415 1573 user_id, ··· 1432 1590 let (user_id, user_hash) = setup_user(&ms); 1433 1591 let rec_ops = ms.record_ops(); 1434 1592 1435 - let collection = Nsid::from("app.bsky.feed.post".to_string()); 1436 - let rkeys: Vec<Rkey> = (0..10).map(|i| Rkey::from(format!("rkey{i:03}"))).collect(); 1593 + let collection = test_nsid("app.bsky.feed.post"); 1594 + let rkeys: Vec<Rkey> = (0..10).map(|i| test_rkey(format!("rkey{i:03}"))).collect(); 1437 1595 let cids: Vec<CidLink> = (0..10).map(|i| test_cid_link(i + 1)).collect(); 1438 1596 1439 1597 let writes: Vec<RecordWrite<'_>> = rkeys ··· 1448 1606 .unwrap(); 1449 1607 batch.commit().unwrap(); 1450 1608 1451 - let cursor = Rkey::from("rkey005".to_string()); 1609 + let cursor = test_rkey("rkey005"); 1452 1610 let results = rec_ops 1453 1611 .list_records(&lrq( 1454 1612 user_id, ··· 1510 1668 k 1511 1669 }; 1512 1670 assert!(key_outside.as_slice() >= upper.as_slice()); 1671 + } 1672 + 1673 + #[test] 1674 + fn referenced_record_cids_reports_a_cid_held_by_another_key() { 1675 + let (_dir, ms) = open_fresh(); 1676 + let (user_id, user_hash) = setup_user(&ms); 1677 + let rec_ops = ms.record_ops(); 1678 + 1679 + let collection = test_nsid("app.bsky.feed.post"); 1680 + let kept = test_rkey("kept"); 1681 + let dropped = test_rkey("dropped"); 1682 + let shared = test_cid_link(9); 1683 + 1684 + let mut batch = ms.database().batch(); 1685 + rec_ops 1686 + .upsert_records( 1687 + &mut batch, 1688 + user_hash, 1689 + &[ 1690 + rw(&collection, &kept, &shared), 1691 + rw(&collection, &dropped, &shared), 1692 + ], 1693 + ) 1694 + .unwrap(); 1695 + batch.commit().unwrap(); 1696 + 1697 + assert_eq!( 1698 + rec_ops 1699 + .referenced_record_cids( 1700 + user_id, 1701 + std::slice::from_ref(&shared), 1702 + &[(&collection, &dropped)] 1703 + ) 1704 + .unwrap(), 1705 + vec![shared.clone()] 1706 + ); 1707 + assert!( 1708 + rec_ops 1709 + .referenced_record_cids( 1710 + user_id, 1711 + &[shared], 1712 + &[(&collection, &dropped), (&collection, &kept)] 1713 + ) 1714 + .unwrap() 1715 + .is_empty() 1716 + ); 1717 + } 1718 + 1719 + #[test] 1720 + fn referenced_record_cids_errors_for_a_repo_with_no_user_hash() { 1721 + let (_dir, ms) = open_fresh(); 1722 + let rec_ops = ms.record_ops(); 1723 + 1724 + assert!(matches!( 1725 + rec_ops.referenced_record_cids(Uuid::new_v4(), &[test_cid_link(1)], &[]), 1726 + Err(MetastoreError::InvalidInput(_)) 1727 + )); 1728 + } 1729 + 1730 + #[test] 1731 + fn two_upserts_of_one_key_in_a_batch_leave_only_the_last_cid_indexed() { 1732 + let (_dir, ms) = open_fresh(); 1733 + let (user_id, user_hash) = setup_user(&ms); 1734 + let rec_ops = ms.record_ops(); 1735 + 1736 + let collection = test_nsid("app.bsky.feed.post"); 1737 + let rkey = test_rkey("twice"); 1738 + let first = test_cid_link(20); 1739 + let second = test_cid_link(21); 1740 + 1741 + let mut batch = ms.database().batch(); 1742 + rec_ops 1743 + .upsert_records( 1744 + &mut batch, 1745 + user_hash, 1746 + &[ 1747 + rw(&collection, &rkey, &first), 1748 + rw(&collection, &rkey, &second), 1749 + ], 1750 + ) 1751 + .unwrap(); 1752 + batch.commit().unwrap(); 1753 + 1754 + assert!( 1755 + rec_ops 1756 + .referenced_record_cids(user_id, &[first], &[]) 1757 + .unwrap() 1758 + .is_empty(), 1759 + "overwritten cid must not stay in the reverse index" 1760 + ); 1761 + assert_eq!( 1762 + rec_ops 1763 + .referenced_record_cids(user_id, std::slice::from_ref(&second), &[]) 1764 + .unwrap(), 1765 + vec![second] 1766 + ); 1767 + } 1768 + 1769 + #[test] 1770 + fn referenced_record_cids_forgets_a_cid_an_upsert_replaced() { 1771 + let (_dir, ms) = open_fresh(); 1772 + let (user_id, user_hash) = setup_user(&ms); 1773 + let rec_ops = ms.record_ops(); 1774 + 1775 + let collection = test_nsid("app.bsky.feed.post"); 1776 + let rkey = test_rkey("churned"); 1777 + let first = test_cid_link(1); 1778 + let second = test_cid_link(2); 1779 + 1780 + let mut batch = ms.database().batch(); 1781 + rec_ops 1782 + .upsert_records(&mut batch, user_hash, &[rw(&collection, &rkey, &first)]) 1783 + .unwrap(); 1784 + batch.commit().unwrap(); 1785 + 1786 + let mut batch = ms.database().batch(); 1787 + rec_ops 1788 + .upsert_records(&mut batch, user_hash, &[rw(&collection, &rkey, &second)]) 1789 + .unwrap(); 1790 + batch.commit().unwrap(); 1791 + 1792 + assert!( 1793 + rec_ops 1794 + .referenced_record_cids(user_id, &[first], &[]) 1795 + .unwrap() 1796 + .is_empty() 1797 + ); 1798 + assert_eq!( 1799 + rec_ops 1800 + .referenced_record_cids(user_id, std::slice::from_ref(&second), &[]) 1801 + .unwrap(), 1802 + vec![second] 1803 + ); 1804 + } 1805 + 1806 + #[test] 1807 + fn referenced_record_cids_forgets_a_deleted_record() { 1808 + let (_dir, ms) = open_fresh(); 1809 + let (user_id, user_hash) = setup_user(&ms); 1810 + let rec_ops = ms.record_ops(); 1811 + 1812 + let collection = test_nsid("app.bsky.feed.post"); 1813 + let rkey = test_rkey("gone"); 1814 + let cid = test_cid_link(3); 1815 + 1816 + let mut batch = ms.database().batch(); 1817 + rec_ops 1818 + .upsert_records(&mut batch, user_hash, &[rw(&collection, &rkey, &cid)]) 1819 + .unwrap(); 1820 + batch.commit().unwrap(); 1821 + 1822 + let mut batch = ms.database().batch(); 1823 + rec_ops 1824 + .delete_records(&mut batch, user_hash, &[rd(&collection, &rkey)]) 1825 + .unwrap(); 1826 + batch.commit().unwrap(); 1827 + 1828 + assert!( 1829 + rec_ops 1830 + .referenced_record_cids(user_id, &[cid], &[]) 1831 + .unwrap() 1832 + .is_empty() 1833 + ); 1834 + } 1835 + 1836 + #[test] 1837 + fn delete_all_records_clears_the_reverse_index() { 1838 + let (_dir, ms) = open_fresh(); 1839 + let (user_id, user_hash) = setup_user(&ms); 1840 + let rec_ops = ms.record_ops(); 1841 + 1842 + let collection = test_nsid("app.bsky.feed.post"); 1843 + let rkey = test_rkey("only"); 1844 + let cid = test_cid_link(4); 1845 + 1846 + let mut batch = ms.database().batch(); 1847 + rec_ops 1848 + .upsert_records(&mut batch, user_hash, &[rw(&collection, &rkey, &cid)]) 1849 + .unwrap(); 1850 + batch.commit().unwrap(); 1851 + 1852 + let mut batch = ms.database().batch(); 1853 + rec_ops.delete_all_records(&mut batch, user_hash).unwrap(); 1854 + batch.commit().unwrap(); 1855 + 1856 + assert!( 1857 + ms.partition(crate::metastore::Partition::RepoData) 1858 + .prefix(record_by_cid_user_prefix(user_hash).as_slice()) 1859 + .next() 1860 + .is_none() 1861 + ); 1862 + assert!( 1863 + rec_ops 1864 + .referenced_record_cids(user_id, &[cid], &[]) 1865 + .unwrap() 1866 + .is_empty() 1867 + ); 1868 + } 1869 + 1870 + fn drop_index_and_marker(ms: &Metastore, repo_data: &Keyspace) { 1871 + let mut batch = ms.database().batch(); 1872 + crate::metastore::scan::delete_all_by_prefix( 1873 + repo_data, 1874 + &mut batch, 1875 + record_by_cid_index_prefix().as_slice(), 1876 + ) 1877 + .unwrap(); 1878 + batch.remove(repo_data, record_by_cid_built_key().as_slice()); 1879 + batch.commit().unwrap(); 1880 + } 1881 + 1882 + #[test] 1883 + fn rebuild_backfills_a_store_written_before_the_index_existed() { 1884 + let (_dir, ms) = open_fresh(); 1885 + let (user_id, user_hash) = setup_user(&ms); 1886 + let rec_ops = ms.record_ops(); 1887 + let repo_data = ms.partition(crate::metastore::Partition::RepoData).clone(); 1888 + 1889 + let collection = test_nsid("app.bsky.feed.post"); 1890 + let rkey = test_rkey("legacy"); 1891 + let cid = test_cid_link(5); 1892 + 1893 + let mut batch = ms.database().batch(); 1894 + rec_ops 1895 + .upsert_records(&mut batch, user_hash, &[rw(&collection, &rkey, &cid)]) 1896 + .unwrap(); 1897 + batch.commit().unwrap(); 1898 + 1899 + drop_index_and_marker(&ms, &repo_data); 1900 + assert!( 1901 + rec_ops 1902 + .referenced_record_cids(user_id, std::slice::from_ref(&cid), &[]) 1903 + .unwrap() 1904 + .is_empty() 1905 + ); 1906 + 1907 + assert_eq!( 1908 + rebuild_record_by_cid_index(ms.database(), &repo_data).unwrap(), 1909 + 1 1910 + ); 1911 + assert_eq!( 1912 + rec_ops 1913 + .referenced_record_cids(user_id, &[cid], &[]) 1914 + .unwrap() 1915 + .len(), 1916 + 1 1917 + ); 1918 + } 1919 + 1920 + #[test] 1921 + fn rebuild_redoes_an_index_left_partial_by_a_crash() { 1922 + let (_dir, ms) = open_fresh(); 1923 + let (user_id, user_hash) = setup_user(&ms); 1924 + let rec_ops = ms.record_ops(); 1925 + let repo_data = ms.partition(crate::metastore::Partition::RepoData).clone(); 1926 + 1927 + let collection = test_nsid("app.bsky.feed.post"); 1928 + let indexed = test_rkey("indexed"); 1929 + let missed = test_rkey("missed"); 1930 + let indexed_cid = test_cid_link(7); 1931 + let missed_cid = test_cid_link(8); 1932 + 1933 + let mut batch = ms.database().batch(); 1934 + rec_ops 1935 + .upsert_records( 1936 + &mut batch, 1937 + user_hash, 1938 + &[ 1939 + rw(&collection, &indexed, &indexed_cid), 1940 + rw(&collection, &missed, &missed_cid), 1941 + ], 1942 + ) 1943 + .unwrap(); 1944 + batch.commit().unwrap(); 1945 + 1946 + drop_index_and_marker(&ms, &repo_data); 1947 + let mut batch = ms.database().batch(); 1948 + batch.insert( 1949 + &repo_data, 1950 + record_by_cid_key( 1951 + user_hash, 1952 + &cid_link_to_bytes(&indexed_cid).unwrap(), 1953 + &collection, 1954 + &indexed, 1955 + ) 1956 + .as_slice(), 1957 + [], 1958 + ); 1959 + batch.commit().unwrap(); 1960 + 1961 + assert_eq!( 1962 + rebuild_record_by_cid_index(ms.database(), &repo_data).unwrap(), 1963 + 2, 1964 + "a partial index without the marker must be rebuilt from scratch" 1965 + ); 1966 + assert_eq!( 1967 + rec_ops 1968 + .referenced_record_cids(user_id, &[indexed_cid, missed_cid], &[]) 1969 + .unwrap() 1970 + .len(), 1971 + 2 1972 + ); 1973 + } 1974 + 1975 + #[test] 1976 + fn rebuild_drops_a_stale_entry_an_older_binary_left_behind() { 1977 + let (_dir, ms) = open_fresh(); 1978 + let (user_id, user_hash) = setup_user(&ms); 1979 + let rec_ops = ms.record_ops(); 1980 + let repo_data = ms.partition(crate::metastore::Partition::RepoData).clone(); 1981 + 1982 + let collection = test_nsid("app.bsky.feed.post"); 1983 + let rkey = test_rkey("churned"); 1984 + let stale_cid = test_cid_link(10); 1985 + let live_cid = test_cid_link(11); 1986 + 1987 + let mut batch = ms.database().batch(); 1988 + rec_ops 1989 + .upsert_records(&mut batch, user_hash, &[rw(&collection, &rkey, &live_cid)]) 1990 + .unwrap(); 1991 + batch.insert( 1992 + &repo_data, 1993 + record_by_cid_key( 1994 + user_hash, 1995 + &cid_link_to_bytes(&stale_cid).unwrap(), 1996 + &collection, 1997 + &rkey, 1998 + ) 1999 + .as_slice(), 2000 + [], 2001 + ); 2002 + batch.remove(&repo_data, record_by_cid_built_key().as_slice()); 2003 + batch.commit().unwrap(); 2004 + 2005 + assert_eq!( 2006 + rebuild_record_by_cid_index(ms.database(), &repo_data).unwrap(), 2007 + 1 2008 + ); 2009 + assert!( 2010 + rec_ops 2011 + .referenced_record_cids(user_id, &[stale_cid], &[]) 2012 + .unwrap() 2013 + .is_empty() 2014 + ); 2015 + assert_eq!( 2016 + rec_ops 2017 + .referenced_record_cids(user_id, &[live_cid], &[]) 2018 + .unwrap() 2019 + .len(), 2020 + 1 2021 + ); 2022 + } 2023 + 2024 + #[test] 2025 + fn rebuild_commits_in_chunks_past_the_batch_threshold() { 2026 + let (_dir, ms) = open_fresh(); 2027 + let (user_id, user_hash) = setup_user(&ms); 2028 + let rec_ops = ms.record_ops(); 2029 + let repo_data = ms.partition(crate::metastore::Partition::RepoData).clone(); 2030 + 2031 + let collection = test_nsid("app.bsky.feed.post"); 2032 + let record_count = 10_001usize; 2033 + let rkeys: Vec<Rkey> = (0..record_count) 2034 + .map(|i| test_rkey(format!("chunk{i:05}"))) 2035 + .collect(); 2036 + let cids: Vec<CidLink> = (0..record_count) 2037 + .map(|i| test_cid_link((i % 251) as u8)) 2038 + .collect(); 2039 + 2040 + let writes: Vec<RecordWrite<'_>> = rkeys 2041 + .iter() 2042 + .zip(cids.iter()) 2043 + .map(|(rkey, cid)| rw(&collection, rkey, cid)) 2044 + .collect(); 2045 + let mut batch = ms.database().batch(); 2046 + rec_ops 2047 + .upsert_records(&mut batch, user_hash, &writes) 2048 + .unwrap(); 2049 + batch.commit().unwrap(); 2050 + 2051 + drop_index_and_marker(&ms, &repo_data); 2052 + 2053 + assert_eq!( 2054 + rebuild_record_by_cid_index(ms.database(), &repo_data).unwrap(), 2055 + record_count 2056 + ); 2057 + assert_eq!( 2058 + rec_ops 2059 + .referenced_record_cids(user_id, std::slice::from_ref(&cids[0]), &[]) 2060 + .unwrap() 2061 + .len(), 2062 + 1 2063 + ); 2064 + assert_eq!( 2065 + rec_ops 2066 + .referenced_record_cids(user_id, std::slice::from_ref(cids.last().unwrap()), &[]) 2067 + .unwrap() 2068 + .len(), 2069 + 1 2070 + ); 2071 + } 2072 + 2073 + #[test] 2074 + fn rebuild_discards_a_large_stale_index_past_the_batch_threshold() { 2075 + let (_dir, ms) = open_fresh(); 2076 + let (user_id, user_hash) = setup_user(&ms); 2077 + let rec_ops = ms.record_ops(); 2078 + let repo_data = ms.partition(crate::metastore::Partition::RepoData).clone(); 2079 + 2080 + let collection = test_nsid("app.bsky.feed.post"); 2081 + let live_rkey = test_rkey("live"); 2082 + let live_cid = test_cid_link(1); 2083 + 2084 + let mut batch = ms.database().batch(); 2085 + rec_ops 2086 + .upsert_records( 2087 + &mut batch, 2088 + user_hash, 2089 + &[rw(&collection, &live_rkey, &live_cid)], 2090 + ) 2091 + .unwrap(); 2092 + batch.commit().unwrap(); 2093 + 2094 + let stale_count = 10_001usize; 2095 + let stale: Vec<(Rkey, Vec<u8>)> = (0..stale_count) 2096 + .map(|i| { 2097 + let rkey = test_rkey(format!("stale{i:05}")); 2098 + let cid = cid_link_to_bytes(&test_cid_link(2)).unwrap(); 2099 + let cid = cid 2100 + .iter() 2101 + .copied() 2102 + .chain((i as u32).to_be_bytes()) 2103 + .collect(); 2104 + (rkey, cid) 2105 + }) 2106 + .collect(); 2107 + let mut batch = ms.database().batch(); 2108 + stale.iter().for_each(|(rkey, cid)| { 2109 + batch.insert( 2110 + &repo_data, 2111 + record_by_cid_key(user_hash, cid, &collection, rkey).as_slice(), 2112 + [], 2113 + ); 2114 + }); 2115 + batch.remove(&repo_data, record_by_cid_built_key().as_slice()); 2116 + batch.commit().unwrap(); 2117 + 2118 + assert_eq!( 2119 + rebuild_record_by_cid_index(ms.database(), &repo_data).unwrap(), 2120 + 1 2121 + ); 2122 + assert_eq!( 2123 + repo_data 2124 + .prefix(record_by_cid_index_prefix().as_slice()) 2125 + .count(), 2126 + 1, 2127 + "rebuild must leave exactly one index entry, so no stale key survived the chunked delete" 2128 + ); 2129 + assert_eq!( 2130 + rec_ops 2131 + .referenced_record_cids(user_id, &[live_cid], &[]) 2132 + .unwrap() 2133 + .len(), 2134 + 1 2135 + ); 2136 + } 2137 + 2138 + #[test] 2139 + fn rebuild_is_a_noop_once_the_marker_exists() { 2140 + let (_dir, ms) = open_fresh(); 2141 + let (_user_id, user_hash) = setup_user(&ms); 2142 + let rec_ops = ms.record_ops(); 2143 + let repo_data = ms.partition(crate::metastore::Partition::RepoData).clone(); 2144 + 2145 + let collection = test_nsid("app.bsky.feed.post"); 2146 + let rkey = test_rkey("present"); 2147 + let cid = test_cid_link(6); 2148 + 2149 + let mut batch = ms.database().batch(); 2150 + rec_ops 2151 + .upsert_records(&mut batch, user_hash, &[rw(&collection, &rkey, &cid)]) 2152 + .unwrap(); 2153 + batch.commit().unwrap(); 2154 + 2155 + assert_eq!( 2156 + rebuild_record_by_cid_index(ms.database(), &repo_data).unwrap(), 2157 + 0 2158 + ); 1513 2159 } 1514 2160 }
+6 -1
crates/tranquil-store/src/metastore/repo_ops.rs
··· 6 6 use super::MetastoreError; 7 7 use super::encoding::KeyReader; 8 8 use super::keys::{KeyTag, UserHash}; 9 - use super::records::record_user_prefix; 9 + use super::records::{record_by_cid_user_prefix, record_user_prefix}; 10 10 use super::repo_meta::{ 11 11 RepoMetaValue, RepoStatus, handle_key, repo_meta_key, repo_meta_prefix, stage_repo_meta_removal, 12 12 }; ··· 620 620 ) -> Result<(), MetastoreError> { 621 621 stage_repo_meta_removal(batch, repo_data, user_hash, handle); 622 622 delete_all_by_prefix(repo_data, batch, record_user_prefix(user_hash).as_slice())?; 623 + delete_all_by_prefix( 624 + repo_data, 625 + batch, 626 + record_by_cid_user_prefix(user_hash).as_slice(), 627 + )?; 623 628 delete_all_by_prefix( 624 629 repo_data, 625 630 batch,
+1
migrations/20260720_records_cid_index.sql
··· 1 + CREATE INDEX IF NOT EXISTS idx_records_repo_record_cid ON records(repo_id, record_cid);