Our Personal Data Server from scratch!
0

Configure Feed

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

1use std::sync::Arc; 2 3use fjall::{Database, Keyspace}; 4use smallvec::SmallVec; 5use uuid::Uuid; 6 7use super::MetastoreError; 8use super::backlink_ops::BacklinkOps; 9use super::backlinks::path_to_discriminant; 10use super::encoding::KeyBuilder; 11use super::event_ops::EventOps; 12use super::keys::{KeyTag, UserHash}; 13use super::record_ops::{RecordDelete, RecordOps, RecordWrite}; 14use super::recovery::{ 15 BacklinkMutation, CommitMutationSet, RecordMutationDelete, RecordMutationUpsert, 16}; 17use super::repo_meta::{RepoMetaValue, repo_meta_key, repo_meta_prefix}; 18use super::repo_ops::{RepoOps, bytes_to_cid_link, cid_link_to_bytes}; 19use super::user_block_ops::UserBlockOps; 20use super::user_blocks::user_block_user_prefix; 21use super::user_hash::UserHashMap; 22use crate::blockstore::TranquilBlockStore; 23use crate::eventlog::EventLogBridge; 24use crate::io::StorageIO; 25 26use tranquil_db_traits::{ 27 ApplyCommitError, ApplyCommitInput, ApplyCommitResult, BrokenGenesisCommit, ImportBlock, 28 ImportRecord, ImportRepoError, RepoEventType, SequenceNumber, UserNeedingRecordBlobsBackfill, 29 UserWithoutBlocks, 30}; 31use tranquil_types::{AtUri, CidLink, Did}; 32 33use serde::{Deserialize, Serialize}; 34 35pub(crate) const RECORD_BLOBS_SCHEMA_VERSION: u8 = 1; 36 37#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] 38pub(crate) struct RecordBlobsValue { 39 pub(crate) blob_cid_bytes: Vec<Vec<u8>>, 40} 41 42impl RecordBlobsValue { 43 pub(crate) fn serialize(&self) -> Vec<u8> { 44 let payload = 45 postcard::to_allocvec(self).expect("RecordBlobsValue serialization cannot fail"); 46 let mut buf = Vec::with_capacity(1 + payload.len()); 47 buf.push(RECORD_BLOBS_SCHEMA_VERSION); 48 buf.extend_from_slice(&payload); 49 buf 50 } 51 52 pub(crate) fn deserialize(bytes: &[u8]) -> Option<Self> { 53 let (&version, payload) = bytes.split_first()?; 54 match version { 55 RECORD_BLOBS_SCHEMA_VERSION => postcard::from_bytes(payload).ok(), 56 _ => None, 57 } 58 } 59} 60 61fn record_blobs_key(user_hash: UserHash, uri: &AtUri) -> SmallVec<[u8; 128]> { 62 KeyBuilder::new() 63 .tag(KeyTag::RECORD_BLOBS) 64 .u64(user_hash.raw()) 65 .string(uri.as_str()) 66 .build() 67} 68 69pub(crate) fn record_blobs_user_prefix(user_hash: UserHash) -> SmallVec<[u8; 128]> { 70 KeyBuilder::new() 71 .tag(KeyTag::RECORD_BLOBS) 72 .u64(user_hash.raw()) 73 .build() 74} 75 76pub struct CommitOps<S: StorageIO> { 77 db: Database, 78 repo_data: Keyspace, 79 user_hashes: Arc<UserHashMap>, 80 repo_ops: RepoOps, 81 record_ops: RecordOps, 82 user_block_ops: UserBlockOps, 83 backlink_ops: BacklinkOps, 84 event_ops: EventOps<S>, 85 blockstore: Option<TranquilBlockStore>, 86} 87 88impl<S: StorageIO> CommitOps<S> { 89 pub fn new( 90 db: Database, 91 repo_data: Keyspace, 92 indexes: Keyspace, 93 user_hashes: Arc<UserHashMap>, 94 bridge: Arc<EventLogBridge<S>>, 95 ) -> Self { 96 let repo_ops = RepoOps::new(repo_data.clone(), Arc::clone(&user_hashes)); 97 let record_ops = RecordOps::new(repo_data.clone(), Arc::clone(&user_hashes)); 98 let user_block_ops = UserBlockOps::new(repo_data.clone(), Arc::clone(&user_hashes)); 99 let backlink_ops = BacklinkOps::new(indexes, Arc::clone(&user_hashes)); 100 let event_ops = EventOps::new(db.clone(), repo_data.clone(), bridge); 101 Self { 102 db, 103 repo_data, 104 user_hashes, 105 repo_ops, 106 record_ops, 107 user_block_ops, 108 backlink_ops, 109 event_ops, 110 blockstore: None, 111 } 112 } 113 114 pub fn with_blockstore(mut self, blockstore: TranquilBlockStore) -> Self { 115 self.blockstore = Some(blockstore); 116 self 117 } 118 119 pub fn apply_commit( 120 &self, 121 input: ApplyCommitInput, 122 ) -> Result<ApplyCommitResult, ApplyCommitError> { 123 let user_hash = self 124 .user_hashes 125 .get(&input.user_id) 126 .ok_or(ApplyCommitError::RepoNotFound)?; 127 128 let key = repo_meta_key(user_hash); 129 let meta = self 130 .repo_data 131 .get(key.as_slice()) 132 .map_err(|e| ApplyCommitError::Database(e.to_string()))? 133 .and_then(|raw| RepoMetaValue::deserialize(&raw)) 134 .ok_or(ApplyCommitError::RepoNotFound)?; 135 136 if let Some(expected) = &input.expected_root_cid { 137 let current = bytes_to_cid_link(&meta.repo_root_cid) 138 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 139 if current != *expected { 140 return Err(ApplyCommitError::ConcurrentModification); 141 } 142 } 143 144 let new_cid_bytes = cid_link_to_bytes(&input.new_root_cid) 145 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 146 147 let is_active = meta.status.is_active(); 148 149 let updated_meta = RepoMetaValue { 150 repo_root_cid: new_cid_bytes.clone(), 151 repo_rev: input.new_rev.clone(), 152 ..meta 153 }; 154 155 let mut batch = self.db.batch(); 156 157 self.repo_ops 158 .write_repo_meta(&mut batch, user_hash, &updated_meta); 159 160 let upserts: Vec<RecordWrite<'_>> = input 161 .record_upserts 162 .iter() 163 .map(|u| RecordWrite { 164 collection: &u.collection, 165 rkey: &u.rkey, 166 cid: &u.cid, 167 }) 168 .collect(); 169 170 self.record_ops 171 .upsert_records(&mut batch, user_hash, &upserts) 172 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 173 174 let deletes: Vec<RecordDelete<'_>> = input 175 .record_deletes 176 .iter() 177 .map(|d| RecordDelete { 178 collection: &d.collection, 179 rkey: &d.rkey, 180 }) 181 .collect(); 182 183 self.record_ops 184 .delete_records(&mut batch, user_hash, &deletes); 185 186 self.user_block_ops 187 .insert_user_blocks(&mut batch, user_hash, &input.new_block_cids, &input.new_rev) 188 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 189 190 self.user_block_ops 191 .delete_user_blocks_by_cid(&mut batch, user_hash, &input.obsolete_block_cids) 192 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 193 194 input.backlinks_to_remove.iter().try_for_each(|uri| { 195 self.backlink_ops 196 .remove_backlinks_by_uri(&mut batch, user_hash, uri) 197 .map_err(|e| ApplyCommitError::Database(e.to_string())) 198 })?; 199 200 self.backlink_ops 201 .add_backlinks(&mut batch, user_hash, &input.backlinks_to_add) 202 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 203 204 let mutation_set = CommitMutationSet { 205 new_root_cid: new_cid_bytes.clone(), 206 new_rev: input.new_rev.clone(), 207 record_upserts: input 208 .record_upserts 209 .iter() 210 .map(|u| { 211 let cid_bytes = cid_link_to_bytes(&u.cid) 212 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 213 Ok(RecordMutationUpsert { 214 collection: u.collection.as_str().to_owned(), 215 rkey: u.rkey.as_str().to_owned(), 216 cid_bytes, 217 }) 218 }) 219 .collect::<Result<Vec<_>, ApplyCommitError>>()?, 220 record_deletes: input 221 .record_deletes 222 .iter() 223 .map(|d| RecordMutationDelete { 224 collection: d.collection.as_str().to_owned(), 225 rkey: d.rkey.as_str().to_owned(), 226 }) 227 .collect(), 228 block_inserts: input.new_block_cids.clone(), 229 block_deletes: input.obsolete_block_cids.clone(), 230 backlink_adds: input 231 .backlinks_to_add 232 .iter() 233 .map(|bl| BacklinkMutation { 234 uri: bl.uri.as_str().to_owned(), 235 path: path_to_discriminant(bl.path), 236 link_to: bl.link_to.clone(), 237 }) 238 .collect(), 239 backlink_remove_uris: input 240 .backlinks_to_remove 241 .iter() 242 .map(|uri| uri.as_str().to_owned()) 243 .collect(), 244 }; 245 let mutation_set_bytes = mutation_set 246 .serialize() 247 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 248 249 let (seq, deferred) = self 250 .event_ops 251 .append_commit_event_into_batch( 252 &mut batch, 253 &input.commit_event, 254 Some(&mutation_set_bytes), 255 ) 256 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 257 258 batch 259 .commit() 260 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 261 262 self.event_ops.complete_broadcast(deferred); 263 264 Ok(ApplyCommitResult { 265 seq: seq.as_i64(), 266 is_account_active: is_active, 267 }) 268 } 269 270 pub fn import_repo_data( 271 &self, 272 user_id: Uuid, 273 blocks: &[ImportBlock], 274 records: &[ImportRecord], 275 expected_root_cid: Option<&CidLink>, 276 ) -> Result<(), ImportRepoError> { 277 let user_hash = self 278 .user_hashes 279 .get(&user_id) 280 .ok_or(ImportRepoError::RepoNotFound)?; 281 282 let key = repo_meta_key(user_hash); 283 let meta = self 284 .repo_data 285 .get(key.as_slice()) 286 .map_err(|e| ImportRepoError::Database(e.to_string()))? 287 .and_then(|raw| RepoMetaValue::deserialize(&raw)) 288 .ok_or(ImportRepoError::RepoNotFound)?; 289 290 if let Some(expected) = expected_root_cid { 291 let current = bytes_to_cid_link(&meta.repo_root_cid) 292 .map_err(|e| ImportRepoError::Database(e.to_string()))?; 293 if current != *expected { 294 return Err(ImportRepoError::ConcurrentModification); 295 } 296 } 297 298 if let Some(bs) = &self.blockstore 299 && !blocks.is_empty() 300 { 301 let block_pairs: Vec<([u8; 36], Vec<u8>)> = blocks 302 .iter() 303 .map(|b| { 304 let cid: [u8; 36] = b.cid_bytes.as_slice().try_into().map_err(|_| { 305 ImportRepoError::Database(format!( 306 "block CID has invalid length: {} (expected 36)", 307 b.cid_bytes.len() 308 )) 309 })?; 310 Ok((cid, b.data.clone())) 311 }) 312 .collect::<Result<Vec<_>, ImportRepoError>>()?; 313 bs.put_blocks_blocking(block_pairs) 314 .map_err(|e| ImportRepoError::Database(e.to_string()))?; 315 } 316 317 let mut batch = self.db.batch(); 318 319 let upserts: Vec<RecordWrite<'_>> = records 320 .iter() 321 .map(|r| RecordWrite { 322 collection: &r.collection, 323 rkey: &r.rkey, 324 cid: &r.record_cid, 325 }) 326 .collect(); 327 328 self.record_ops 329 .upsert_records(&mut batch, user_hash, &upserts) 330 .map_err(|e| ImportRepoError::Database(e.to_string()))?; 331 332 batch 333 .commit() 334 .map_err(|e| ImportRepoError::Database(e.to_string())) 335 } 336 337 pub fn insert_record_blobs( 338 &self, 339 repo_id: Uuid, 340 record_uris: &[AtUri], 341 blob_cids: &[CidLink], 342 ) -> Result<(), MetastoreError> { 343 let user_hash = self 344 .user_hashes 345 .get(&repo_id) 346 .ok_or(MetastoreError::InvalidInput("unknown user_id"))?; 347 348 let blob_bytes: Vec<Vec<u8>> = blob_cids 349 .iter() 350 .map(cid_link_to_bytes) 351 .collect::<Result<_, _>>()?; 352 353 let serialized = RecordBlobsValue { 354 blob_cid_bytes: blob_bytes, 355 } 356 .serialize(); 357 358 let mut batch = self.db.batch(); 359 record_uris.iter().for_each(|uri| { 360 let key = record_blobs_key(user_hash, uri); 361 batch.insert(&self.repo_data, key.as_slice(), serialized.as_slice()); 362 }); 363 batch.commit().map_err(MetastoreError::Fjall) 364 } 365 366 pub fn get_users_needing_record_blobs_backfill( 367 &self, 368 limit: i64, 369 ) -> Result<Vec<UserNeedingRecordBlobsBackfill>, MetastoreError> { 370 let limit_usize = usize::try_from(limit).unwrap_or(usize::MAX); 371 372 self.scan_users_missing_prefix( 373 record_blobs_user_prefix, 374 |meta, user_id| { 375 let did = meta 376 .did 377 .map(Did::from) 378 .ok_or(MetastoreError::CorruptData("repo_meta missing did field"))?; 379 Ok(UserNeedingRecordBlobsBackfill { user_id, did }) 380 }, 381 limit_usize, 382 ) 383 } 384 385 pub fn get_broken_genesis_commits(&self) -> Result<Vec<BrokenGenesisCommit>, MetastoreError> { 386 const PAGE_SIZE: usize = 4096; 387 self.collect_broken_genesis_page(SequenceNumber::ZERO, Vec::new(), PAGE_SIZE) 388 } 389 390 fn collect_broken_genesis_page( 391 &self, 392 cursor: SequenceNumber, 393 acc: Vec<BrokenGenesisCommit>, 394 page_size: usize, 395 ) -> Result<Vec<BrokenGenesisCommit>, MetastoreError> { 396 let limit = i64::try_from(page_size).unwrap_or(i64::MAX); 397 let events = self 398 .event_ops 399 .get_events_since_seq(cursor, Some(limit)) 400 .map_err(|_| MetastoreError::CorruptData("failed to read events"))?; 401 402 let page_len = events.len(); 403 let page_high_seq = events.last().map(|e| e.seq).unwrap_or(cursor); 404 405 let results = events.into_iter().fold(acc, |mut results, e| { 406 if e.event_type == RepoEventType::Commit 407 && e.prev_cid.is_none() 408 && e.commit_cid.is_none() 409 { 410 results.push(BrokenGenesisCommit { 411 seq: e.seq, 412 did: e.did, 413 commit_cid: e.commit_cid, 414 }); 415 } 416 results 417 }); 418 419 match page_len < page_size { 420 true => Ok(results), 421 false => self.collect_broken_genesis_page(page_high_seq, results, page_size), 422 } 423 } 424 425 pub fn get_users_without_blocks(&self) -> Result<Vec<UserWithoutBlocks>, MetastoreError> { 426 const MAX_RESULTS: usize = 10_000; 427 428 self.scan_users_missing_prefix( 429 user_block_user_prefix, 430 |meta, user_id| { 431 let root_cid = bytes_to_cid_link(&meta.repo_root_cid)?; 432 Ok(UserWithoutBlocks { 433 user_id, 434 repo_root_cid: root_cid, 435 repo_rev: match meta.repo_rev.is_empty() { 436 true => None, 437 false => Some(meta.repo_rev), 438 }, 439 }) 440 }, 441 MAX_RESULTS, 442 ) 443 } 444 445 fn scan_users_missing_prefix<T, F, P>( 446 &self, 447 make_prefix: P, 448 build_result: F, 449 limit: usize, 450 ) -> Result<Vec<T>, MetastoreError> 451 where 452 F: Fn(RepoMetaValue, Uuid) -> Result<T, MetastoreError>, 453 P: Fn(UserHash) -> SmallVec<[u8; 128]>, 454 { 455 let prefix = repo_meta_prefix(); 456 457 self.repo_data 458 .prefix(prefix.as_slice()) 459 .filter_map(|guard| { 460 let (key_bytes, val_bytes) = match guard.into_inner() { 461 Ok(pair) => pair, 462 Err(e) => return Some(Err(MetastoreError::Fjall(e))), 463 }; 464 465 let user_hash = match parse_user_hash_from_key(&key_bytes) { 466 Some(h) => h, 467 None => return Some(Err(MetastoreError::CorruptData("invalid repo_meta key"))), 468 }; 469 470 let check_prefix = make_prefix(user_hash); 471 let has_entries = match self.repo_data.prefix(check_prefix.as_slice()).next() { 472 Some(guard) => match guard.into_inner() { 473 Ok(_) => true, 474 Err(e) => return Some(Err(MetastoreError::Fjall(e))), 475 }, 476 None => false, 477 }; 478 479 match has_entries { 480 true => None, 481 false => { 482 let meta = match RepoMetaValue::deserialize(&val_bytes) { 483 Some(v) => v, 484 None => { 485 return Some(Err(MetastoreError::CorruptData( 486 "invalid repo_meta value", 487 ))); 488 } 489 }; 490 let user_id = match self.user_hashes.get_uuid(&user_hash) { 491 Some(id) => id, 492 None => { 493 return Some(Err(MetastoreError::CorruptData( 494 "user_hash has no reverse mapping", 495 ))); 496 } 497 }; 498 Some(build_result(meta, user_id)) 499 } 500 } 501 }) 502 .take(limit) 503 .collect() 504 } 505} 506 507fn parse_user_hash_from_key(key_bytes: &[u8]) -> Option<UserHash> { 508 use super::encoding::KeyReader; 509 let mut reader = KeyReader::new(key_bytes); 510 let _tag = reader.tag()?; 511 let hash = reader.u64()?; 512 Some(UserHash::from_raw(hash)) 513} 514 515#[cfg(test)] 516mod tests { 517 use super::*; 518 use crate::eventlog::{EventLog, EventLogConfig}; 519 use crate::io::RealIO; 520 use crate::metastore::{Metastore, MetastoreConfig}; 521 use tranquil_db_traits::CommitEventData; 522 use tranquil_types::{Handle, Nsid, Rkey}; 523 524 struct TestHarness { 525 _metastore_dir: tempfile::TempDir, 526 _eventlog_dir: tempfile::TempDir, 527 metastore: Metastore, 528 bridge: Arc<EventLogBridge<RealIO>>, 529 } 530 531 fn setup() -> TestHarness { 532 let metastore_dir = tempfile::TempDir::new().unwrap(); 533 let eventlog_dir = tempfile::TempDir::new().unwrap(); 534 let segments_dir = eventlog_dir.path().join("segments"); 535 std::fs::create_dir_all(&segments_dir).unwrap(); 536 537 let metastore = Metastore::open( 538 metastore_dir.path(), 539 MetastoreConfig { 540 cache_size_bytes: 64 * 1024 * 1024, 541 }, 542 ) 543 .unwrap(); 544 545 let event_log = EventLog::open( 546 EventLogConfig { 547 segments_dir, 548 ..EventLogConfig::default() 549 }, 550 RealIO::new(), 551 ) 552 .unwrap(); 553 554 let bridge = Arc::new(EventLogBridge::new(Arc::new(event_log))); 555 556 TestHarness { 557 _metastore_dir: metastore_dir, 558 _eventlog_dir: eventlog_dir, 559 metastore, 560 bridge, 561 } 562 } 563 564 fn test_cid_link(seed: u8) -> CidLink { 565 let digest: [u8; 32] = std::array::from_fn(|i| seed.wrapping_add(i as u8)); 566 let mh = multihash::Multihash::<64>::wrap(0x12, &digest).unwrap(); 567 let c = cid::Cid::new_v1(0x71, mh); 568 CidLink::from_cid(&c) 569 } 570 571 fn test_did(name: &str) -> Did { 572 Did::from(format!("did:plc:{name}")) 573 } 574 575 fn test_handle(name: &str) -> Handle { 576 Handle::from(format!("{name}.test.invalid")) 577 } 578 579 fn make_commit_ops(h: &TestHarness) -> CommitOps<RealIO> { 580 use crate::metastore::partitions::Partition; 581 CommitOps::new( 582 h.metastore.database().clone(), 583 h.metastore.partition(Partition::RepoData).clone(), 584 h.metastore.partition(Partition::Indexes).clone(), 585 Arc::clone(h.metastore.user_hashes()), 586 Arc::clone(&h.bridge), 587 ) 588 } 589 590 fn create_test_repo(h: &TestHarness, name: &str, seed: u8) -> (Uuid, Did, CidLink) { 591 let user_id = Uuid::new_v4(); 592 let did = test_did(name); 593 let handle = test_handle(name); 594 let cid = test_cid_link(seed); 595 h.metastore 596 .repo_ops() 597 .create_repo(h.metastore.database(), user_id, &did, &handle, &cid, "rev0") 598 .unwrap(); 599 (user_id, did, cid) 600 } 601 602 #[test] 603 fn apply_commit_updates_records_and_meta() { 604 let h = setup(); 605 let ops = make_commit_ops(&h); 606 let (user_id, did, root_cid) = create_test_repo(&h, "alice", 1); 607 608 let new_root = test_cid_link(2); 609 let record_cid = test_cid_link(3); 610 let collection = Nsid::from("app.bsky.feed.post".to_string()); 611 let rkey = Rkey::from("3k2abc".to_string()); 612 613 let input = ApplyCommitInput { 614 user_id, 615 did: did.clone(), 616 expected_root_cid: Some(root_cid.clone()), 617 new_root_cid: new_root.clone(), 618 new_rev: "rev1".to_string(), 619 new_block_cids: vec![vec![0x01, 0x02]], 620 obsolete_block_cids: vec![], 621 record_upserts: vec![tranquil_db_traits::RecordUpsert { 622 collection: collection.clone(), 623 rkey: rkey.clone(), 624 cid: record_cid.clone(), 625 }], 626 record_deletes: vec![], 627 backlinks_to_add: vec![], 628 backlinks_to_remove: vec![], 629 commit_event: CommitEventData { 630 did: did.clone(), 631 event_type: RepoEventType::Commit, 632 commit_cid: Some(new_root.clone()), 633 prev_cid: Some(root_cid.clone()), 634 ops: None, 635 blobs: None, 636 blocks_cids: None, 637 prev_data_cid: None, 638 rev: Some("rev1".to_string()), 639 }, 640 }; 641 642 let result = ops.apply_commit(input).unwrap(); 643 assert!(result.seq > 0); 644 assert!(result.is_account_active); 645 646 let repo = h.metastore.repo_ops().get_repo(user_id).unwrap().unwrap(); 647 assert_eq!(repo.repo_root_cid, new_root); 648 assert_eq!(repo.repo_rev.as_deref(), Some("rev1")); 649 650 let found_cid = h 651 .metastore 652 .record_ops() 653 .get_record_cid(user_id, &collection, &rkey) 654 .unwrap() 655 .unwrap(); 656 assert_eq!(found_cid, record_cid); 657 } 658 659 #[test] 660 fn apply_commit_cas_rejects_stale_root() { 661 let h = setup(); 662 let ops = make_commit_ops(&h); 663 let (user_id, did, _root_cid) = create_test_repo(&h, "bob", 10); 664 665 let stale_root = test_cid_link(99); 666 let new_root = test_cid_link(11); 667 668 let input = ApplyCommitInput { 669 user_id, 670 did, 671 expected_root_cid: Some(stale_root), 672 new_root_cid: new_root, 673 new_rev: "rev1".to_string(), 674 new_block_cids: vec![], 675 obsolete_block_cids: vec![], 676 record_upserts: vec![], 677 record_deletes: vec![], 678 backlinks_to_add: vec![], 679 backlinks_to_remove: vec![], 680 commit_event: CommitEventData { 681 did: test_did("bob"), 682 event_type: RepoEventType::Commit, 683 commit_cid: None, 684 prev_cid: None, 685 ops: None, 686 blobs: None, 687 blocks_cids: None, 688 prev_data_cid: None, 689 rev: Some("rev1".to_string()), 690 }, 691 }; 692 693 let result = ops.apply_commit(input); 694 assert_eq!( 695 result.unwrap_err(), 696 ApplyCommitError::ConcurrentModification 697 ); 698 } 699 700 #[test] 701 fn apply_commit_returns_repo_not_found_for_unknown_user() { 702 let h = setup(); 703 let ops = make_commit_ops(&h); 704 705 let input = ApplyCommitInput { 706 user_id: Uuid::new_v4(), 707 did: test_did("nobody"), 708 expected_root_cid: None, 709 new_root_cid: test_cid_link(1), 710 new_rev: "rev1".to_string(), 711 new_block_cids: vec![], 712 obsolete_block_cids: vec![], 713 record_upserts: vec![], 714 record_deletes: vec![], 715 backlinks_to_add: vec![], 716 backlinks_to_remove: vec![], 717 commit_event: CommitEventData { 718 did: test_did("nobody"), 719 event_type: RepoEventType::Commit, 720 commit_cid: None, 721 prev_cid: None, 722 ops: None, 723 blobs: None, 724 blocks_cids: None, 725 prev_data_cid: None, 726 rev: None, 727 }, 728 }; 729 730 assert_eq!( 731 ops.apply_commit(input).unwrap_err(), 732 ApplyCommitError::RepoNotFound 733 ); 734 } 735 736 #[test] 737 fn apply_commit_record_deletes() { 738 let h = setup(); 739 let ops = make_commit_ops(&h); 740 let (user_id, did, root_cid) = create_test_repo(&h, "carol", 20); 741 742 let mid_root = test_cid_link(21); 743 let record_cid = test_cid_link(22); 744 let collection = Nsid::from("app.bsky.feed.post".to_string()); 745 let rkey = Rkey::from("3k2del".to_string()); 746 747 let insert_input = ApplyCommitInput { 748 user_id, 749 did: did.clone(), 750 expected_root_cid: Some(root_cid.clone()), 751 new_root_cid: mid_root.clone(), 752 new_rev: "rev1".to_string(), 753 new_block_cids: vec![], 754 obsolete_block_cids: vec![], 755 record_upserts: vec![tranquil_db_traits::RecordUpsert { 756 collection: collection.clone(), 757 rkey: rkey.clone(), 758 cid: record_cid, 759 }], 760 record_deletes: vec![], 761 backlinks_to_add: vec![], 762 backlinks_to_remove: vec![], 763 commit_event: CommitEventData { 764 did: did.clone(), 765 event_type: RepoEventType::Commit, 766 commit_cid: Some(mid_root.clone()), 767 prev_cid: None, 768 ops: None, 769 blobs: None, 770 blocks_cids: None, 771 prev_data_cid: None, 772 rev: Some("rev1".to_string()), 773 }, 774 }; 775 ops.apply_commit(insert_input).unwrap(); 776 777 assert!( 778 h.metastore 779 .record_ops() 780 .get_record_cid(user_id, &collection, &rkey) 781 .unwrap() 782 .is_some() 783 ); 784 785 let final_root = test_cid_link(23); 786 let delete_input = ApplyCommitInput { 787 user_id, 788 did: did.clone(), 789 expected_root_cid: Some(mid_root.clone()), 790 new_root_cid: final_root.clone(), 791 new_rev: "rev2".to_string(), 792 new_block_cids: vec![], 793 obsolete_block_cids: vec![], 794 record_upserts: vec![], 795 record_deletes: vec![tranquil_db_traits::RecordDelete { 796 collection: collection.clone(), 797 rkey: rkey.clone(), 798 }], 799 backlinks_to_add: vec![], 800 backlinks_to_remove: vec![], 801 commit_event: CommitEventData { 802 did: did.clone(), 803 event_type: RepoEventType::Commit, 804 commit_cid: Some(final_root.clone()), 805 prev_cid: Some(mid_root.clone()), 806 ops: None, 807 blobs: None, 808 blocks_cids: None, 809 prev_data_cid: None, 810 rev: Some("rev2".to_string()), 811 }, 812 }; 813 ops.apply_commit(delete_input).unwrap(); 814 815 assert!( 816 h.metastore 817 .record_ops() 818 .get_record_cid(user_id, &collection, &rkey) 819 .unwrap() 820 .is_none() 821 ); 822 } 823 824 #[test] 825 fn apply_commit_event_visible_after_commit() { 826 let h = setup(); 827 let ops = make_commit_ops(&h); 828 let (user_id, did, root_cid) = create_test_repo(&h, "dave", 30); 829 830 let new_root = test_cid_link(31); 831 let input = ApplyCommitInput { 832 user_id, 833 did: did.clone(), 834 expected_root_cid: Some(root_cid.clone()), 835 new_root_cid: new_root.clone(), 836 new_rev: "rev1".to_string(), 837 new_block_cids: vec![], 838 obsolete_block_cids: vec![], 839 record_upserts: vec![], 840 record_deletes: vec![], 841 backlinks_to_add: vec![], 842 backlinks_to_remove: vec![], 843 commit_event: CommitEventData { 844 did: did.clone(), 845 event_type: RepoEventType::Commit, 846 commit_cid: Some(new_root.clone()), 847 prev_cid: Some(root_cid.clone()), 848 ops: None, 849 blobs: None, 850 blocks_cids: None, 851 prev_data_cid: None, 852 rev: Some("rev1".to_string()), 853 }, 854 }; 855 856 let result = ops.apply_commit(input).unwrap(); 857 let seq = SequenceNumber::from_raw(result.seq); 858 859 let event = ops.event_ops.get_event_by_seq(seq).unwrap().unwrap(); 860 assert_eq!(event.did, did); 861 assert_eq!(event.event_type, RepoEventType::Commit); 862 assert_eq!(event.rev.as_deref(), Some("rev1")); 863 } 864 865 #[test] 866 fn import_repo_data_inserts_records() { 867 let h = setup(); 868 let ops = make_commit_ops(&h); 869 let (user_id, _did, root_cid) = create_test_repo(&h, "eve", 40); 870 871 let collection = Nsid::from("app.bsky.feed.post".to_string()); 872 let rkey = Rkey::from("3k2import".to_string()); 873 let record_cid = test_cid_link(41); 874 875 ops.import_repo_data( 876 user_id, 877 &[], 878 &[ImportRecord { 879 collection: collection.clone(), 880 rkey: rkey.clone(), 881 record_cid: record_cid.clone(), 882 }], 883 Some(&root_cid), 884 ) 885 .unwrap(); 886 887 let found = h 888 .metastore 889 .record_ops() 890 .get_record_cid(user_id, &collection, &rkey) 891 .unwrap() 892 .unwrap(); 893 assert_eq!(found, record_cid); 894 } 895 896 #[test] 897 fn import_repo_data_cas_rejects_stale_root() { 898 let h = setup(); 899 let ops = make_commit_ops(&h); 900 let (user_id, _did, _root_cid) = create_test_repo(&h, "frank", 50); 901 902 let stale = test_cid_link(99); 903 let result = ops.import_repo_data(user_id, &[], &[], Some(&stale)); 904 assert_eq!(result.unwrap_err(), ImportRepoError::ConcurrentModification); 905 } 906 907 #[test] 908 fn insert_record_blobs_and_backfill_query() { 909 let h = setup(); 910 let ops = make_commit_ops(&h); 911 let (user_id_a, did_a, _) = create_test_repo(&h, "grace", 60); 912 let (user_id_b, _did_b, _) = create_test_repo(&h, "henry", 61); 913 914 let needing = ops.get_users_needing_record_blobs_backfill(100).unwrap(); 915 assert_eq!(needing.len(), 2); 916 917 let uri = AtUri::from_parts(did_a.as_str(), "app.bsky.feed.post", "3k2abc"); 918 let blob_cid = test_cid_link(62); 919 ops.insert_record_blobs(user_id_a, &[uri], &[blob_cid]) 920 .unwrap(); 921 922 let needing_after = ops.get_users_needing_record_blobs_backfill(100).unwrap(); 923 assert_eq!(needing_after.len(), 1); 924 assert_eq!(needing_after[0].user_id, user_id_b); 925 } 926 927 #[test] 928 fn get_users_without_blocks_returns_users_with_no_blocks() { 929 let h = setup(); 930 let ops = make_commit_ops(&h); 931 let (user_id_a, did_a, root_a) = create_test_repo(&h, "ivan", 70); 932 let (user_id_b, _did_b, _root_b) = create_test_repo(&h, "julia", 71); 933 934 let new_root = test_cid_link(72); 935 let input = ApplyCommitInput { 936 user_id: user_id_a, 937 did: did_a.clone(), 938 expected_root_cid: Some(root_a), 939 new_root_cid: new_root.clone(), 940 new_rev: "rev1".to_string(), 941 new_block_cids: vec![vec![0x01, 0x02, 0x03]], 942 obsolete_block_cids: vec![], 943 record_upserts: vec![], 944 record_deletes: vec![], 945 backlinks_to_add: vec![], 946 backlinks_to_remove: vec![], 947 commit_event: CommitEventData { 948 did: did_a.clone(), 949 event_type: RepoEventType::Commit, 950 commit_cid: Some(new_root), 951 prev_cid: None, 952 ops: None, 953 blobs: None, 954 blocks_cids: None, 955 prev_data_cid: None, 956 rev: Some("rev1".to_string()), 957 }, 958 }; 959 ops.apply_commit(input).unwrap(); 960 961 let without = ops.get_users_without_blocks().unwrap(); 962 assert_eq!(without.len(), 1); 963 assert_eq!(without[0].user_id, user_id_b); 964 } 965 966 #[test] 967 fn apply_commit_without_expected_root_skips_cas() { 968 let h = setup(); 969 let ops = make_commit_ops(&h); 970 let (user_id, did, _root_cid) = create_test_repo(&h, "kate", 80); 971 972 let new_root = test_cid_link(81); 973 let input = ApplyCommitInput { 974 user_id, 975 did: did.clone(), 976 expected_root_cid: None, 977 new_root_cid: new_root.clone(), 978 new_rev: "rev_force".to_string(), 979 new_block_cids: vec![], 980 obsolete_block_cids: vec![], 981 record_upserts: vec![], 982 record_deletes: vec![], 983 backlinks_to_add: vec![], 984 backlinks_to_remove: vec![], 985 commit_event: CommitEventData { 986 did, 987 event_type: RepoEventType::Commit, 988 commit_cid: None, 989 prev_cid: None, 990 ops: None, 991 blobs: None, 992 blocks_cids: None, 993 prev_data_cid: None, 994 rev: Some("rev_force".to_string()), 995 }, 996 }; 997 998 let result = ops.apply_commit(input).unwrap(); 999 assert!(result.seq > 0); 1000 } 1001 1002 #[test] 1003 fn apply_commit_update_preserves_new_backlinks() { 1004 use crate::metastore::backlinks::backlink_target_prefix; 1005 use crate::metastore::partitions::Partition; 1006 1007 let h = setup(); 1008 let ops = make_commit_ops(&h); 1009 let (user_id, did, root_cid) = create_test_repo(&h, "backlink_upd", 90); 1010 1011 let collection = Nsid::from("app.bsky.feed.like".to_string()); 1012 let rkey = Rkey::from("3k2like1".to_string()); 1013 let record_cid = test_cid_link(91); 1014 let record_uri = AtUri::from_parts(did.as_str(), collection.as_str(), rkey.as_str()); 1015 1016 let mid_root = test_cid_link(92); 1017 let create_input = ApplyCommitInput { 1018 user_id, 1019 did: did.clone(), 1020 expected_root_cid: Some(root_cid.clone()), 1021 new_root_cid: mid_root.clone(), 1022 new_rev: "rev1".to_string(), 1023 new_block_cids: vec![], 1024 obsolete_block_cids: vec![], 1025 record_upserts: vec![tranquil_db_traits::RecordUpsert { 1026 collection: collection.clone(), 1027 rkey: rkey.clone(), 1028 cid: record_cid.clone(), 1029 }], 1030 record_deletes: vec![], 1031 backlinks_to_add: vec![tranquil_db_traits::Backlink { 1032 uri: record_uri.clone(), 1033 path: tranquil_db_traits::BacklinkPath::SubjectUri, 1034 link_to: "at://did:plc:target_a/app.bsky.feed.post/p1".to_string(), 1035 }], 1036 backlinks_to_remove: vec![], 1037 commit_event: CommitEventData { 1038 did: did.clone(), 1039 event_type: RepoEventType::Commit, 1040 commit_cid: Some(mid_root.clone()), 1041 prev_cid: Some(root_cid.clone()), 1042 ops: None, 1043 blobs: None, 1044 blocks_cids: None, 1045 prev_data_cid: None, 1046 rev: Some("rev1".to_string()), 1047 }, 1048 }; 1049 ops.apply_commit(create_input).unwrap(); 1050 1051 let indexes = h.metastore.partition(Partition::Indexes); 1052 let target_a_prefix = backlink_target_prefix("at://did:plc:target_a/app.bsky.feed.post/p1"); 1053 let count_a_before = indexes 1054 .prefix(target_a_prefix.as_slice()) 1055 .map(|g| g.into_inner().expect("scan must not fail")) 1056 .fold(0, |acc, _| acc + 1); 1057 assert_eq!(count_a_before, 1); 1058 1059 let final_root = test_cid_link(93); 1060 let new_record_cid = test_cid_link(94); 1061 let update_input = ApplyCommitInput { 1062 user_id, 1063 did: did.clone(), 1064 expected_root_cid: Some(mid_root.clone()), 1065 new_root_cid: final_root.clone(), 1066 new_rev: "rev2".to_string(), 1067 new_block_cids: vec![], 1068 obsolete_block_cids: vec![], 1069 record_upserts: vec![tranquil_db_traits::RecordUpsert { 1070 collection: collection.clone(), 1071 rkey: rkey.clone(), 1072 cid: new_record_cid.clone(), 1073 }], 1074 record_deletes: vec![], 1075 backlinks_to_add: vec![tranquil_db_traits::Backlink { 1076 uri: record_uri.clone(), 1077 path: tranquil_db_traits::BacklinkPath::SubjectUri, 1078 link_to: "at://did:plc:target_b/app.bsky.feed.post/p2".to_string(), 1079 }], 1080 backlinks_to_remove: vec![record_uri], 1081 commit_event: CommitEventData { 1082 did: did.clone(), 1083 event_type: RepoEventType::Commit, 1084 commit_cid: Some(final_root.clone()), 1085 prev_cid: Some(mid_root.clone()), 1086 ops: None, 1087 blobs: None, 1088 blocks_cids: None, 1089 prev_data_cid: None, 1090 rev: Some("rev2".to_string()), 1091 }, 1092 }; 1093 ops.apply_commit(update_input).unwrap(); 1094 1095 let count_a_after = indexes 1096 .prefix(target_a_prefix.as_slice()) 1097 .map(|g| g.into_inner().expect("scan must not fail")) 1098 .fold(0, |acc, _| acc + 1); 1099 assert_eq!(count_a_after, 0); 1100 1101 let target_b_prefix = backlink_target_prefix("at://did:plc:target_b/app.bsky.feed.post/p2"); 1102 let count_b = indexes 1103 .prefix(target_b_prefix.as_slice()) 1104 .map(|g| g.into_inner().expect("scan must not fail")) 1105 .fold(0, |acc, _| acc + 1); 1106 assert_eq!(count_b, 1); 1107 } 1108 1109 #[test] 1110 fn crash_recovery_replays_mutation_set() { 1111 let metastore_dir = tempfile::TempDir::new().unwrap(); 1112 let eventlog_dir = tempfile::TempDir::new().unwrap(); 1113 let segments_dir = eventlog_dir.path().join("segments"); 1114 std::fs::create_dir_all(&segments_dir).unwrap(); 1115 1116 let user_id = Uuid::new_v4(); 1117 let did = test_did("crash_alice"); 1118 let handle = test_handle("crash_alice"); 1119 let initial_root = test_cid_link(200); 1120 let new_root = test_cid_link(201); 1121 let record_cid = test_cid_link(202); 1122 let collection = Nsid::from("app.bsky.feed.post".to_string()); 1123 let rkey = Rkey::from("3k2crash".to_string()); 1124 1125 let event_log = EventLog::open( 1126 EventLogConfig { 1127 segments_dir: segments_dir.clone(), 1128 ..EventLogConfig::default() 1129 }, 1130 RealIO::new(), 1131 ) 1132 .unwrap(); 1133 let event_log = Arc::new(event_log); 1134 let bridge = Arc::new(EventLogBridge::new(Arc::clone(&event_log))); 1135 1136 { 1137 let metastore = Metastore::open( 1138 metastore_dir.path(), 1139 MetastoreConfig { 1140 cache_size_bytes: 64 * 1024 * 1024, 1141 }, 1142 ) 1143 .unwrap(); 1144 1145 metastore 1146 .repo_ops() 1147 .create_repo( 1148 metastore.database(), 1149 user_id, 1150 &did, 1151 &handle, 1152 &initial_root, 1153 "rev0", 1154 ) 1155 .unwrap(); 1156 metastore.persist().unwrap(); 1157 1158 let ops = make_commit_ops_from(&metastore, &bridge); 1159 let input = ApplyCommitInput { 1160 user_id, 1161 did: did.clone(), 1162 expected_root_cid: Some(initial_root.clone()), 1163 new_root_cid: new_root.clone(), 1164 new_rev: "rev1".to_string(), 1165 new_block_cids: vec![vec![0xAA, 0xBB]], 1166 obsolete_block_cids: vec![], 1167 record_upserts: vec![tranquil_db_traits::RecordUpsert { 1168 collection: collection.clone(), 1169 rkey: rkey.clone(), 1170 cid: record_cid.clone(), 1171 }], 1172 record_deletes: vec![], 1173 backlinks_to_add: vec![], 1174 backlinks_to_remove: vec![], 1175 commit_event: CommitEventData { 1176 did: did.clone(), 1177 event_type: RepoEventType::Commit, 1178 commit_cid: Some(new_root.clone()), 1179 prev_cid: Some(initial_root.clone()), 1180 ops: None, 1181 blobs: None, 1182 blocks_cids: None, 1183 prev_data_cid: None, 1184 rev: Some("rev1".to_string()), 1185 }, 1186 }; 1187 1188 let result = ops.apply_commit(input).unwrap(); 1189 assert!(result.seq > 0); 1190 metastore.persist().unwrap(); 1191 } 1192 1193 { 1194 let metastore = Metastore::open( 1195 metastore_dir.path(), 1196 MetastoreConfig { 1197 cache_size_bytes: 64 * 1024 * 1024, 1198 }, 1199 ) 1200 .unwrap(); 1201 1202 let event_ops = metastore.event_ops(Arc::clone(&bridge)); 1203 1204 event_ops.write_last_applied_cursor_direct(0).unwrap(); 1205 metastore.persist().unwrap(); 1206 } 1207 1208 { 1209 let metastore = Metastore::open( 1210 metastore_dir.path(), 1211 MetastoreConfig { 1212 cache_size_bytes: 64 * 1024 * 1024, 1213 }, 1214 ) 1215 .unwrap(); 1216 1217 let repo_before = metastore.repo_ops().get_repo(user_id).unwrap().unwrap(); 1218 assert_eq!(repo_before.repo_root_cid, new_root); 1219 1220 let event_ops = metastore.event_ops(Arc::clone(&bridge)); 1221 let cursor_before = event_ops.read_last_applied_cursor().unwrap(); 1222 assert_eq!(cursor_before, Some(0)); 1223 1224 let indexes = metastore 1225 .partition(crate::metastore::partitions::Partition::Indexes) 1226 .clone(); 1227 let recovered = event_ops.recover_metastore_mutations(&indexes).unwrap(); 1228 assert!(recovered > 0, "should replay at least one event"); 1229 1230 let cursor_after = event_ops.read_last_applied_cursor().unwrap(); 1231 assert!(cursor_after.unwrap_or(0) > 0); 1232 } 1233 } 1234 1235 #[test] 1236 fn crash_recovery_with_uncommitted_batch() { 1237 let metastore_dir = tempfile::TempDir::new().unwrap(); 1238 let eventlog_dir = tempfile::TempDir::new().unwrap(); 1239 let segments_dir = eventlog_dir.path().join("segments"); 1240 std::fs::create_dir_all(&segments_dir).unwrap(); 1241 1242 let user_id = Uuid::new_v4(); 1243 let did = test_did("crash_bob"); 1244 let handle = test_handle("crash_bob"); 1245 let initial_root = test_cid_link(210); 1246 let new_root = test_cid_link(211); 1247 let record_cid = test_cid_link(212); 1248 let collection = Nsid::from("app.bsky.feed.post".to_string()); 1249 let rkey = Rkey::from("3k2bob".to_string()); 1250 1251 let event_log = EventLog::open( 1252 EventLogConfig { 1253 segments_dir: segments_dir.clone(), 1254 ..EventLogConfig::default() 1255 }, 1256 RealIO::new(), 1257 ) 1258 .unwrap(); 1259 let event_log = Arc::new(event_log); 1260 let bridge = Arc::new(EventLogBridge::new(Arc::clone(&event_log))); 1261 1262 { 1263 let metastore = Metastore::open( 1264 metastore_dir.path(), 1265 MetastoreConfig { 1266 cache_size_bytes: 64 * 1024 * 1024, 1267 }, 1268 ) 1269 .unwrap(); 1270 1271 metastore 1272 .repo_ops() 1273 .create_repo( 1274 metastore.database(), 1275 user_id, 1276 &did, 1277 &handle, 1278 &initial_root, 1279 "rev0", 1280 ) 1281 .unwrap(); 1282 metastore.persist().unwrap(); 1283 1284 let event_ops = metastore.event_ops(Arc::clone(&bridge)); 1285 let mutation_set = super::CommitMutationSet { 1286 new_root_cid: super::cid_link_to_bytes(&new_root).unwrap(), 1287 new_rev: "rev1".to_string(), 1288 record_upserts: vec![super::RecordMutationUpsert { 1289 collection: collection.as_str().to_owned(), 1290 rkey: rkey.as_str().to_owned(), 1291 cid_bytes: super::cid_link_to_bytes(&record_cid).unwrap(), 1292 }], 1293 record_deletes: vec![], 1294 block_inserts: vec![vec![0xCC, 0xDD]], 1295 block_deletes: vec![], 1296 backlink_adds: vec![], 1297 backlink_remove_uris: vec![], 1298 }; 1299 let ms_bytes = mutation_set.serialize().unwrap(); 1300 1301 let commit_data = CommitEventData { 1302 did: did.clone(), 1303 event_type: RepoEventType::Commit, 1304 commit_cid: Some(new_root.clone()), 1305 prev_cid: Some(initial_root.clone()), 1306 ops: None, 1307 blobs: None, 1308 blocks_cids: None, 1309 prev_data_cid: None, 1310 rev: Some("rev1".to_string()), 1311 }; 1312 1313 let mut batch = metastore.database().batch(); 1314 let (_seq, deferred) = event_ops 1315 .append_commit_event_into_batch(&mut batch, &commit_data, Some(&ms_bytes)) 1316 .unwrap(); 1317 1318 event_ops.complete_broadcast(deferred); 1319 1320 drop(batch); 1321 1322 metastore.persist().unwrap(); 1323 } 1324 1325 { 1326 let metastore = Metastore::open( 1327 metastore_dir.path(), 1328 MetastoreConfig { 1329 cache_size_bytes: 64 * 1024 * 1024, 1330 }, 1331 ) 1332 .unwrap(); 1333 1334 let repo = metastore.repo_ops().get_repo(user_id).unwrap().unwrap(); 1335 assert_eq!(repo.repo_root_cid, initial_root); 1336 1337 let record = metastore 1338 .record_ops() 1339 .get_record_cid(user_id, &collection, &rkey) 1340 .unwrap(); 1341 assert!(record.is_none()); 1342 1343 let event_ops = metastore.event_ops(Arc::clone(&bridge)); 1344 let indexes = metastore 1345 .partition(crate::metastore::partitions::Partition::Indexes) 1346 .clone(); 1347 let recovered = event_ops.recover_metastore_mutations(&indexes).unwrap(); 1348 assert_eq!(recovered, 1); 1349 1350 let repo_after = metastore.repo_ops().get_repo(user_id).unwrap().unwrap(); 1351 assert_eq!(repo_after.repo_root_cid, new_root); 1352 assert_eq!(repo_after.repo_rev.as_deref(), Some("rev1")); 1353 1354 let record_after = metastore 1355 .record_ops() 1356 .get_record_cid(user_id, &collection, &rkey) 1357 .unwrap(); 1358 assert_eq!(record_after, Some(record_cid)); 1359 } 1360 } 1361 1362 fn make_commit_ops_from( 1363 metastore: &Metastore, 1364 bridge: &Arc<EventLogBridge<RealIO>>, 1365 ) -> CommitOps<RealIO> { 1366 use crate::metastore::partitions::Partition; 1367 CommitOps::new( 1368 metastore.database().clone(), 1369 metastore.partition(Partition::RepoData).clone(), 1370 metastore.partition(Partition::Indexes).clone(), 1371 Arc::clone(metastore.user_hashes()), 1372 Arc::clone(bridge), 1373 ) 1374 } 1375 1376 #[test] 1377 fn apply_commit_backlinks_isolated_by_collection() { 1378 use crate::metastore::backlinks::{backlink_by_user_prefix, backlink_target_prefix}; 1379 use crate::metastore::partitions::Partition; 1380 1381 let h = setup(); 1382 let ops = make_commit_ops(&h); 1383 let (user_id, did, root_cid) = create_test_repo(&h, "col_iso", 95); 1384 1385 let col_like = Nsid::from("app.bsky.feed.like".to_string()); 1386 let col_repost = Nsid::from("app.bsky.feed.repost".to_string()); 1387 let rkey = Rkey::from("same_rkey".to_string()); 1388 let target = "at://did:plc:someone/app.bsky.feed.post/p1"; 1389 1390 let mid_root = test_cid_link(96); 1391 let uri_like = AtUri::from_parts(did.as_str(), col_like.as_str(), rkey.as_str()); 1392 let uri_repost = AtUri::from_parts(did.as_str(), col_repost.as_str(), rkey.as_str()); 1393 1394 let input = ApplyCommitInput { 1395 user_id, 1396 did: did.clone(), 1397 expected_root_cid: Some(root_cid.clone()), 1398 new_root_cid: mid_root.clone(), 1399 new_rev: "rev1".to_string(), 1400 new_block_cids: vec![], 1401 obsolete_block_cids: vec![], 1402 record_upserts: vec![ 1403 tranquil_db_traits::RecordUpsert { 1404 collection: col_like.clone(), 1405 rkey: rkey.clone(), 1406 cid: test_cid_link(97), 1407 }, 1408 tranquil_db_traits::RecordUpsert { 1409 collection: col_repost.clone(), 1410 rkey: rkey.clone(), 1411 cid: test_cid_link(98), 1412 }, 1413 ], 1414 record_deletes: vec![], 1415 backlinks_to_add: vec![ 1416 tranquil_db_traits::Backlink { 1417 uri: uri_like.clone(), 1418 path: tranquil_db_traits::BacklinkPath::SubjectUri, 1419 link_to: target.to_string(), 1420 }, 1421 tranquil_db_traits::Backlink { 1422 uri: uri_repost.clone(), 1423 path: tranquil_db_traits::BacklinkPath::SubjectUri, 1424 link_to: target.to_string(), 1425 }, 1426 ], 1427 backlinks_to_remove: vec![], 1428 commit_event: CommitEventData { 1429 did: did.clone(), 1430 event_type: RepoEventType::Commit, 1431 commit_cid: Some(mid_root.clone()), 1432 prev_cid: Some(root_cid.clone()), 1433 ops: None, 1434 blobs: None, 1435 blocks_cids: None, 1436 prev_data_cid: None, 1437 rev: Some("rev1".to_string()), 1438 }, 1439 }; 1440 ops.apply_commit(input).unwrap(); 1441 1442 let indexes = h.metastore.partition(Partition::Indexes); 1443 let user_hash = h.metastore.user_hashes().get(&user_id).unwrap(); 1444 let target_prefix = backlink_target_prefix(target); 1445 let user_prefix = backlink_by_user_prefix(user_hash); 1446 1447 assert_eq!( 1448 indexes 1449 .prefix(target_prefix.as_slice()) 1450 .map(|g| g.into_inner().expect("scan must not fail")) 1451 .fold(0, |acc, _| acc + 1), 1452 2 1453 ); 1454 assert_eq!( 1455 indexes 1456 .prefix(user_prefix.as_slice()) 1457 .map(|g| g.into_inner().expect("scan must not fail")) 1458 .fold(0, |acc, _| acc + 1), 1459 2 1460 ); 1461 1462 let final_root = test_cid_link(99); 1463 let remove_like = ApplyCommitInput { 1464 user_id, 1465 did: did.clone(), 1466 expected_root_cid: Some(mid_root.clone()), 1467 new_root_cid: final_root.clone(), 1468 new_rev: "rev2".to_string(), 1469 new_block_cids: vec![], 1470 obsolete_block_cids: vec![], 1471 record_upserts: vec![], 1472 record_deletes: vec![], 1473 backlinks_to_add: vec![], 1474 backlinks_to_remove: vec![uri_like], 1475 commit_event: CommitEventData { 1476 did: did.clone(), 1477 event_type: RepoEventType::Commit, 1478 commit_cid: Some(final_root.clone()), 1479 prev_cid: Some(mid_root.clone()), 1480 ops: None, 1481 blobs: None, 1482 blocks_cids: None, 1483 prev_data_cid: None, 1484 rev: Some("rev2".to_string()), 1485 }, 1486 }; 1487 ops.apply_commit(remove_like).unwrap(); 1488 1489 assert_eq!( 1490 indexes 1491 .prefix(target_prefix.as_slice()) 1492 .map(|g| g.into_inner().expect("scan must not fail")) 1493 .fold(0, |acc, _| acc + 1), 1494 1 1495 ); 1496 assert_eq!( 1497 indexes 1498 .prefix(user_prefix.as_slice()) 1499 .map(|g| g.into_inner().expect("scan must not fail")) 1500 .fold(0, |acc, _| acc + 1), 1501 1 1502 ); 1503 } 1504}