diff --git a/crates/tranquil-pds/tests/gc_after_delete.rs b/crates/tranquil-pds/tests/gc_after_delete.rs index 8169821..d237f89 100644 --- a/crates/tranquil-pds/tests/gc_after_delete.rs +++ b/crates/tranquil-pds/tests/gc_after_delete.rs @@ -5,7 +5,7 @@ use common::*; use helpers::*; use reqwest::StatusCode; use serde_json::{Value, json}; -use tranquil_types::Did; +use tranquil_types::{Did, Nsid, Rkey}; #[tokio::test] async fn test_delete_record_marks_blocks_obsolete() { @@ -13,22 +13,17 @@ async fn test_delete_record_marks_blocks_obsolete() { let base = base_url().await; let repos = get_test_repos().await; let (did, jwt) = setup_new_user("gc-after-delete").await; + let did = Did::new(did).expect("setup_new_user returned a valid DID"); let user_id = repos .user - .get_id_by_did(&Did::new(did.clone()).unwrap()) + .get_id_by_did(&did) .await .expect("DB error") .expect("User not found"); - let count_baseline = repos - .repo - .count_user_blocks(user_id) - .await - .expect("count_user_blocks failed"); - - let collection = "app.bsky.feed.post"; - let rkey = format!("gc_test_{}", Utc::now().timestamp_millis()); + let collection = Nsid::new("app.bsky.feed.post".to_string()).expect("valid NSID"); + let rkey = Rkey::new(format!("gc_test_{}", Utc::now().timestamp_millis())).expect("valid rkey"); let create_payload = json!({ "repo": did, "collection": collection, @@ -65,17 +60,7 @@ async fn test_delete_record_marks_blocks_obsolete() { .expect("createRecord response missing cid") .to_string(); - let count_after_create = repos - .repo - .count_user_blocks(user_id) - .await - .expect("count_user_blocks failed"); - assert!( - count_after_create > count_baseline, - "user_blocks count did not grow after createRecord (baseline={}, after_create={})", - count_baseline, - count_after_create - ); + assert_user_blocks_matches_repo(user_id, "createRecord").await; let delete_payload = json!({ "repo": did, @@ -96,28 +81,13 @@ async fn test_delete_record_marks_blocks_obsolete() { delete_res.text().await ); - let count_after_delete = repos - .repo - .count_user_blocks(user_id) - .await - .expect("count_user_blocks failed"); - - assert!( - count_after_delete < count_after_create, - "user_blocks count did not shrink after deleteRecord \ - (baseline={}, after_create={}, after_delete={}). \ - The delete path produced no obsolete CIDs beyond the prior commit root, \ - which is the regression this test guards against.", - count_baseline, - count_after_create, - count_after_delete - ); + assert_user_blocks_matches_repo(user_id, "deleteRecord").await; let get_res = client .get(format!("{}/xrpc/com.atproto.repo.getRecord", base)) .query(&[ ("repo", did.as_str()), - ("collection", collection), + ("collection", collection.as_str()), ("rkey", rkey.as_str()), ]) .send() @@ -138,16 +108,18 @@ async fn test_update_record_marks_old_record_block_obsolete() { let base = base_url().await; let repos = get_test_repos().await; let (did, jwt) = setup_new_user("gc-after-update").await; + let did = Did::new(did).expect("setup_new_user returned a valid DID"); let user_id = repos .user - .get_id_by_did(&Did::new(did.clone()).unwrap()) + .get_id_by_did(&did) .await .expect("DB error") .expect("User not found"); - let collection = "app.bsky.feed.post"; - let rkey = format!("gc_update_{}", Utc::now().timestamp_millis()); + let collection = Nsid::new("app.bsky.feed.post".to_string()).expect("valid NSID"); + let rkey = + Rkey::new(format!("gc_update_{}", Utc::now().timestamp_millis())).expect("valid rkey"); let put_v1 = json!({ "repo": did, @@ -168,11 +140,7 @@ async fn test_update_record_marks_old_record_block_obsolete() { .expect("Failed to send putRecord v1"); assert_eq!(res.status(), StatusCode::OK, "first putRecord failed"); - let count_after_create = repos - .repo - .count_user_blocks(user_id) - .await - .expect("count_user_blocks failed"); + assert_user_blocks_matches_repo(user_id, "the first putRecord").await; let put_v2 = json!({ "repo": did, @@ -193,133 +161,7 @@ async fn test_update_record_marks_old_record_block_obsolete() { .expect("Failed to send putRecord v2"); assert_eq!(res.status(), StatusCode::OK, "second putRecord failed"); - let count_after_update = repos - .repo - .count_user_blocks(user_id) - .await - .expect("count_user_blocks failed"); - - assert!( - count_after_update <= count_after_create + 1, - "user_blocks count grew by more than 1 after putRecord update \ - (after_create={}, after_update={}). The previous version's record block \ - should have been marked obsolete; instead it appears to be leaking.", - count_after_create, - count_after_update - ); -} - -#[tokio::test] -async fn test_delete_in_populated_repo_marks_merged_subtree_blocks_obsolete() { - let client = client(); - let base = base_url().await; - let repos = get_test_repos().await; - let (did, jwt) = setup_new_user("gc-merge").await; - - let user_id = repos - .user - .get_id_by_did(&Did::new(did.clone()).unwrap()) - .await - .expect("DB error") - .expect("User not found"); - - let collection = "app.bsky.feed.post"; - let record_count = 64usize; - let now_ms = Utc::now().timestamp_millis(); - - let rkeys: Vec = (0..record_count) - .map(|i| format!("gc_merge_{}_{:04}", now_ms, i)) - .collect(); - - let create_results = - futures::future::try_join_all(rkeys.iter().enumerate().map(|(i, rkey)| { - let client = client.clone(); - let jwt = jwt.clone(); - let did = did.clone(); - let base = base.to_string(); - let payload = json!({ - "repo": did, - "collection": collection, - "rkey": rkey, - "record": { - "$type": collection, - "text": format!("seed record {}", i), - "createdAt": Utc::now().to_rfc3339() - } - }); - async move { - let res = client - .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) - .bearer_auth(&jwt) - .json(&payload) - .send() - .await - .expect("Failed to send createRecord"); - if res.status() != StatusCode::OK { - return Err(format!("seed createRecord failed: {}", res.status())); - } - Ok::<(), String>(()) - } - })) - .await; - create_results.expect("seeding records failed"); - - let count_after_seed = repos - .repo - .count_user_blocks(user_id) - .await - .expect("count_user_blocks failed"); - - let target_rkey = &rkeys[record_count / 2]; - let delete_payload = json!({ - "repo": did, - "collection": collection, - "rkey": target_rkey, - }); - let delete_res = client - .post(format!("{}/xrpc/com.atproto.repo.deleteRecord", base)) - .bearer_auth(&jwt) - .json(&delete_payload) - .send() - .await - .expect("Failed to send deleteRecord"); - assert_eq!( - delete_res.status(), - StatusCode::OK, - "deleteRecord did not return 200: {:?}", - delete_res.text().await - ); - - let count_after_delete = repos - .repo - .count_user_blocks(user_id) - .await - .expect("count_user_blocks failed"); - assert!( - count_after_delete < count_after_seed, - "user_blocks did not shrink after deleting from a populated repo \ - (after_seed={}, after_delete={}). The path-walk-based obsolete \ - calculation does not capture sibling subtree blocks orphaned by \ - delete-merge; only an MST-diff-based calculation does.", - count_after_seed, - count_after_delete - ); - - let get_res = client - .get(format!("{}/xrpc/com.atproto.repo.getRecord", base)) - .query(&[ - ("repo", did.as_str()), - ("collection", collection), - ("rkey", target_rkey.as_str()), - ]) - .send() - .await - .expect("Failed to send getRecord"); - assert!( - !get_res.status().is_success(), - "deleted record is still resolvable via getRecord (status={})", - get_res.status(), - ); + assert_user_blocks_matches_repo(user_id, "the second putRecord").await; } #[tokio::test] @@ -339,9 +181,11 @@ async fn test_delete_decrements_tranquil_store_refcounts() { .as_tranquil_store() .expect("tranquil-store backend selected but block_store is not TranquilStore"); let (did, jwt) = setup_new_user("gc-store-decrement").await; + let did = Did::new(did).expect("setup_new_user returned a valid DID"); - let collection = "app.bsky.feed.post"; - let rkey = format!("gc_store_{}", Utc::now().timestamp_millis()); + let collection = Nsid::new("app.bsky.feed.post".to_string()).expect("valid NSID"); + let rkey = + Rkey::new(format!("gc_store_{}", Utc::now().timestamp_millis())).expect("valid rkey"); let create_res = client .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) diff --git a/crates/tranquil-pds/tests/helpers/mod.rs b/crates/tranquil-pds/tests/helpers/mod.rs index 534e9a3..6f6f6e3 100644 --- a/crates/tranquil-pds/tests/helpers/mod.rs +++ b/crates/tranquil-pds/tests/helpers/mod.rs @@ -480,6 +480,69 @@ pub fn get_multikey_from_signing_key(signing_key: &k256::ecdsa::SigningKey) -> S multibase::encode(multibase::Base::Base58Btc, buf) } +async fn reachable_blocks(user_id: uuid::Uuid) -> std::collections::BTreeSet { + use jacquard_repo::storage::BlockStore; + + let repos = super::common::get_test_repos().await; + let store = super::common::get_test_block_store().await; + + let root_str = repos + .repo + .get_repo_root_cid_by_user_id(user_id) + .await + .expect("DB error fetching repo root") + .expect("repo root not found"); + let root_cid = cid::Cid::try_from(root_str.as_str()).expect("repo root isn't a valid CID"); + let commit_bytes = store + .get(&root_cid) + .await + .expect("block store error fetching commit") + .expect("commit block not in block store"); + let data_cid = jacquard_repo::commit::Commit::from_cbor(&commit_bytes) + .expect("commit block doesn't parse") + .data; + + let mst = jacquard_repo::mst::Mst::load(std::sync::Arc::new(store.clone()), data_cid, None); + let mut cids = tranquil_pds::repo_ops::reachable_tree_cids(&mst) + .await + .expect("walking the new MST failed"); + cids.insert(root_cid); + cids +} + +async fn recorded_blocks(user_id: uuid::Uuid) -> std::collections::BTreeSet { + super::common::get_test_repos() + .await + .repo + .get_user_block_cids_since_rev(user_id, None) + .await + .expect("get_user_block_cids_since_rev failed") + .iter() + .map(|b| cid::Cid::try_from(b.as_slice()).expect("invalid CID in user_blocks")) + .collect() +} + +#[allow(dead_code)] +pub async fn assert_user_blocks_matches_repo(user_id: uuid::Uuid, phase: &str) { + let reachable = reachable_blocks(user_id).await; + let recorded = recorded_blocks(user_id).await; + let missing: Vec = reachable + .difference(&recorded) + .map(cid::Cid::to_string) + .collect(); + let stale: Vec = recorded + .difference(&reachable) + .map(cid::Cid::to_string) + .collect(); + assert!( + missing.is_empty() && stale.is_empty(), + "user_blocks doesn't match the blocks reachable from the repo root after {phase}. \ + reachable={} recorded={} missing={missing:?} stale={stale:?}", + reachable.len(), + recorded.len(), + ); +} + #[allow(dead_code)] pub async fn get_user_signing_key(did: &str) -> Option> { let repos = super::common::get_test_repos().await; diff --git a/crates/tranquil-pds/tests/user_blocks_reachability.rs b/crates/tranquil-pds/tests/user_blocks_reachability.rs new file mode 100644 index 0000000..b94cfd9 --- /dev/null +++ b/crates/tranquil-pds/tests/user_blocks_reachability.rs @@ -0,0 +1,387 @@ +mod common; +mod helpers; +use chrono::Utc; +use common::*; +use helpers::*; +use reqwest::StatusCode; +use serde_json::json; +use std::sync::LazyLock; +use tranquil_types::{Did, Nsid, Rkey}; + +static COLLECTION: LazyLock = + LazyLock::new(|| Nsid::new("app.bsky.feed.post".to_string()).expect("valid NSID")); + +async fn create_record(did: &Did, jwt: &str, rkey: &Rkey, text: String) -> cid::Cid { + create_record_at(did, jwt, rkey, text, Utc::now().to_rfc3339()).await +} + +async fn create_record_at( + did: &Did, + jwt: &str, + rkey: &Rkey, + text: String, + created_at: String, +) -> cid::Cid { + let res = client() + .post(format!( + "{}/xrpc/com.atproto.repo.createRecord", + base_url().await + )) + .bearer_auth(jwt) + .json(&json!({ + "repo": did, + "collection": &*COLLECTION, + "rkey": rkey, + "record": { + "$type": &*COLLECTION, + "text": text, + "createdAt": created_at + } + })) + .send() + .await + .expect("Failed to send createRecord"); + assert_eq!( + res.status(), + StatusCode::OK, + "createRecord for {rkey} didn't return 200: {:?}", + res.text().await + ); + let body: serde_json::Value = res.json().await.expect("createRecord response isn't JSON"); + let cid_str = body["cid"] + .as_str() + .expect("createRecord response missing cid"); + cid::Cid::try_from(cid_str).expect("createRecord returned an invalid cid") +} + +async fn refcount_of(cid: &cid::Cid) -> u32 { + get_test_block_store() + .await + .as_tranquil_store() + .expect("tranquil-store backend selected but block_store isn't TranquilStore") + .refcount_of(cid) + .expect("refcount_of failed") + .unwrap_or(0) +} + +async fn delete_record(did: &Did, jwt: &str, rkey: &Rkey) { + let res = client() + .post(format!( + "{}/xrpc/com.atproto.repo.deleteRecord", + base_url().await + )) + .bearer_auth(jwt) + .json(&json!({ + "repo": did, + "collection": &*COLLECTION, + "rkey": rkey, + })) + .send() + .await + .expect("Failed to send deleteRecord"); + assert_eq!( + res.status(), + StatusCode::OK, + "deleteRecord for {rkey} didn't return 200: {:?}", + res.text().await + ); +} + +async fn assert_record_gone(did: &Did, rkey: &Rkey) { + let res = client() + .get(format!( + "{}/xrpc/com.atproto.repo.getRecord", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", COLLECTION.as_str()), + ("rkey", rkey.as_str()), + ]) + .send() + .await + .expect("Failed to send getRecord"); + assert!( + !res.status().is_success(), + "deleted record {rkey} is still resolvable via getRecord: {}", + res.status() + ); +} + +async fn user_id_for(did: &Did) -> uuid::Uuid { + get_test_repos() + .await + .user + .get_id_by_did(did) + .await + .expect("DB error looking up the user id") + .expect("User not found") +} + +#[tokio::test] +async fn deleting_from_a_populated_repo_keeps_user_blocks_equal_to_reachable_set() { + let (did, jwt) = setup_new_user("user-blocks-reachability").await; + let did = Did::new(did).expect("setup_new_user returned a valid DID"); + let user_id = user_id_for(&did).await; + + let record_count = 64usize; + let now_ms = Utc::now().timestamp_millis(); + let rkeys: Vec = (0..record_count) + .map(|i| Rkey::new(format!("reach_{}_{:04}", now_ms, i)).expect("valid rkey")) + .collect(); + + futures::future::join_all( + rkeys + .iter() + .enumerate() + .map(|(i, rkey)| create_record(&did, &jwt, rkey, format!("seed record {}", i))), + ) + .await; + + assert_user_blocks_matches_repo(user_id, "64 creates").await; + + let target = &rkeys[record_count / 2]; + delete_record(&did, &jwt, target).await; + + assert_user_blocks_matches_repo(user_id, "deleting one record").await; + assert_record_gone(&did, target).await; +} + +#[tokio::test] +async fn applying_a_multi_op_batch_keeps_user_blocks_equal_to_reachable_set() { + let (did, jwt) = setup_new_user("user-blocks-apply-writes").await; + let did = Did::new(did).expect("setup_new_user returned a valid DID"); + let user_id = user_id_for(&did).await; + + let now_ms = Utc::now().timestamp_millis(); + let rkeys: Vec = (0..3) + .map(|i| Rkey::new(format!("batch_{}_{}", now_ms, i)).expect("valid rkey")) + .collect(); + let created_at = "2026-01-01T00:00:00Z".to_string(); + + futures::future::join_all( + rkeys.iter().map(|rkey| { + create_record_at(&did, &jwt, rkey, "shared".to_string(), created_at.clone()) + }), + ) + .await; + assert_user_blocks_matches_repo(user_id, "three records sharing one leaf block").await; + + let writes = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#delete", + "collection": &*COLLECTION, + "rkey": rkeys[0], + }, + { + "$type": "com.atproto.repo.applyWrites#update", + "collection": &*COLLECTION, + "rkey": rkeys[1], + "value": { + "$type": &*COLLECTION, + "text": "updated in the same commit", + "createdAt": created_at + }, + }, + { + "$type": "com.atproto.repo.applyWrites#create", + "collection": &*COLLECTION, + "rkey": Rkey::new(format!("batch_{}_new", now_ms)).expect("valid rkey"), + "value": { + "$type": &*COLLECTION, + "text": "created in the same commit", + "createdAt": created_at + }, + }, + ] + }); + let res = client() + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&jwt) + .json(&writes) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!( + res.status(), + StatusCode::OK, + "applyWrites didn't return 200: {:?}", + res.text().await + ); + + assert_user_blocks_matches_repo(user_id, "a delete, update, & create in one commit").await; + assert_record_gone(&did, &rkeys[0]).await; +} + +#[tokio::test] +async fn emptying_the_repo_keeps_user_blocks_equal_to_reachable_set() { + let (did, jwt) = setup_new_user("user-blocks-empty-tree").await; + let did = Did::new(did).expect("setup_new_user returned a valid DID"); + let user_id = user_id_for(&did).await; + + let now_ms = Utc::now().timestamp_millis(); + let only = Rkey::new(format!("empty_{}_only", now_ms)).expect("valid rkey"); + + create_record(&did, &jwt, &only, "the only record".to_string()).await; + assert_user_blocks_matches_repo(user_id, "creating the only record").await; + + delete_record(&did, &jwt, &only).await; + assert_user_blocks_matches_repo(user_id, "deleting the last record").await; + assert_record_gone(&did, &only).await; +} + +#[tokio::test] +async fn reverting_the_tree_to_a_stored_shape_still_records_its_blocks() { + let (did, jwt) = setup_new_user("user-blocks-resurrect").await; + let did = Did::new(did).expect("setup_new_user returned a valid DID"); + let user_id = user_id_for(&did).await; + + let now_ms = Utc::now().timestamp_millis(); + let kept = Rkey::new(format!("resurrect_{}_kept", now_ms)).expect("valid rkey"); + let churned = Rkey::new(format!("resurrect_{}_churned", now_ms)).expect("valid rkey"); + + create_record(&did, &jwt, &kept, "first record".to_string()).await; + assert_user_blocks_matches_repo(user_id, "creating the first record").await; + + create_record(&did, &jwt, &churned, "second record".to_string()).await; + assert_user_blocks_matches_repo(user_id, "creating the second record").await; + + delete_record(&did, &jwt, &churned).await; + assert_user_blocks_matches_repo(user_id, "deleting back to the one-record tree").await; + + create_record(&did, &jwt, &churned, "second record".to_string()).await; + assert_user_blocks_matches_repo(user_id, "recreating the deleted record").await; +} + +#[tokio::test] +async fn deleting_one_of_two_identical_records_keeps_the_shared_block() { + let (did, jwt) = setup_new_user("user-blocks-shared-leaf").await; + let did = Did::new(did).expect("setup_new_user returned a valid DID"); + let user_id = user_id_for(&did).await; + + let now_ms = Utc::now().timestamp_millis(); + let kept = Rkey::new(format!("shared_{}_kept", now_ms)).expect("valid rkey"); + let dropped = Rkey::new(format!("shared_{}_dropped", now_ms)).expect("valid rkey"); + let created_at = "2026-01-01T00:00:00Z".to_string(); + + let shared = create_record_at( + &did, + &jwt, + &kept, + "identical content".to_string(), + created_at.clone(), + ) + .await; + assert_eq!( + create_record_at( + &did, + &jwt, + &dropped, + "identical content".to_string(), + created_at, + ) + .await, + shared, + "two records with identical content must produce one block" + ); + assert_user_blocks_matches_repo(user_id, "creating two identical records").await; + + delete_record(&did, &jwt, &dropped).await; + + assert_user_blocks_matches_repo(user_id, "deleting one of two identical records").await; + + if is_store_backend() { + assert_eq!( + refcount_of(&shared).await, + 1, + "deleting one of two records sharing a block must drop exactly one reference" + ); + } + + delete_record(&did, &jwt, &kept).await; + assert_user_blocks_matches_repo(user_id, "deleting the second of two identical records").await; + + if is_store_backend() { + assert_eq!( + refcount_of(&shared).await, + 0, + "the shared block must reach refcount 0 once the last record referencing it is gone" + ); + } +} + +#[tokio::test] +async fn deleting_two_identical_records_in_one_commit_drops_both_references() { + let (did, jwt) = setup_new_user("user-blocks-shared-leaf-batch").await; + let did = Did::new(did).expect("setup_new_user returned a valid DID"); + let user_id = user_id_for(&did).await; + + let now_ms = Utc::now().timestamp_millis(); + let first = Rkey::new(format!("batchshared_{}_first", now_ms)).expect("valid rkey"); + let second = Rkey::new(format!("batchshared_{}_second", now_ms)).expect("valid rkey"); + let created_at = "2026-01-01T00:00:00Z".to_string(); + + let shared = create_record_at( + &did, + &jwt, + &first, + "identical content".to_string(), + created_at.clone(), + ) + .await; + create_record_at( + &did, + &jwt, + &second, + "identical content".to_string(), + created_at, + ) + .await; + assert_user_blocks_matches_repo(user_id, "creating two identical records").await; + + let res = client() + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&jwt) + .json(&json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#delete", + "collection": &*COLLECTION, + "rkey": first, + }, + { + "$type": "com.atproto.repo.applyWrites#delete", + "collection": &*COLLECTION, + "rkey": second, + }, + ] + })) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!( + res.status(), + StatusCode::OK, + "applyWrites didn't return 200: {:?}", + res.text().await + ); + + assert_user_blocks_matches_repo(user_id, "deleting both identical records in one commit").await; + + if is_store_backend() { + assert_eq!( + refcount_of(&shared).await, + 0, + "one commit dropping both references to a block must decrement it twice" + ); + } +}