From af3821514fba3b197826e8e3b3a99dabded058c0 Mon Sep 17 00:00:00 2001 From: Lewis Date: Mon, 27 Apr 2026 11:34:42 +0300 Subject: [PATCH] test(tranquil-pds): same-rkey batch coverage and inductive inverse for in-batch dups Lewis: May this revision serve well! --- crates/tranquil-pds/tests/lifecycle_record.rs | 2 +- .../tests/mst_inductive_firehose.rs | 59 +- crates/tranquil-pds/tests/repo_batch.rs | 1020 +++++++++++++++++ 3 files changed, 1079 insertions(+), 2 deletions(-) diff --git a/crates/tranquil-pds/tests/lifecycle_record.rs b/crates/tranquil-pds/tests/lifecycle_record.rs index 5659c31..3077855 100644 --- a/crates/tranquil-pds/tests/lifecycle_record.rs +++ b/crates/tranquil-pds/tests/lifecycle_record.rs @@ -456,7 +456,7 @@ async fn test_apply_writes_batch() { "writes": [ { "$type": "com.atproto.repo.applyWrites#create", "collection": "app.bsky.feed.post", "rkey": "batch-post-1", "value": { "$type": "app.bsky.feed.post", "text": "First batch post", "createdAt": now } }, { "$type": "com.atproto.repo.applyWrites#create", "collection": "app.bsky.feed.post", "rkey": "batch-post-2", "value": { "$type": "app.bsky.feed.post", "text": "Second batch post", "createdAt": now } }, - { "$type": "com.atproto.repo.applyWrites#create", "collection": "app.bsky.actor.profile", "rkey": "self", "value": { "$type": "app.bsky.actor.profile", "displayName": "Batch User" } } + { "$type": "com.atproto.repo.applyWrites#update", "collection": "app.bsky.actor.profile", "rkey": "self", "value": { "$type": "app.bsky.actor.profile", "displayName": "Batch User" } } ] }); let apply_res = client diff --git a/crates/tranquil-pds/tests/mst_inductive_firehose.rs b/crates/tranquil-pds/tests/mst_inductive_firehose.rs index 713ce16..3b87ffd 100644 --- a/crates/tranquil-pds/tests/mst_inductive_firehose.rs +++ b/crates/tranquil-pds/tests/mst_inductive_firehose.rs @@ -124,7 +124,7 @@ async fn verify_inductive_inverse(event: &SequencedEvent) -> Result<(Cid, Cid), let new_data_cid = new_commit_data_cid(&storage, &commit_cid).await?; let mut mst = Mst::load(storage.clone(), new_data_cid, None); - for op_value in ops_json(event)? { + for op_value in ops_json(event)?.iter().rev() { let verified = parse_op_to_verified(op_value)?; let inverted = mst .invert_op(verified.clone()) @@ -546,6 +546,63 @@ async fn inductive_inverse_verifies_every_commit() { report_failures(non_genesis.len(), &failures, "any inverse"); } +#[tokio::test] +async fn inductive_inverse_handles_same_rkey_in_batch() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + + let now = chrono::Utc::now().to_rfc3339(); + let rkey = rkey_for("dup", 0); + create_record(&client, &token, &did, COLLECTION, &rkey).await; + + let writes = vec![ + json!({ + "$type": "com.atproto.repo.applyWrites#update", + "collection": COLLECTION, + "rkey": rkey, + "value": { + "$type": COLLECTION, + "text": "v1", + "createdAt": now, + } + }), + json!({ + "$type": "com.atproto.repo.applyWrites#update", + "collection": COLLECTION, + "rkey": rkey, + "value": { + "$type": COLLECTION, + "text": "v2", + "createdAt": now, + } + }), + ]; + apply_writes_batch(&client, &token, &did, writes).await; + + let our = our_commit_events(&did).await; + let dup_event = our + .iter() + .find(|e| { + ops_json(e) + .map(|arr| { + arr.iter() + .filter(|op| op["action"].as_str() == Some("update")) + .count() + == 2 + }) + .unwrap_or(false) + }) + .expect("commit event with two same-rkey updates"); + + let (exp, got) = verify_inductive_inverse(dup_event) + .await + .expect("inverse verify should succeed for same-rkey batch"); + assert_eq!( + exp, got, + "inverse root mismatch for same-rkey batch: exp={exp} got={got}" + ); +} + #[tokio::test] async fn prev_cid_chain_walks_to_genesis() { let client = client(); diff --git a/crates/tranquil-pds/tests/repo_batch.rs b/crates/tranquil-pds/tests/repo_batch.rs index 034cea3..787d8d9 100644 --- a/crates/tranquil-pds/tests/repo_batch.rs +++ b/crates/tranquil-pds/tests/repo_batch.rs @@ -3,6 +3,8 @@ use chrono::Utc; use common::*; use reqwest::StatusCode; use serde_json::{Value, json}; +use tranquil_db_traits::{Backlink, BacklinkPath}; +use tranquil_pds::types::{AtUri, Did, Nsid}; #[tokio::test] async fn test_apply_writes_create() { @@ -312,3 +314,1021 @@ async fn test_apply_writes_empty_writes() { .expect("Failed to send request"); assert_eq!(res.status(), StatusCode::BAD_REQUEST); } + +#[tokio::test] +async fn test_apply_writes_delete_then_create_same_rkey() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let rkey = format!("recreate_{}", Utc::now().timestamp_millis()); + + let create_payload = json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": rkey, + "record": { + "$type": "app.bsky.feed.post", + "text": "original", + "createdAt": now + } + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.putRecord", + base_url().await + )) + .bearer_auth(&token) + .json(&create_payload) + .send() + .await + .expect("Failed to create"); + assert_eq!(res.status(), StatusCode::OK); + + let recreate_payload = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#delete", + "collection": "app.bsky.feed.post", + "rkey": rkey + }, + { + "$type": "com.atproto.repo.applyWrites#create", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "recreated", + "createdAt": now + } + } + ] + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&token) + .json(&recreate_payload) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!(res.status(), StatusCode::OK); + + let get_res = client + .get(format!( + "{}/xrpc/com.atproto.repo.getRecord", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", rkey.as_str()), + ]) + .send() + .await + .expect("Failed to fetch record"); + assert_eq!( + get_res.status(), + StatusCode::OK, + "repo.getRecord must return the recreated record, not 404" + ); + let body: Value = get_res.json().await.expect("Response was not valid JSON"); + assert_eq!(body["value"]["text"], json!("recreated")); + + let list_res = client + .get(format!( + "{}/xrpc/com.atproto.repo.listRecords", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ]) + .send() + .await + .expect("Failed to list records"); + assert_eq!(list_res.status(), StatusCode::OK); + let list_body: Value = list_res.json().await.expect("listRecords not valid JSON"); + let records = list_body["records"] + .as_array() + .expect("records must be array"); + let expected_uri = format!("at://{}/app.bsky.feed.post/{}", did, rkey); + assert!( + records.iter().any(|r| r["uri"] == json!(expected_uri)), + "listRecords must include the recreated rkey" + ); +} + +#[tokio::test] +async fn test_apply_writes_create_then_delete_same_rkey() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let rkey = format!("transient_{}", Utc::now().timestamp_millis()); + + let payload = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#create", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "transient", + "createdAt": now + } + }, + { + "$type": "com.atproto.repo.applyWrites#delete", + "collection": "app.bsky.feed.post", + "rkey": rkey + } + ] + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&token) + .json(&payload) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!(res.status(), StatusCode::OK); + + let get_res = client + .get(format!( + "{}/xrpc/com.atproto.repo.getRecord", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", rkey.as_str()), + ]) + .send() + .await + .expect("Failed to fetch record"); + assert_eq!( + get_res.status(), + StatusCode::NOT_FOUND, + "create+delete on same rkey must leave no record" + ); + + let list_res = client + .get(format!( + "{}/xrpc/com.atproto.repo.listRecords", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ]) + .send() + .await + .expect("Failed to list records"); + assert_eq!(list_res.status(), StatusCode::OK); + let list_body: Value = list_res.json().await.expect("listRecords not valid JSON"); + let records = list_body["records"] + .as_array() + .expect("records must be array"); + let expected_uri = format!("at://{}/app.bsky.feed.post/{}", did, rkey); + assert!( + !records.iter().any(|r| r["uri"] == json!(expected_uri)), + "listRecords must not include a created-then-deleted rkey" + ); +} + +async fn repo_id_for_did(did: &str) -> uuid::Uuid { + let repos = get_test_repos().await; + let parsed = Did::new(did).expect("valid did"); + repos + .user + .get_id_by_did(&parsed) + .await + .expect("lookup user_id") + .expect("user exists") +} + +async fn follow_uris_pointing_to(repo_id: uuid::Uuid, target_did: &str) -> Vec { + let repos = get_test_repos().await; + let probe = Backlink { + uri: AtUri::from_parts("did:plc:probe", "app.bsky.graph.follow", "probe"), + path: BacklinkPath::Subject, + link_to: target_did.to_string(), + }; + let collection = Nsid::new("app.bsky.graph.follow").expect("valid nsid"); + repos + .backlink + .get_backlink_conflicts(repo_id, &collection, &[probe]) + .await + .expect("backlink query") + .into_iter() + .map(|u| u.as_str().to_string()) + .collect() +} + +#[tokio::test] +async fn test_apply_writes_update_then_update_same_rkey() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let rkey = format!("double_update_{}", Utc::now().timestamp_millis()); + + let create_payload = json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": rkey, + "record": { + "$type": "app.bsky.feed.post", + "text": "v0", + "createdAt": now + } + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.putRecord", + base_url().await + )) + .bearer_auth(&token) + .json(&create_payload) + .send() + .await + .expect("Failed to create"); + assert_eq!(res.status(), StatusCode::OK); + + let payload = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#update", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "v1", + "createdAt": now + } + }, + { + "$type": "com.atproto.repo.applyWrites#update", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "v2", + "createdAt": now + } + } + ] + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&token) + .json(&payload) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!(res.status(), StatusCode::OK); + let body: Value = res.json().await.expect("Response was not valid JSON"); + let final_cid = body["results"] + .as_array() + .and_then(|r| r.last()) + .and_then(|r| r["cid"].as_str()) + .expect("last result must carry a cid") + .to_string(); + + let get_res = client + .get(format!( + "{}/xrpc/com.atproto.repo.getRecord", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", rkey.as_str()), + ]) + .send() + .await + .expect("Failed to fetch record"); + assert_eq!(get_res.status(), StatusCode::OK); + let stored: Value = get_res.json().await.expect("Response was not valid JSON"); + assert_eq!(stored["value"]["text"], json!("v2")); + assert_eq!(stored["cid"], json!(final_cid)); +} + +#[tokio::test] +async fn test_apply_writes_update_then_delete_same_rkey() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let rkey = format!("update_delete_{}", Utc::now().timestamp_millis()); + + let create_payload = json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": rkey, + "record": { + "$type": "app.bsky.feed.post", + "text": "v0", + "createdAt": now + } + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.putRecord", + base_url().await + )) + .bearer_auth(&token) + .json(&create_payload) + .send() + .await + .expect("Failed to create"); + assert_eq!(res.status(), StatusCode::OK); + + let payload = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#update", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "v1", + "createdAt": now + } + }, + { + "$type": "com.atproto.repo.applyWrites#delete", + "collection": "app.bsky.feed.post", + "rkey": rkey + } + ] + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&token) + .json(&payload) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!(res.status(), StatusCode::OK); + + let get_res = client + .get(format!( + "{}/xrpc/com.atproto.repo.getRecord", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", rkey.as_str()), + ]) + .send() + .await + .expect("Failed to fetch record"); + assert_eq!(get_res.status(), StatusCode::NOT_FOUND); +} + +#[tokio::test] +async fn test_apply_writes_delete_then_update_same_rkey_rejected() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let rkey = format!("delete_update_{}", Utc::now().timestamp_millis()); + + let create_payload = json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": rkey, + "record": { + "$type": "app.bsky.feed.post", + "text": "v0", + "createdAt": now + } + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.putRecord", + base_url().await + )) + .bearer_auth(&token) + .json(&create_payload) + .send() + .await + .expect("Failed to seed record"); + assert_eq!(res.status(), StatusCode::OK); + + let payload = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#delete", + "collection": "app.bsky.feed.post", + "rkey": rkey + }, + { + "$type": "com.atproto.repo.applyWrites#update", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "v1", + "createdAt": now + } + } + ] + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&token) + .json(&payload) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!( + res.status(), + StatusCode::BAD_REQUEST, + "update of a record deleted earlier in the same batch must be rejected" + ); + + let get_res = client + .get(format!( + "{}/xrpc/com.atproto.repo.getRecord", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", rkey.as_str()), + ]) + .send() + .await + .expect("Failed to fetch record"); + assert_eq!(get_res.status(), StatusCode::OK); + let body: Value = get_res.json().await.expect("Response was not valid JSON"); + assert_eq!( + body["value"]["text"], + json!("v0"), + "rejected batch must leave the seed record untouched" + ); +} + +#[tokio::test] +async fn test_apply_writes_create_update_delete_same_rkey() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let rkey = format!("triple_{}", Utc::now().timestamp_millis()); + + let payload = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#create", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "v0", + "createdAt": now + } + }, + { + "$type": "com.atproto.repo.applyWrites#update", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "v1", + "createdAt": now + } + }, + { + "$type": "com.atproto.repo.applyWrites#delete", + "collection": "app.bsky.feed.post", + "rkey": rkey + } + ] + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&token) + .json(&payload) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!(res.status(), StatusCode::OK); + + let get_res = client + .get(format!( + "{}/xrpc/com.atproto.repo.getRecord", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", rkey.as_str()), + ]) + .send() + .await + .expect("Failed to fetch record"); + assert_eq!(get_res.status(), StatusCode::NOT_FOUND); +} + +#[tokio::test] +async fn test_apply_writes_distinct_rkeys_not_conflated() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let stamp = Utc::now().timestamp_millis(); + let rkey_a = format!("distinct_a_{}", stamp); + let rkey_b = format!("distinct_b_{}", stamp); + + let create_a = json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": rkey_a, + "record": { + "$type": "app.bsky.feed.post", + "text": "a0", + "createdAt": now + } + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.putRecord", + base_url().await + )) + .bearer_auth(&token) + .json(&create_a) + .send() + .await + .expect("Failed to seed rkey_a"); + assert_eq!(res.status(), StatusCode::OK); + + let payload = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#delete", + "collection": "app.bsky.feed.post", + "rkey": rkey_a + }, + { + "$type": "com.atproto.repo.applyWrites#create", + "collection": "app.bsky.feed.post", + "rkey": rkey_a, + "value": { + "$type": "app.bsky.feed.post", + "text": "a1", + "createdAt": now + } + }, + { + "$type": "com.atproto.repo.applyWrites#create", + "collection": "app.bsky.feed.post", + "rkey": rkey_b, + "value": { + "$type": "app.bsky.feed.post", + "text": "b", + "createdAt": now + } + } + ] + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&token) + .json(&payload) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!(res.status(), StatusCode::OK); + + let fetch = |rkey: String| { + let did = did.clone(); + let client = client.clone(); + async move { + let res = client + .get(format!( + "{}/xrpc/com.atproto.repo.getRecord", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", rkey.as_str()), + ]) + .send() + .await + .expect("Failed to fetch record"); + assert_eq!(res.status(), StatusCode::OK); + let body: Value = res.json().await.expect("Response was not valid JSON"); + body["value"]["text"].as_str().unwrap().to_string() + } + }; + + assert_eq!(fetch(rkey_a.clone()).await, "a1"); + assert_eq!(fetch(rkey_b.clone()).await, "b"); +} + +#[tokio::test] +async fn test_apply_writes_follow_create_then_delete_no_orphan_backlink() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let rkey = format!("follow_orphan_{}", Utc::now().timestamp_millis()); + let target = "did:plc:orphantargettestabcdefgh"; + + let payload = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#create", + "collection": "app.bsky.graph.follow", + "rkey": rkey, + "value": { + "$type": "app.bsky.graph.follow", + "subject": target, + "createdAt": now + } + }, + { + "$type": "com.atproto.repo.applyWrites#delete", + "collection": "app.bsky.graph.follow", + "rkey": rkey + } + ] + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&token) + .json(&payload) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!(res.status(), StatusCode::OK); + + let repo_id = repo_id_for_did(&did).await; + let lingering = follow_uris_pointing_to(repo_id, target).await; + assert!( + lingering.is_empty(), + "follow record was deleted in same batch but backlink survived: {:?}", + lingering + ); +} + +#[tokio::test] +async fn test_apply_writes_follow_update_update_final_link_wins() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let rkey = format!("follow_doubleupdate_{}", Utc::now().timestamp_millis()); + let target_initial = "did:plc:initialtarget1234567890"; + let target_intermediate = "did:plc:intermediatetarget12345"; + let target_final = "did:plc:finaltarget1234567890ab"; + + let create_payload = json!({ + "repo": did, + "collection": "app.bsky.graph.follow", + "rkey": rkey, + "record": { + "$type": "app.bsky.graph.follow", + "subject": target_initial, + "createdAt": now + } + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.putRecord", + base_url().await + )) + .bearer_auth(&token) + .json(&create_payload) + .send() + .await + .expect("Failed to create initial follow"); + assert_eq!(res.status(), StatusCode::OK); + + let payload = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#update", + "collection": "app.bsky.graph.follow", + "rkey": rkey, + "value": { + "$type": "app.bsky.graph.follow", + "subject": target_intermediate, + "createdAt": now + } + }, + { + "$type": "com.atproto.repo.applyWrites#update", + "collection": "app.bsky.graph.follow", + "rkey": rkey, + "value": { + "$type": "app.bsky.graph.follow", + "subject": target_final, + "createdAt": now + } + } + ] + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&token) + .json(&payload) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!(res.status(), StatusCode::OK); + + let repo_id = repo_id_for_did(&did).await; + let expected_uri = format!("at://{}/app.bsky.graph.follow/{}", did, rkey); + + let final_links = follow_uris_pointing_to(repo_id, target_final).await; + assert_eq!( + final_links, + vec![expected_uri.clone()], + "backlink must point to final subject" + ); + + let intermediate_links = follow_uris_pointing_to(repo_id, target_intermediate).await; + assert!( + intermediate_links.is_empty(), + "intermediate subject must not retain a backlink: {:?}", + intermediate_links + ); + + let initial_links = follow_uris_pointing_to(repo_id, target_initial).await; + assert!( + initial_links.is_empty(), + "initial subject backlink must be cleared: {:?}", + initial_links + ); +} + +#[tokio::test] +async fn test_apply_writes_create_create_same_rkey_rejected() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let rkey = format!("dup_create_{}", Utc::now().timestamp_millis()); + + let payload = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#create", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "first", + "createdAt": now + } + }, + { + "$type": "com.atproto.repo.applyWrites#create", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "second", + "createdAt": now + } + } + ] + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&token) + .json(&payload) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!( + res.status(), + StatusCode::BAD_REQUEST, + "duplicate create on same rkey within a batch must be rejected" + ); + + let get_res = client + .get(format!( + "{}/xrpc/com.atproto.repo.getRecord", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", rkey.as_str()), + ]) + .send() + .await + .expect("Failed to fetch record"); + assert_eq!( + get_res.status(), + StatusCode::NOT_FOUND, + "rejected batch must not have produced a record" + ); +} + +#[tokio::test] +async fn test_apply_writes_update_then_create_same_rkey_rejected() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let rkey = format!("update_create_{}", Utc::now().timestamp_millis()); + + let create_payload = json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": rkey, + "record": { + "$type": "app.bsky.feed.post", + "text": "v0", + "createdAt": now + } + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.putRecord", + base_url().await + )) + .bearer_auth(&token) + .json(&create_payload) + .send() + .await + .expect("Failed to seed record"); + assert_eq!(res.status(), StatusCode::OK); + + let payload = json!({ + "repo": did, + "writes": [ + { + "$type": "com.atproto.repo.applyWrites#update", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "v1", + "createdAt": now + } + }, + { + "$type": "com.atproto.repo.applyWrites#create", + "collection": "app.bsky.feed.post", + "rkey": rkey, + "value": { + "$type": "app.bsky.feed.post", + "text": "v2", + "createdAt": now + } + } + ] + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.applyWrites", + base_url().await + )) + .bearer_auth(&token) + .json(&payload) + .send() + .await + .expect("Failed to send applyWrites"); + assert_eq!( + res.status(), + StatusCode::BAD_REQUEST, + "create over a record live in this batch must be rejected" + ); + + let get_res = client + .get(format!( + "{}/xrpc/com.atproto.repo.getRecord", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", rkey.as_str()), + ]) + .send() + .await + .expect("Failed to fetch record"); + assert_eq!(get_res.status(), StatusCode::OK); + let body: Value = get_res.json().await.expect("Response was not valid JSON"); + assert_eq!( + body["value"]["text"], + json!("v0"), + "rejected batch must leave the seed record untouched" + ); +} + +#[tokio::test] +async fn test_create_record_rejects_existing_rkey() { + let client = client(); + let (token, did) = create_account_and_login(&client).await; + let now = Utc::now().to_rfc3339(); + let rkey = format!("existing_{}", Utc::now().timestamp_millis()); + + let seed = json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": rkey, + "record": { + "$type": "app.bsky.feed.post", + "text": "first", + "createdAt": now + } + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.createRecord", + base_url().await + )) + .bearer_auth(&token) + .json(&seed) + .send() + .await + .expect("Failed initial create"); + assert_eq!(res.status(), StatusCode::OK); + + let dup = json!({ + "repo": did, + "collection": "app.bsky.feed.post", + "rkey": rkey, + "record": { + "$type": "app.bsky.feed.post", + "text": "second", + "createdAt": now + } + }); + let res = client + .post(format!( + "{}/xrpc/com.atproto.repo.createRecord", + base_url().await + )) + .bearer_auth(&token) + .json(&dup) + .send() + .await + .expect("Failed duplicate create"); + assert_eq!( + res.status(), + StatusCode::BAD_REQUEST, + "createRecord on an existing rkey must be rejected" + ); + + let get_res = client + .get(format!( + "{}/xrpc/com.atproto.repo.getRecord", + base_url().await + )) + .query(&[ + ("repo", did.as_str()), + ("collection", "app.bsky.feed.post"), + ("rkey", rkey.as_str()), + ]) + .send() + .await + .expect("Failed to fetch record"); + assert_eq!(get_res.status(), StatusCode::OK); + let body: Value = get_res.json().await.expect("Response was not valid JSON"); + assert_eq!( + body["value"]["text"], + json!("first"), + "duplicate create must not have overwritten the original" + ); +}