forked from
tranquil.farm/tranquil-pds
Our Personal Data Server from scratch!
53 kB
1504 lines
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}