tests: assert user_blocks matches reachable set after every write

Lewis: May this revision serve well! <lu5a@proton.me>
This commit is contained in:
Lewis
2026-07-25 08:27:40 +03:00
committed by Tangled
parent 1e2311f8fc
commit c3a8240154
3 changed files with 469 additions and 175 deletions
+19 -175
View File
@@ -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<String> = (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))
+63
View File
@@ -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<cid::Cid> {
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<cid::Cid> {
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<String> = reachable
.difference(&recorded)
.map(cid::Cid::to_string)
.collect();
let stale: Vec<String> = 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<Vec<u8>> {
let repos = super::common::get_test_repos().await;
@@ -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<Nsid> =
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<Rkey> = (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<Rkey> = (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"
);
}
}