Our Personal Data Server from scratch!
0

Configure Feed

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

tests: assert user_blocks matches reachable set after every write

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

author
Lewis
date (Jul 24, 2026, 8:45 PM +0300) commit 84df045c parent 57ba4fc1 change-id kpkvsnoq
+469 -175
+19 -175
crates/tranquil-pds/tests/gc_after_delete.rs
··· 5 5 use helpers::*; 6 6 use reqwest::StatusCode; 7 7 use serde_json::{Value, json}; 8 - use tranquil_types::Did; 8 + use tranquil_types::{Did, Nsid, Rkey}; 9 9 10 10 #[tokio::test] 11 11 async fn test_delete_record_marks_blocks_obsolete() { ··· 13 13 let base = base_url().await; 14 14 let repos = get_test_repos().await; 15 15 let (did, jwt) = setup_new_user("gc-after-delete").await; 16 + let did = Did::new(did).expect("setup_new_user returned a valid DID"); 16 17 17 18 let user_id = repos 18 19 .user 19 - .get_id_by_did(&Did::new(did.clone()).unwrap()) 20 + .get_id_by_did(&did) 20 21 .await 21 22 .expect("DB error") 22 23 .expect("User not found"); 23 24 24 - let count_baseline = repos 25 - .repo 26 - .count_user_blocks(user_id) 27 - .await 28 - .expect("count_user_blocks failed"); 29 - 30 - let collection = "app.bsky.feed.post"; 31 - let rkey = format!("gc_test_{}", Utc::now().timestamp_millis()); 25 + let collection = Nsid::new("app.bsky.feed.post".to_string()).expect("valid NSID"); 26 + let rkey = Rkey::new(format!("gc_test_{}", Utc::now().timestamp_millis())).expect("valid rkey"); 32 27 let create_payload = json!({ 33 28 "repo": did, 34 29 "collection": collection, ··· 65 60 .expect("createRecord response missing cid") 66 61 .to_string(); 67 62 68 - let count_after_create = repos 69 - .repo 70 - .count_user_blocks(user_id) 71 - .await 72 - .expect("count_user_blocks failed"); 73 - assert!( 74 - count_after_create > count_baseline, 75 - "user_blocks count did not grow after createRecord (baseline={}, after_create={})", 76 - count_baseline, 77 - count_after_create 78 - ); 63 + assert_user_blocks_matches_repo(user_id, "createRecord").await; 79 64 80 65 let delete_payload = json!({ 81 66 "repo": did, ··· 96 81 delete_res.text().await 97 82 ); 98 83 99 - let count_after_delete = repos 100 - .repo 101 - .count_user_blocks(user_id) 102 - .await 103 - .expect("count_user_blocks failed"); 104 - 105 - assert!( 106 - count_after_delete < count_after_create, 107 - "user_blocks count did not shrink after deleteRecord \ 108 - (baseline={}, after_create={}, after_delete={}). \ 109 - The delete path produced no obsolete CIDs beyond the prior commit root, \ 110 - which is the regression this test guards against.", 111 - count_baseline, 112 - count_after_create, 113 - count_after_delete 114 - ); 84 + assert_user_blocks_matches_repo(user_id, "deleteRecord").await; 115 85 116 86 let get_res = client 117 87 .get(format!("{}/xrpc/com.atproto.repo.getRecord", base)) 118 88 .query(&[ 119 89 ("repo", did.as_str()), 120 - ("collection", collection), 90 + ("collection", collection.as_str()), 121 91 ("rkey", rkey.as_str()), 122 92 ]) 123 93 .send() ··· 138 108 let base = base_url().await; 139 109 let repos = get_test_repos().await; 140 110 let (did, jwt) = setup_new_user("gc-after-update").await; 111 + let did = Did::new(did).expect("setup_new_user returned a valid DID"); 141 112 142 113 let user_id = repos 143 114 .user 144 - .get_id_by_did(&Did::new(did.clone()).unwrap()) 115 + .get_id_by_did(&did) 145 116 .await 146 117 .expect("DB error") 147 118 .expect("User not found"); 148 119 149 - let collection = "app.bsky.feed.post"; 150 - let rkey = format!("gc_update_{}", Utc::now().timestamp_millis()); 120 + let collection = Nsid::new("app.bsky.feed.post".to_string()).expect("valid NSID"); 121 + let rkey = 122 + Rkey::new(format!("gc_update_{}", Utc::now().timestamp_millis())).expect("valid rkey"); 151 123 152 124 let put_v1 = json!({ 153 125 "repo": did, ··· 168 140 .expect("Failed to send putRecord v1"); 169 141 assert_eq!(res.status(), StatusCode::OK, "first putRecord failed"); 170 142 171 - let count_after_create = repos 172 - .repo 173 - .count_user_blocks(user_id) 174 - .await 175 - .expect("count_user_blocks failed"); 143 + assert_user_blocks_matches_repo(user_id, "the first putRecord").await; 176 144 177 145 let put_v2 = json!({ 178 146 "repo": did, ··· 193 161 .expect("Failed to send putRecord v2"); 194 162 assert_eq!(res.status(), StatusCode::OK, "second putRecord failed"); 195 163 196 - let count_after_update = repos 197 - .repo 198 - .count_user_blocks(user_id) 199 - .await 200 - .expect("count_user_blocks failed"); 201 - 202 - assert!( 203 - count_after_update <= count_after_create + 1, 204 - "user_blocks count grew by more than 1 after putRecord update \ 205 - (after_create={}, after_update={}). The previous version's record block \ 206 - should have been marked obsolete; instead it appears to be leaking.", 207 - count_after_create, 208 - count_after_update 209 - ); 210 - } 211 - 212 - #[tokio::test] 213 - async fn test_delete_in_populated_repo_marks_merged_subtree_blocks_obsolete() { 214 - let client = client(); 215 - let base = base_url().await; 216 - let repos = get_test_repos().await; 217 - let (did, jwt) = setup_new_user("gc-merge").await; 218 - 219 - let user_id = repos 220 - .user 221 - .get_id_by_did(&Did::new(did.clone()).unwrap()) 222 - .await 223 - .expect("DB error") 224 - .expect("User not found"); 225 - 226 - let collection = "app.bsky.feed.post"; 227 - let record_count = 64usize; 228 - let now_ms = Utc::now().timestamp_millis(); 229 - 230 - let rkeys: Vec<String> = (0..record_count) 231 - .map(|i| format!("gc_merge_{}_{:04}", now_ms, i)) 232 - .collect(); 233 - 234 - let create_results = 235 - futures::future::try_join_all(rkeys.iter().enumerate().map(|(i, rkey)| { 236 - let client = client.clone(); 237 - let jwt = jwt.clone(); 238 - let did = did.clone(); 239 - let base = base.to_string(); 240 - let payload = json!({ 241 - "repo": did, 242 - "collection": collection, 243 - "rkey": rkey, 244 - "record": { 245 - "$type": collection, 246 - "text": format!("seed record {}", i), 247 - "createdAt": Utc::now().to_rfc3339() 248 - } 249 - }); 250 - async move { 251 - let res = client 252 - .post(format!("{}/xrpc/com.atproto.repo.createRecord", base)) 253 - .bearer_auth(&jwt) 254 - .json(&payload) 255 - .send() 256 - .await 257 - .expect("Failed to send createRecord"); 258 - if res.status() != StatusCode::OK { 259 - return Err(format!("seed createRecord failed: {}", res.status())); 260 - } 261 - Ok::<(), String>(()) 262 - } 263 - })) 264 - .await; 265 - create_results.expect("seeding records failed"); 266 - 267 - let count_after_seed = repos 268 - .repo 269 - .count_user_blocks(user_id) 270 - .await 271 - .expect("count_user_blocks failed"); 272 - 273 - let target_rkey = &rkeys[record_count / 2]; 274 - let delete_payload = json!({ 275 - "repo": did, 276 - "collection": collection, 277 - "rkey": target_rkey, 278 - }); 279 - let delete_res = client 280 - .post(format!("{}/xrpc/com.atproto.repo.deleteRecord", base)) 281 - .bearer_auth(&jwt) 282 - .json(&delete_payload) 283 - .send() 284 - .await 285 - .expect("Failed to send deleteRecord"); 286 - assert_eq!( 287 - delete_res.status(), 288 - StatusCode::OK, 289 - "deleteRecord did not return 200: {:?}", 290 - delete_res.text().await 291 - ); 292 - 293 - let count_after_delete = repos 294 - .repo 295 - .count_user_blocks(user_id) 296 - .await 297 - .expect("count_user_blocks failed"); 298 - assert!( 299 - count_after_delete < count_after_seed, 300 - "user_blocks did not shrink after deleting from a populated repo \ 301 - (after_seed={}, after_delete={}). The path-walk-based obsolete \ 302 - calculation does not capture sibling subtree blocks orphaned by \ 303 - delete-merge; only an MST-diff-based calculation does.", 304 - count_after_seed, 305 - count_after_delete 306 - ); 307 - 308 - let get_res = client 309 - .get(format!("{}/xrpc/com.atproto.repo.getRecord", base)) 310 - .query(&[ 311 - ("repo", did.as_str()), 312 - ("collection", collection), 313 - ("rkey", target_rkey.as_str()), 314 - ]) 315 - .send() 316 - .await 317 - .expect("Failed to send getRecord"); 318 - assert!( 319 - !get_res.status().is_success(), 320 - "deleted record is still resolvable via getRecord (status={})", 321 - get_res.status(), 322 - ); 164 + assert_user_blocks_matches_repo(user_id, "the second putRecord").await; 323 165 } 324 166 325 167 #[tokio::test] ··· 339 181 .as_tranquil_store() 340 182 .expect("tranquil-store backend selected but block_store is not TranquilStore"); 341 183 let (did, jwt) = setup_new_user("gc-store-decrement").await; 184 + let did = Did::new(did).expect("setup_new_user returned a valid DID"); 342 185 343 - let collection = "app.bsky.feed.post"; 344 - let rkey = format!("gc_store_{}", Utc::now().timestamp_millis()); 186 + let collection = Nsid::new("app.bsky.feed.post".to_string()).expect("valid NSID"); 187 + let rkey = 188 + Rkey::new(format!("gc_store_{}", Utc::now().timestamp_millis())).expect("valid rkey"); 345 189 346 190 let create_res = client 347 191 .post(format!("{}/xrpc/com.atproto.repo.createRecord", base))
+63
crates/tranquil-pds/tests/helpers/mod.rs
··· 480 480 multibase::encode(multibase::Base::Base58Btc, buf) 481 481 } 482 482 483 + async fn reachable_blocks(user_id: uuid::Uuid) -> std::collections::BTreeSet<cid::Cid> { 484 + use jacquard_repo::storage::BlockStore; 485 + 486 + let repos = super::common::get_test_repos().await; 487 + let store = super::common::get_test_block_store().await; 488 + 489 + let root_str = repos 490 + .repo 491 + .get_repo_root_cid_by_user_id(user_id) 492 + .await 493 + .expect("DB error fetching repo root") 494 + .expect("repo root not found"); 495 + let root_cid = cid::Cid::try_from(root_str.as_str()).expect("repo root isn't a valid CID"); 496 + let commit_bytes = store 497 + .get(&root_cid) 498 + .await 499 + .expect("block store error fetching commit") 500 + .expect("commit block not in block store"); 501 + let data_cid = jacquard_repo::commit::Commit::from_cbor(&commit_bytes) 502 + .expect("commit block doesn't parse") 503 + .data; 504 + 505 + let mst = jacquard_repo::mst::Mst::load(std::sync::Arc::new(store.clone()), data_cid, None); 506 + let mut cids = tranquil_pds::repo_ops::reachable_tree_cids(&mst) 507 + .await 508 + .expect("walking the new MST failed"); 509 + cids.insert(root_cid); 510 + cids 511 + } 512 + 513 + async fn recorded_blocks(user_id: uuid::Uuid) -> std::collections::BTreeSet<cid::Cid> { 514 + super::common::get_test_repos() 515 + .await 516 + .repo 517 + .get_user_block_cids_since_rev(user_id, None) 518 + .await 519 + .expect("get_user_block_cids_since_rev failed") 520 + .iter() 521 + .map(|b| cid::Cid::try_from(b.as_slice()).expect("invalid CID in user_blocks")) 522 + .collect() 523 + } 524 + 525 + #[allow(dead_code)] 526 + pub async fn assert_user_blocks_matches_repo(user_id: uuid::Uuid, phase: &str) { 527 + let reachable = reachable_blocks(user_id).await; 528 + let recorded = recorded_blocks(user_id).await; 529 + let missing: Vec<String> = reachable 530 + .difference(&recorded) 531 + .map(cid::Cid::to_string) 532 + .collect(); 533 + let stale: Vec<String> = recorded 534 + .difference(&reachable) 535 + .map(cid::Cid::to_string) 536 + .collect(); 537 + assert!( 538 + missing.is_empty() && stale.is_empty(), 539 + "user_blocks doesn't match the blocks reachable from the repo root after {phase}. \ 540 + reachable={} recorded={} missing={missing:?} stale={stale:?}", 541 + reachable.len(), 542 + recorded.len(), 543 + ); 544 + } 545 + 483 546 #[allow(dead_code)] 484 547 pub async fn get_user_signing_key(did: &str) -> Option<Vec<u8>> { 485 548 let repos = super::common::get_test_repos().await;
+387
crates/tranquil-pds/tests/user_blocks_reachability.rs
··· 1 + mod common; 2 + mod helpers; 3 + use chrono::Utc; 4 + use common::*; 5 + use helpers::*; 6 + use reqwest::StatusCode; 7 + use serde_json::json; 8 + use std::sync::LazyLock; 9 + use tranquil_types::{Did, Nsid, Rkey}; 10 + 11 + static COLLECTION: LazyLock<Nsid> = 12 + LazyLock::new(|| Nsid::new("app.bsky.feed.post".to_string()).expect("valid NSID")); 13 + 14 + async fn create_record(did: &Did, jwt: &str, rkey: &Rkey, text: String) -> cid::Cid { 15 + create_record_at(did, jwt, rkey, text, Utc::now().to_rfc3339()).await 16 + } 17 + 18 + async fn create_record_at( 19 + did: &Did, 20 + jwt: &str, 21 + rkey: &Rkey, 22 + text: String, 23 + created_at: String, 24 + ) -> cid::Cid { 25 + let res = client() 26 + .post(format!( 27 + "{}/xrpc/com.atproto.repo.createRecord", 28 + base_url().await 29 + )) 30 + .bearer_auth(jwt) 31 + .json(&json!({ 32 + "repo": did, 33 + "collection": &*COLLECTION, 34 + "rkey": rkey, 35 + "record": { 36 + "$type": &*COLLECTION, 37 + "text": text, 38 + "createdAt": created_at 39 + } 40 + })) 41 + .send() 42 + .await 43 + .expect("Failed to send createRecord"); 44 + assert_eq!( 45 + res.status(), 46 + StatusCode::OK, 47 + "createRecord for {rkey} didn't return 200: {:?}", 48 + res.text().await 49 + ); 50 + let body: serde_json::Value = res.json().await.expect("createRecord response isn't JSON"); 51 + let cid_str = body["cid"] 52 + .as_str() 53 + .expect("createRecord response missing cid"); 54 + cid::Cid::try_from(cid_str).expect("createRecord returned an invalid cid") 55 + } 56 + 57 + async fn refcount_of(cid: &cid::Cid) -> u32 { 58 + get_test_block_store() 59 + .await 60 + .as_tranquil_store() 61 + .expect("tranquil-store backend selected but block_store isn't TranquilStore") 62 + .refcount_of(cid) 63 + .expect("refcount_of failed") 64 + .unwrap_or(0) 65 + } 66 + 67 + async fn delete_record(did: &Did, jwt: &str, rkey: &Rkey) { 68 + let res = client() 69 + .post(format!( 70 + "{}/xrpc/com.atproto.repo.deleteRecord", 71 + base_url().await 72 + )) 73 + .bearer_auth(jwt) 74 + .json(&json!({ 75 + "repo": did, 76 + "collection": &*COLLECTION, 77 + "rkey": rkey, 78 + })) 79 + .send() 80 + .await 81 + .expect("Failed to send deleteRecord"); 82 + assert_eq!( 83 + res.status(), 84 + StatusCode::OK, 85 + "deleteRecord for {rkey} didn't return 200: {:?}", 86 + res.text().await 87 + ); 88 + } 89 + 90 + async fn assert_record_gone(did: &Did, rkey: &Rkey) { 91 + let res = client() 92 + .get(format!( 93 + "{}/xrpc/com.atproto.repo.getRecord", 94 + base_url().await 95 + )) 96 + .query(&[ 97 + ("repo", did.as_str()), 98 + ("collection", COLLECTION.as_str()), 99 + ("rkey", rkey.as_str()), 100 + ]) 101 + .send() 102 + .await 103 + .expect("Failed to send getRecord"); 104 + assert!( 105 + !res.status().is_success(), 106 + "deleted record {rkey} is still resolvable via getRecord: {}", 107 + res.status() 108 + ); 109 + } 110 + 111 + async fn user_id_for(did: &Did) -> uuid::Uuid { 112 + get_test_repos() 113 + .await 114 + .user 115 + .get_id_by_did(did) 116 + .await 117 + .expect("DB error looking up the user id") 118 + .expect("User not found") 119 + } 120 + 121 + #[tokio::test] 122 + async fn deleting_from_a_populated_repo_keeps_user_blocks_equal_to_reachable_set() { 123 + let (did, jwt) = setup_new_user("user-blocks-reachability").await; 124 + let did = Did::new(did).expect("setup_new_user returned a valid DID"); 125 + let user_id = user_id_for(&did).await; 126 + 127 + let record_count = 64usize; 128 + let now_ms = Utc::now().timestamp_millis(); 129 + let rkeys: Vec<Rkey> = (0..record_count) 130 + .map(|i| Rkey::new(format!("reach_{}_{:04}", now_ms, i)).expect("valid rkey")) 131 + .collect(); 132 + 133 + futures::future::join_all( 134 + rkeys 135 + .iter() 136 + .enumerate() 137 + .map(|(i, rkey)| create_record(&did, &jwt, rkey, format!("seed record {}", i))), 138 + ) 139 + .await; 140 + 141 + assert_user_blocks_matches_repo(user_id, "64 creates").await; 142 + 143 + let target = &rkeys[record_count / 2]; 144 + delete_record(&did, &jwt, target).await; 145 + 146 + assert_user_blocks_matches_repo(user_id, "deleting one record").await; 147 + assert_record_gone(&did, target).await; 148 + } 149 + 150 + #[tokio::test] 151 + async fn applying_a_multi_op_batch_keeps_user_blocks_equal_to_reachable_set() { 152 + let (did, jwt) = setup_new_user("user-blocks-apply-writes").await; 153 + let did = Did::new(did).expect("setup_new_user returned a valid DID"); 154 + let user_id = user_id_for(&did).await; 155 + 156 + let now_ms = Utc::now().timestamp_millis(); 157 + let rkeys: Vec<Rkey> = (0..3) 158 + .map(|i| Rkey::new(format!("batch_{}_{}", now_ms, i)).expect("valid rkey")) 159 + .collect(); 160 + let created_at = "2026-01-01T00:00:00Z".to_string(); 161 + 162 + futures::future::join_all( 163 + rkeys.iter().map(|rkey| { 164 + create_record_at(&did, &jwt, rkey, "shared".to_string(), created_at.clone()) 165 + }), 166 + ) 167 + .await; 168 + assert_user_blocks_matches_repo(user_id, "three records sharing one leaf block").await; 169 + 170 + let writes = json!({ 171 + "repo": did, 172 + "writes": [ 173 + { 174 + "$type": "com.atproto.repo.applyWrites#delete", 175 + "collection": &*COLLECTION, 176 + "rkey": rkeys[0], 177 + }, 178 + { 179 + "$type": "com.atproto.repo.applyWrites#update", 180 + "collection": &*COLLECTION, 181 + "rkey": rkeys[1], 182 + "value": { 183 + "$type": &*COLLECTION, 184 + "text": "updated in the same commit", 185 + "createdAt": created_at 186 + }, 187 + }, 188 + { 189 + "$type": "com.atproto.repo.applyWrites#create", 190 + "collection": &*COLLECTION, 191 + "rkey": Rkey::new(format!("batch_{}_new", now_ms)).expect("valid rkey"), 192 + "value": { 193 + "$type": &*COLLECTION, 194 + "text": "created in the same commit", 195 + "createdAt": created_at 196 + }, 197 + }, 198 + ] 199 + }); 200 + let res = client() 201 + .post(format!( 202 + "{}/xrpc/com.atproto.repo.applyWrites", 203 + base_url().await 204 + )) 205 + .bearer_auth(&jwt) 206 + .json(&writes) 207 + .send() 208 + .await 209 + .expect("Failed to send applyWrites"); 210 + assert_eq!( 211 + res.status(), 212 + StatusCode::OK, 213 + "applyWrites didn't return 200: {:?}", 214 + res.text().await 215 + ); 216 + 217 + assert_user_blocks_matches_repo(user_id, "a delete, update, & create in one commit").await; 218 + assert_record_gone(&did, &rkeys[0]).await; 219 + } 220 + 221 + #[tokio::test] 222 + async fn emptying_the_repo_keeps_user_blocks_equal_to_reachable_set() { 223 + let (did, jwt) = setup_new_user("user-blocks-empty-tree").await; 224 + let did = Did::new(did).expect("setup_new_user returned a valid DID"); 225 + let user_id = user_id_for(&did).await; 226 + 227 + let now_ms = Utc::now().timestamp_millis(); 228 + let only = Rkey::new(format!("empty_{}_only", now_ms)).expect("valid rkey"); 229 + 230 + create_record(&did, &jwt, &only, "the only record".to_string()).await; 231 + assert_user_blocks_matches_repo(user_id, "creating the only record").await; 232 + 233 + delete_record(&did, &jwt, &only).await; 234 + assert_user_blocks_matches_repo(user_id, "deleting the last record").await; 235 + assert_record_gone(&did, &only).await; 236 + } 237 + 238 + #[tokio::test] 239 + async fn reverting_the_tree_to_a_stored_shape_still_records_its_blocks() { 240 + let (did, jwt) = setup_new_user("user-blocks-resurrect").await; 241 + let did = Did::new(did).expect("setup_new_user returned a valid DID"); 242 + let user_id = user_id_for(&did).await; 243 + 244 + let now_ms = Utc::now().timestamp_millis(); 245 + let kept = Rkey::new(format!("resurrect_{}_kept", now_ms)).expect("valid rkey"); 246 + let churned = Rkey::new(format!("resurrect_{}_churned", now_ms)).expect("valid rkey"); 247 + 248 + create_record(&did, &jwt, &kept, "first record".to_string()).await; 249 + assert_user_blocks_matches_repo(user_id, "creating the first record").await; 250 + 251 + create_record(&did, &jwt, &churned, "second record".to_string()).await; 252 + assert_user_blocks_matches_repo(user_id, "creating the second record").await; 253 + 254 + delete_record(&did, &jwt, &churned).await; 255 + assert_user_blocks_matches_repo(user_id, "deleting back to the one-record tree").await; 256 + 257 + create_record(&did, &jwt, &churned, "second record".to_string()).await; 258 + assert_user_blocks_matches_repo(user_id, "recreating the deleted record").await; 259 + } 260 + 261 + #[tokio::test] 262 + async fn deleting_one_of_two_identical_records_keeps_the_shared_block() { 263 + let (did, jwt) = setup_new_user("user-blocks-shared-leaf").await; 264 + let did = Did::new(did).expect("setup_new_user returned a valid DID"); 265 + let user_id = user_id_for(&did).await; 266 + 267 + let now_ms = Utc::now().timestamp_millis(); 268 + let kept = Rkey::new(format!("shared_{}_kept", now_ms)).expect("valid rkey"); 269 + let dropped = Rkey::new(format!("shared_{}_dropped", now_ms)).expect("valid rkey"); 270 + let created_at = "2026-01-01T00:00:00Z".to_string(); 271 + 272 + let shared = create_record_at( 273 + &did, 274 + &jwt, 275 + &kept, 276 + "identical content".to_string(), 277 + created_at.clone(), 278 + ) 279 + .await; 280 + assert_eq!( 281 + create_record_at( 282 + &did, 283 + &jwt, 284 + &dropped, 285 + "identical content".to_string(), 286 + created_at, 287 + ) 288 + .await, 289 + shared, 290 + "two records with identical content must produce one block" 291 + ); 292 + assert_user_blocks_matches_repo(user_id, "creating two identical records").await; 293 + 294 + delete_record(&did, &jwt, &dropped).await; 295 + 296 + assert_user_blocks_matches_repo(user_id, "deleting one of two identical records").await; 297 + 298 + if is_store_backend() { 299 + assert_eq!( 300 + refcount_of(&shared).await, 301 + 1, 302 + "deleting one of two records sharing a block must drop exactly one reference" 303 + ); 304 + } 305 + 306 + delete_record(&did, &jwt, &kept).await; 307 + assert_user_blocks_matches_repo(user_id, "deleting the second of two identical records").await; 308 + 309 + if is_store_backend() { 310 + assert_eq!( 311 + refcount_of(&shared).await, 312 + 0, 313 + "the shared block must reach refcount 0 once the last record referencing it is gone" 314 + ); 315 + } 316 + } 317 + 318 + #[tokio::test] 319 + async fn deleting_two_identical_records_in_one_commit_drops_both_references() { 320 + let (did, jwt) = setup_new_user("user-blocks-shared-leaf-batch").await; 321 + let did = Did::new(did).expect("setup_new_user returned a valid DID"); 322 + let user_id = user_id_for(&did).await; 323 + 324 + let now_ms = Utc::now().timestamp_millis(); 325 + let first = Rkey::new(format!("batchshared_{}_first", now_ms)).expect("valid rkey"); 326 + let second = Rkey::new(format!("batchshared_{}_second", now_ms)).expect("valid rkey"); 327 + let created_at = "2026-01-01T00:00:00Z".to_string(); 328 + 329 + let shared = create_record_at( 330 + &did, 331 + &jwt, 332 + &first, 333 + "identical content".to_string(), 334 + created_at.clone(), 335 + ) 336 + .await; 337 + create_record_at( 338 + &did, 339 + &jwt, 340 + &second, 341 + "identical content".to_string(), 342 + created_at, 343 + ) 344 + .await; 345 + assert_user_blocks_matches_repo(user_id, "creating two identical records").await; 346 + 347 + let res = client() 348 + .post(format!( 349 + "{}/xrpc/com.atproto.repo.applyWrites", 350 + base_url().await 351 + )) 352 + .bearer_auth(&jwt) 353 + .json(&json!({ 354 + "repo": did, 355 + "writes": [ 356 + { 357 + "$type": "com.atproto.repo.applyWrites#delete", 358 + "collection": &*COLLECTION, 359 + "rkey": first, 360 + }, 361 + { 362 + "$type": "com.atproto.repo.applyWrites#delete", 363 + "collection": &*COLLECTION, 364 + "rkey": second, 365 + }, 366 + ] 367 + })) 368 + .send() 369 + .await 370 + .expect("Failed to send applyWrites"); 371 + assert_eq!( 372 + res.status(), 373 + StatusCode::OK, 374 + "applyWrites didn't return 200: {:?}", 375 + res.text().await 376 + ); 377 + 378 + assert_user_blocks_matches_repo(user_id, "deleting both identical records in one commit").await; 379 + 380 + if is_store_backend() { 381 + assert_eq!( 382 + refcount_of(&shared).await, 383 + 0, 384 + "one commit dropping both references to a block must decrement it twice" 385 + ); 386 + } 387 + }