From 764b951f2d7b370655113214964efb9a72299a15 Mon Sep 17 00:00:00 2001 From: Lewis Date: Tue, 19 May 2026 22:21:33 +0300 Subject: [PATCH] refactor(xrpc,sim): newtypes in sim & tests Lewis: May this revision serve well! --- crates/bobbin-sim/src/trace_capture.rs | 8 +- .../src/workloads/cancel_mid_hydration.rs | 14 +- .../concurrent_reads_during_replay.rs | 16 +- .../bobbin-sim/src/workloads/frame_burst.rs | 14 +- .../workloads/hydrant_disconnect_barrage.rs | 14 +- .../src/workloads/slingshot_flap.rs | 14 +- crates/bobbin-sim/src/workloads/util.rs | 23 ++- crates/xrpc/tests/bulk.rs | 136 ++++++++++------ crates/xrpc/tests/cold_start.rs | 148 ++++++++++-------- crates/xrpc/tests/knot_proxy.rs | 98 +++++++----- crates/xrpc/tests/search.rs | 137 +++++++++++----- 11 files changed, 407 insertions(+), 215 deletions(-) diff --git a/crates/bobbin-sim/src/trace_capture.rs b/crates/bobbin-sim/src/trace_capture.rs index 0184268..d955d36 100644 --- a/crates/bobbin-sim/src/trace_capture.rs +++ b/crates/bobbin-sim/src/trace_capture.rs @@ -1,6 +1,8 @@ use std::sync::Arc; use std::sync::Mutex; +use jacquard_common::DefaultStr; +use jacquard_common::types::nsid::Nsid; use tracing::field::{Field, Visit}; use tracing::{Event, Subscriber}; use tracing_subscriber::Layer; @@ -51,7 +53,7 @@ where "cursor={} regime={} nsid={} edge_count={} parallelism={}", visitor.cursor.unwrap_or(u64::MAX), visitor.regime.as_deref().unwrap_or(""), - visitor.nsid.as_deref().unwrap_or(""), + visitor.nsid.as_ref().map(Nsid::as_ref).unwrap_or(""), visitor.edge_count.unwrap_or(0), visitor.parallelism.unwrap_or(0), ); @@ -63,7 +65,7 @@ where struct StageVisitor { cursor: Option, regime: Option, - nsid: Option, + nsid: Option>, edge_count: Option, parallelism: Option, } @@ -87,7 +89,7 @@ impl Visit for StageVisitor { fn record_str(&mut self, field: &Field, value: &str) { match field.name() { "regime" => self.regime = Some(value.to_owned()), - "nsid" => self.nsid = Some(value.to_owned()), + "nsid" => self.nsid = Nsid::new_owned(value).ok(), _ => {} } } diff --git a/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs b/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs index 6c885da..d112ac5 100644 --- a/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs +++ b/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs @@ -9,6 +9,8 @@ use bobbin_runtime::{ }; use bytes::Bytes; use http::StatusCode; +use jacquard_common::DefaultStr; +use jacquard_common::types::nsid::Nsid; use tokio::sync::mpsc; use url::Url; @@ -19,6 +21,10 @@ use crate::workloads::util::{format_rkey, format_tid, parse_repo_lookup}; const NAME: &str = "cancel-mid-hydration"; const REPO_COLLECTION: &str = "sh.tangled.repo"; +fn repo_collection_nsid() -> Nsid { + Nsid::new_static(REPO_COLLECTION).expect("REPO_COLLECTION literal is a valid NSID") +} + #[derive(Clone, Debug)] pub struct CancelMidHydrationConfig { pub frames: usize, @@ -247,18 +253,18 @@ struct LatentSlingshot { impl MemHttpResponder for LatentSlingshot { fn respond(&self, request: &HttpRequest) -> MemHttpResponse { self.slingshot_calls.fetch_add(1, Ordering::Relaxed); - let owner_rkey = parse_repo_lookup(&request.url, REPO_COLLECTION); + let owner_rkey = parse_repo_lookup(&request.url, &repo_collection_nsid()); let body = match owner_rkey { Some((owner, rkey)) => { - let repo_did = owner.replace("owner-", "repo-"); + let repo_did = owner.as_ref().replace("owner-", "repo-"); serde_json::json!({ - "uri": format!("at://{owner}/{REPO_COLLECTION}/{rkey}"), + "uri": format!("at://{}/{REPO_COLLECTION}/{}", owner.as_ref(), rkey.as_ref()), "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", "value": { "$type": REPO_COLLECTION, "createdAt": "2026-05-01T00:00:00Z", "knot": "oyster.cafe", - "name": "abalone", + "name": "scallop", "repoDid": repo_did, } }) diff --git a/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs b/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs index b328958..30a40ae 100644 --- a/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs +++ b/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs @@ -6,6 +6,7 @@ use std::time::Duration; use bobbin_edge_index::{PageCursor, PageLimit}; use bobbin_runtime::{Clock, MemHttpResponder, MemWsResponder, MemWsServerFuture, WsMessage}; use bobbin_types::ids::{EdgeKey, SubjectRef}; +use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; use tokio::sync::mpsc; @@ -16,7 +17,10 @@ use crate::workload::{Workload, WorkloadCtx, WorkloadHooks}; use crate::workloads::util::{AssertNoSlingshot, format_rkey, format_tid}; const NAME: &str = "concurrent-reads-during-replay"; -const TARGET_DID: &str = "did:plc:lyna"; + +fn target_did() -> Did { + Did::new_static("did:plc:lyna").expect("literal is a valid DID") +} #[derive(Clone, Debug)] pub struct ConcurrentReadsDuringReplayConfig { @@ -85,7 +89,7 @@ impl Workload for ConcurrentReadsDuringReplay { let read_key = EdgeKey::new( Nsid::new_static("sh.tangled.graph.follow").unwrap(), - SubjectRef::Did(Did::new_owned(TARGET_DID).unwrap()), + SubjectRef::Did(target_did()), ); let final_key = read_key.clone(); @@ -192,16 +196,18 @@ impl Workload for ConcurrentReadsDuringReplay { } fn build_frame_log(count: usize) -> Vec { + let target = target_did(); (0..count) .map(|i| { let id = (i + 1) as u64; - let actor = format!("did:plc:actor-{i}"); + let actor = Did::::new_owned(&format!("did:plc:nautilus-{i}")) + .expect("valid DID"); serde_json::json!({ "id": id, "type": "record", "record": { "live": false, - "did": actor, + "did": actor.as_ref(), "rev": format_tid(i), "collection": "sh.tangled.graph.follow", "rkey": format_rkey(i), @@ -209,7 +215,7 @@ fn build_frame_log(count: usize) -> Vec { "record": { "$type": "sh.tangled.graph.follow", "createdAt": "2026-05-01T00:00:00Z", - "subject": TARGET_DID, + "subject": target.as_ref(), } } }) diff --git a/crates/bobbin-sim/src/workloads/frame_burst.rs b/crates/bobbin-sim/src/workloads/frame_burst.rs index 807d398..c9654e6 100644 --- a/crates/bobbin-sim/src/workloads/frame_burst.rs +++ b/crates/bobbin-sim/src/workloads/frame_burst.rs @@ -4,6 +4,8 @@ use std::sync::atomic::{AtomicU64, Ordering}; use std::time::Duration; use bobbin_runtime::{MemHttpResponder, MemWsResponder, MemWsServerFuture, WsMessage}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; use tokio::sync::mpsc; use url::Url; @@ -12,7 +14,10 @@ use crate::workload::{Workload, WorkloadCtx, WorkloadHooks}; use crate::workloads::util::{AssertNoSlingshot, format_rkey, format_tid}; const NAME: &str = "frame-burst"; -const AUTHORITY_DID: &str = "did:plc:abalone"; + +fn authority_did() -> Did { + Did::new_static("did:plc:squid").expect("literal is a valid DID") +} pub struct FrameBurst { frames: usize, @@ -120,14 +125,17 @@ impl Workload for FrameBurst { } fn build_frame_script(count: usize) -> Vec { + let authority = authority_did(); (0..count) .map(|i| { + let repo_did = + Did::::new_owned(&format!("did:plc:squid-{i}")).expect("valid DID"); serde_json::json!({ "id": (i + 1) as u64, "type": "record", "record": { "live": false, - "did": AUTHORITY_DID, + "did": authority.as_ref(), "rev": format_tid(i), "collection": "sh.tangled.repo", "rkey": format_rkey(i), @@ -137,7 +145,7 @@ fn build_frame_script(count: usize) -> Vec { "createdAt": "2026-05-01T00:00:00Z", "knot": "oyster.cafe", "name": format!("repo-{i}"), - "repoDid": format!("did:plc:abalone-{i}") + "repoDid": repo_did.as_ref() } } }) diff --git a/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs b/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs index 3545630..926a80c 100644 --- a/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs +++ b/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs @@ -4,6 +4,8 @@ use std::sync::atomic::{AtomicU64, Ordering}; use std::time::Duration; use bobbin_runtime::{MemHttpResponder, MemWsResponder, MemWsServerFuture, WsMessage}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; use tokio::sync::mpsc; use url::Url; @@ -12,7 +14,10 @@ use crate::workload::{Workload, WorkloadCtx, WorkloadHooks}; use crate::workloads::util::{AssertNoSlingshot, format_rkey, format_tid}; const NAME: &str = "hydrant-disconnect-barrage"; -const AUTHORITY_DID: &str = "did:plc:abalone"; + +fn authority_did() -> Did { + Did::new_static("did:plc:squid").expect("literal is a valid DID") +} #[derive(Clone, Debug)] pub struct HydrantDisconnectBarrageConfig { @@ -143,15 +148,18 @@ impl Workload for HydrantDisconnectBarrage { } fn build_frame_log(count: usize) -> Vec<(u64, String)> { + let authority = authority_did(); (0..count) .map(|i| { let id = (i + 1) as u64; + let repo_did = + Did::::new_owned(&format!("did:plc:squid-{i}")).expect("valid DID"); let body = serde_json::json!({ "id": id, "type": "record", "record": { "live": false, - "did": AUTHORITY_DID, + "did": authority.as_ref(), "rev": format_tid(i), "collection": "sh.tangled.repo", "rkey": format_rkey(i), @@ -161,7 +169,7 @@ fn build_frame_log(count: usize) -> Vec<(u64, String)> { "createdAt": "2026-05-01T00:00:00Z", "knot": "oyster.cafe", "name": format!("repo-{i}"), - "repoDid": format!("did:plc:abalone-{i}") + "repoDid": repo_did.as_ref() } } }) diff --git a/crates/bobbin-sim/src/workloads/slingshot_flap.rs b/crates/bobbin-sim/src/workloads/slingshot_flap.rs index e82f0f5..1356215 100644 --- a/crates/bobbin-sim/src/workloads/slingshot_flap.rs +++ b/crates/bobbin-sim/src/workloads/slingshot_flap.rs @@ -8,6 +8,8 @@ use bobbin_runtime::{ }; use bytes::Bytes; use http::StatusCode; +use jacquard_common::DefaultStr; +use jacquard_common::types::nsid::Nsid; use jacquard_common::types::tid::Tid; use tokio::sync::mpsc; use url::Url; @@ -19,6 +21,10 @@ use crate::workloads::util::{format_rkey, format_tid, parse_repo_lookup}; const NAME: &str = "slingshot-flap"; const REPO_COLLECTION: &str = "sh.tangled.repo"; +fn repo_collection_nsid() -> Nsid { + Nsid::new_static(REPO_COLLECTION).expect("REPO_COLLECTION literal is a valid NSID") +} + #[derive(Clone, Debug)] pub struct SlingshotFlapConfig { pub cross_did_stars: usize, @@ -379,18 +385,18 @@ impl MemHttpResponder for FlapSlingshot { }; } - let owner_rkey = parse_repo_lookup(&request.url, REPO_COLLECTION); + let owner_rkey = parse_repo_lookup(&request.url, &repo_collection_nsid()); let body = match owner_rkey { Some((owner, rkey)) => { - let repo_did = owner.replace("owner-", "repo-"); + let repo_did = owner.as_ref().replace("owner-", "repo-"); serde_json::json!({ - "uri": format!("at://{owner}/{REPO_COLLECTION}/{rkey}"), + "uri": format!("at://{}/{REPO_COLLECTION}/{}", owner.as_ref(), rkey.as_ref()), "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", "value": { "$type": REPO_COLLECTION, "createdAt": "2026-05-01T00:00:00Z", "knot": "oyster.cafe", - "name": "abalone", + "name": "scallop", "repoDid": repo_did, } }) diff --git a/crates/bobbin-sim/src/workloads/util.rs b/crates/bobbin-sim/src/workloads/util.rs index 8022b17..c5c13d1 100644 --- a/crates/bobbin-sim/src/workloads/util.rs +++ b/crates/bobbin-sim/src/workloads/util.rs @@ -4,6 +4,10 @@ use std::time::Duration; use bobbin_runtime::{HttpRequest, MemHttpBody, MemHttpResponder, MemHttpResponse}; use http::StatusCode; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::nsid::Nsid; +use jacquard_common::types::recordkey::Rkey; use url::Url; const ALPHABET_RKEY: &[u8] = b"abcdefghijklmnopqrstuvwxyz234567"; @@ -29,19 +33,22 @@ fn encode_padded(mut idx: usize, alphabet: &[u8], pad: u8) -> String { String::from_utf8(buf.to_vec()).unwrap() } -pub fn parse_repo_lookup(url: &Url, expected_collection: &str) -> Option<(String, String)> { - let mut repo: Option = None; - let mut collection: Option = None; - let mut rkey: Option = None; +pub fn parse_repo_lookup( + url: &Url, + expected_collection: &Nsid, +) -> Option<(Did, Rkey)> { + let mut repo: Option> = None; + let mut collection: Option> = None; + let mut rkey: Option> = None; for (k, v) in url.query_pairs() { match k.as_ref() { - "repo" => repo = Some(v.into_owned()), - "collection" => collection = Some(v.into_owned()), - "rkey" => rkey = Some(v.into_owned()), + "repo" => repo = Did::new_owned(&v).ok(), + "collection" => collection = Nsid::new_owned(&v).ok(), + "rkey" => rkey = Rkey::new_owned(&v).ok(), _ => {} } } - if collection.as_deref() != Some(expected_collection) { + if collection.as_ref()? != expected_collection { return None; } Some((repo?, rkey?)) diff --git a/crates/xrpc/tests/bulk.rs b/crates/xrpc/tests/bulk.rs index 75b25ed..b57b6c2 100644 --- a/crates/xrpc/tests/bulk.rs +++ b/crates/xrpc/tests/bulk.rs @@ -10,6 +10,11 @@ use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader}; use bobbin_slingshot_client::SlingshotClient; use bobbin_xrpc::{AppState, router}; use http::{Request, StatusCode}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::handle::Handle; +use jacquard_common::types::nsid::Nsid; +use jacquard_common::types::recordkey::Rkey; use serde_json::{Value, json}; use tower::ServiceExt; use url::Url; @@ -19,6 +24,22 @@ use wiremock::{Mock, MockServer, ResponseTemplate}; const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; +fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() +} + +fn rkey(s: &str) -> Rkey { + Rkey::new_owned(s).unwrap() +} + +fn nsid(s: &'static str) -> Nsid { + Nsid::new_static(s).unwrap() +} + +fn handle(s: &str) -> Handle { + Handle::new_owned(s).unwrap() +} + struct Harness { server: MockServer, state: AppState, @@ -52,13 +73,24 @@ impl Harness { Self { server, state } } - async fn mount(&self, did: &str, collection: &str, rkey: &str, value: Value) { - let uri = format!("at://{did}/{collection}/{rkey}"); + async fn mount( + &self, + did: &Did, + collection: &Nsid, + rkey: &Rkey, + value: Value, + ) { + let uri = format!( + "at://{}/{}/{}", + did.as_ref(), + collection.as_ref(), + rkey.as_ref() + ); Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) - .and(query_param("repo", did)) - .and(query_param("collection", collection)) - .and(query_param("rkey", rkey)) + .and(query_param("repo", did.as_ref())) + .and(query_param("collection", collection.as_ref())) + .and(query_param("rkey", rkey.as_ref())) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "uri": uri, "cid": CID, @@ -68,12 +100,17 @@ impl Harness { .await; } - async fn mount_404(&self, did: &str, collection: &str, rkey: &str) { + async fn mount_404( + &self, + did: &Did, + collection: &Nsid, + rkey: &Rkey, + ) { Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) - .and(query_param("repo", did)) - .and(query_param("collection", collection)) - .and(query_param("rkey", rkey)) + .and(query_param("repo", did.as_ref())) + .and(query_param("collection", collection.as_ref())) + .and(query_param("rkey", rkey.as_ref())) .respond_with( ResponseTemplate::new(404) .set_body_json(json!({"error": "RecordNotFound", "message": "missing"})), @@ -106,16 +143,16 @@ async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) { (status, parsed) } -fn issue_body(repo_did: &str, title: &str) -> Value { +fn issue_body(repo_did: &Did, title: &str) -> Value { json!({ "$type": "sh.tangled.repo.issue", - "repo": repo_did, + "repo": repo_did.as_ref(), "title": title, "createdAt": "2026-05-01T00:00:00Z" }) } -fn pull_body(target_repo: &str, title: &str) -> Value { +fn pull_body(target_repo: &Did, title: &str) -> Value { json!({ "$type": "sh.tangled.repo.pull", "title": title, @@ -123,7 +160,7 @@ fn pull_body(target_repo: &str, title: &str) -> Value { "rounds": [], "target": { "branch": "main", - "repo": target_repo + "repo": target_repo.as_ref() } }) } @@ -137,11 +174,11 @@ fn repo_body(name: &str) -> Value { }) } -fn profile_body(handle: &str) -> Value { +fn profile_body(handle: &Handle) -> Value { json!({ "$type": "sh.tangled.actor.profile", "bluesky": false, - "preferredHandle": handle + "preferredHandle": handle.as_ref() }) } @@ -149,16 +186,16 @@ fn profile_body(handle: &str) -> Value { async fn get_repos_returns_all_resolved_records() { let h = Harness::new().await; h.mount( - "did:plc:nel", - "sh.tangled.repo", - "abalone", + &did("did:plc:nel"), + &nsid("sh.tangled.repo"), + &rkey("abalone"), repo_body("abalone"), ) .await; h.mount( - "did:plc:teq", - "sh.tangled.repo", - "limpet", + &did("did:plc:teq"), + &nsid("sh.tangled.repo"), + &rkey("limpet"), repo_body("limpet"), ) .await; @@ -191,17 +228,17 @@ async fn get_repos_returns_all_resolved_records() { async fn get_profiles_returns_all_resolved_profiles() { let h = Harness::new().await; h.mount( - "did:plc:nel", - "sh.tangled.actor.profile", - "self", - profile_body("witchcraft.systems"), + &did("did:plc:nel"), + &nsid("sh.tangled.actor.profile"), + &rkey("self"), + profile_body(&handle("witchcraft.systems")), ) .await; h.mount( - "did:plc:teq", - "sh.tangled.actor.profile", - "self", - profile_body("olaren.dev"), + &did("did:plc:teq"), + &nsid("sh.tangled.actor.profile"), + &rkey("self"), + profile_body(&handle("olaren.dev")), ) .await; let app = router(h.state.clone()); @@ -226,19 +263,19 @@ async fn get_profiles_returns_all_resolved_profiles() { #[tokio::test] async fn get_issues_returns_all_resolved_issues() { let h = Harness::new().await; - let repo = "did:plc:abalone"; + let repo = did("did:plc:abalone"); h.mount( - "did:plc:nel", - "sh.tangled.repo.issue", - "i1", - issue_body(repo, "first"), + &did("did:plc:nel"), + &nsid("sh.tangled.repo.issue"), + &rkey("i1"), + issue_body(&repo, "first"), ) .await; h.mount( - "did:plc:olaren", - "sh.tangled.repo.issue", - "i2", - issue_body(repo, "second"), + &did("did:plc:olaren"), + &nsid("sh.tangled.repo.issue"), + &rkey("i2"), + issue_body(&repo, "second"), ) .await; let app = router(h.state.clone()); @@ -263,12 +300,12 @@ async fn get_issues_returns_all_resolved_issues() { #[tokio::test] async fn get_pulls_returns_all_resolved_pulls() { let h = Harness::new().await; - let target = "did:plc:abalone"; + let target = did("did:plc:abalone"); h.mount( - "did:plc:nel", - "sh.tangled.repo.pull", - "p1", - pull_body(target, "patch one"), + &did("did:plc:nel"), + &nsid("sh.tangled.repo.pull"), + &rkey("p1"), + pull_body(&target, "patch one"), ) .await; let app = router(h.state.clone()); @@ -292,13 +329,18 @@ async fn get_pulls_returns_all_resolved_pulls() { async fn missing_records_are_dropped_silently() { let h = Harness::new().await; h.mount( - "did:plc:nel", - "sh.tangled.repo", - "abalone", + &did("did:plc:nel"), + &nsid("sh.tangled.repo"), + &rkey("abalone"), repo_body("abalone"), ) .await; - h.mount_404("did:plc:teq", "sh.tangled.repo", "ghost").await; + h.mount_404( + &did("did:plc:teq"), + &nsid("sh.tangled.repo"), + &rkey("ghost"), + ) + .await; let app = router(h.state.clone()); let (status, body) = json_response( app.oneshot(bulk_request( diff --git a/crates/xrpc/tests/cold_start.rs b/crates/xrpc/tests/cold_start.rs index ab6d6ba..526350b 100644 --- a/crates/xrpc/tests/cold_start.rs +++ b/crates/xrpc/tests/cold_start.rs @@ -11,6 +11,7 @@ use bobbin_slingshot_client::SlingshotClient; use bobbin_xrpc::{AppState, router}; use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; +use jacquard_common::types::nsid::Nsid; use jacquard_common::types::recordkey::Rkey; use futures::stream::{self, StreamExt}; use http::{Request, StatusCode}; @@ -23,10 +24,22 @@ use wiremock::{Mock, MockServer, ResponseTemplate}; const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; -async fn fresh_app(server_uri: &str) -> AppState { +fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() +} + +fn rkey(s: &str) -> Rkey { + Rkey::new_owned(s).unwrap() +} + +fn nsid(s: &'static str) -> Nsid { + Nsid::new_static(s).unwrap() +} + +async fn fresh_app(server_uri: &Url) -> AppState { AppState::new( Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), - SlingshotClient::with_default_http(Url::parse(server_uri).unwrap()).unwrap(), + SlingshotClient::with_default_http(server_uri.clone()).unwrap(), Arc::new(EdgeStore::new(RuntimeHasher::default())), Arc::new(StateIndex::new(RuntimeHasher::default())), Arc::new(StateIndex::new(RuntimeHasher::default())), @@ -46,21 +59,32 @@ async fn fresh_app(server_uri: &str) -> AppState { ) } -async fn mount_record(server: &MockServer, did: &str, collection: &str, rkey: &str, value: Value) { - let uri = format!("at://{did}/{collection}/{rkey}"); +async fn mount_record( + server: &MockServer, + did: &Did, + collection: &Nsid, + rkey: &Rkey, + value: Value, +) { + let uri = format!( + "at://{}/{}/{}", + did.as_ref(), + collection.as_ref(), + rkey.as_ref() + ); let body = json!({ "uri": uri, "cid": CID, "value": value }); Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) - .and(query_param("repo", did)) - .and(query_param("collection", collection)) - .and(query_param("rkey", rkey)) + .and(query_param("repo", did.as_ref())) + .and(query_param("collection", collection.as_ref())) + .and(query_param("rkey", rkey.as_ref())) .respond_with(ResponseTemplate::new(200).set_body_json(body)) .mount(server) .await; } -fn xrpc_request(endpoint: &str, param: &str, at_uri: &str) -> Request { - let encoded: String = byte_serialize(at_uri.as_bytes()).collect(); +fn xrpc_request(endpoint: &str, param: &str, value: &str) -> Request { + let encoded: String = byte_serialize(value.as_bytes()).collect(); Request::builder() .uri(format!("/xrpc/{endpoint}?{param}={encoded}")) .body(Body::empty()) @@ -77,13 +101,13 @@ async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) { #[tokio::test] async fn cold_start_serves_all_four_point_lookups() { let server = MockServer::start().await; - let did = "did:plc:clam"; + let clam = did("did:plc:clam"); mount_record( &server, - did, - "sh.tangled.repo", - "r1", + &clam, + &nsid("sh.tangled.repo"), + &rkey("r1"), json!({ "$type": "sh.tangled.repo", "name": "clam", @@ -95,9 +119,9 @@ async fn cold_start_serves_all_four_point_lookups() { mount_record( &server, - did, - "sh.tangled.actor.profile", - "self", + &clam, + &nsid("sh.tangled.actor.profile"), + &rkey("self"), json!({ "$type": "sh.tangled.actor.profile", "bluesky": false, @@ -108,9 +132,9 @@ async fn cold_start_serves_all_four_point_lookups() { mount_record( &server, - did, - "sh.tangled.repo.issue", - "i1", + &clam, + &nsid("sh.tangled.repo.issue"), + &rkey("i1"), json!({ "$type": "sh.tangled.repo.issue", "repo": "did:plc:limpet", @@ -122,9 +146,9 @@ async fn cold_start_serves_all_four_point_lookups() { mount_record( &server, - did, - "sh.tangled.repo.pull", - "p1", + &clam, + &nsid("sh.tangled.repo.pull"), + &rkey("p1"), json!({ "$type": "sh.tangled.repo.pull", "title": "ship", @@ -135,35 +159,35 @@ async fn cold_start_serves_all_four_point_lookups() { ) .await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let cases = [ ( "sh.tangled.repo.getRepo", "repo", - format!("at://{did}/sh.tangled.repo/r1"), + format!("at://{}/sh.tangled.repo/r1", clam.as_ref()), "knot", json!("oyster.cafe"), ), ( "sh.tangled.actor.getProfile", "actor", - format!("at://{did}/sh.tangled.actor.profile/self"), + format!("at://{}/sh.tangled.actor.profile/self", clam.as_ref()), "description", json!("clam shell"), ), ( "sh.tangled.repo.getIssue", "issue", - format!("at://{did}/sh.tangled.repo.issue/i1"), + format!("at://{}/sh.tangled.repo.issue/i1", clam.as_ref()), "title", json!("broken"), ), ( "sh.tangled.repo.getPull", "pull", - format!("at://{did}/sh.tangled.repo.pull/p1"), + format!("at://{}/sh.tangled.repo.pull/p1", clam.as_ref()), "title", json!("ship"), ), @@ -193,12 +217,12 @@ async fn cold_start_serves_all_four_point_lookups() { #[tokio::test] async fn second_call_is_served_from_lru() { let server = MockServer::start().await; - let did = "did:plc:uni"; + let uni = did("did:plc:uni"); let mock = Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) - .and(query_param("repo", did)) + .and(query_param("repo", uni.as_ref())) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ - "uri": format!("at://{did}/sh.tangled.repo/r1"), + "uri": format!("at://{}/sh.tangled.repo/r1", uni.as_ref()), "cid": CID, "value": { "$type": "sh.tangled.repo", @@ -211,9 +235,9 @@ async fn second_call_is_served_from_lru() { .mount_as_scoped(&server) .await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); - let req_uri = format!("at://{did}/sh.tangled.repo/r1"); + let req_uri = format!("at://{}/sh.tangled.repo/r1", uni.as_ref()); stream::iter(0..3) .for_each(|_| { @@ -235,7 +259,7 @@ async fn second_call_is_served_from_lru() { #[tokio::test] async fn collection_mismatch_is_400() { let server = MockServer::start().await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let resp = app .oneshot(xrpc_request( @@ -251,7 +275,7 @@ async fn collection_mismatch_is_400() { #[tokio::test] async fn handle_authority_is_400() { let server = MockServer::start().await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let resp = app .oneshot(xrpc_request( @@ -272,7 +296,7 @@ async fn slingshot_404_propagates_as_404() { .respond_with(ResponseTemplate::new(404)) .mount(&server) .await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let resp = app .oneshot(xrpc_request( @@ -290,9 +314,9 @@ async fn wrong_record_type_is_502() { let server = MockServer::start().await; mount_record( &server, - "did:plc:clam", - "sh.tangled.repo", - "r1", + &did("did:plc:clam"), + &nsid("sh.tangled.repo"), + &rkey("r1"), json!({ "$type": "sh.tangled.knot", "knot": "oyster.cafe", @@ -300,7 +324,7 @@ async fn wrong_record_type_is_502() { }), ) .await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let resp = app .oneshot(xrpc_request( @@ -333,7 +357,7 @@ async fn wrong_type_does_not_poison_cache() { .expect(2) .mount_as_scoped(&server) .await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let req = || { xrpc_request( @@ -352,7 +376,7 @@ async fn wrong_type_does_not_poison_cache() { #[tokio::test] async fn missing_uri_param_returns_json_envelope() { let server = MockServer::start().await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let req = Request::builder() .uri("/xrpc/sh.tangled.repo.getRepo") @@ -368,7 +392,7 @@ async fn missing_uri_param_returns_json_envelope() { #[tokio::test] async fn malformed_at_uri_returns_400_envelope() { let server = MockServer::start().await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let resp = app .oneshot(xrpc_request( @@ -400,7 +424,7 @@ async fn upstream_uri_mismatch_routes_to_invalid_record() { }))) .mount(&server) .await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let resp = app .oneshot(xrpc_request( @@ -431,7 +455,7 @@ async fn upstream_garbage_cid_routes_to_invalid_record() { }))) .mount(&server) .await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let resp = app .oneshot(xrpc_request( @@ -459,7 +483,7 @@ async fn oversize_upstream_body_routes_to_upstream_failed() { ) .mount(&server) .await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let resp = app .oneshot(xrpc_request( @@ -482,7 +506,7 @@ async fn upstream_503_routes_to_upstream_failed() { .respond_with(ResponseTemplate::new(503)) .mount(&server) .await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let resp = app .oneshot(xrpc_request( @@ -500,32 +524,28 @@ async fn upstream_503_routes_to_upstream_failed() { #[tokio::test] async fn get_repo_by_repo_did_returns_observed_record() { let server = MockServer::start().await; - let owner_did = "did:plc:scallop"; - let rkey = "r1"; - let repo_did = "did:plc:limpet"; + let owner_did = did("did:plc:scallop"); + let rk = rkey("r1"); + let repo_did = did("did:plc:limpet"); mount_record( &server, - owner_did, - "sh.tangled.repo", - rkey, + &owner_did, + &nsid("sh.tangled.repo"), + &rk, json!({ "$type": "sh.tangled.repo", "name": "scallop", "knot": "oyster.cafe", "createdAt": "2026-05-01T00:00:00Z", - "repoDid": repo_did, + "repoDid": repo_did.as_ref(), }), ) .await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; state .resolver - .observe( - Did::::new_owned(owner_did).unwrap(), - Rkey::::new_owned(rkey).unwrap(), - Some(Did::::new_owned(repo_did).unwrap()), - ) + .observe(owner_did.clone(), rk.clone(), Some(repo_did.clone())) .await; let app = router(state); @@ -533,7 +553,7 @@ async fn get_repo_by_repo_did_returns_observed_record() { .oneshot(xrpc_request( "sh.tangled.repo.getRepoByRepoDid", "repoDid", - repo_did, + repo_did.as_ref(), )) .await .unwrap(); @@ -541,16 +561,16 @@ async fn get_repo_by_repo_did_returns_observed_record() { assert_eq!(status, StatusCode::OK); assert_eq!( body["uri"], - format!("at://{owner_did}/sh.tangled.repo/{rkey}") + format!("at://{}/sh.tangled.repo/{}", owner_did.as_ref(), rk.as_ref()) ); assert_eq!(body["value"]["name"], "scallop"); - assert_eq!(body["value"]["repoDid"], repo_did); + assert_eq!(body["value"]["repoDid"], repo_did.as_ref()); } #[tokio::test] async fn get_repo_by_repo_did_404_when_unobserved() { let server = MockServer::start().await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let resp = app .oneshot(xrpc_request( @@ -566,7 +586,7 @@ async fn get_repo_by_repo_did_404_when_unobserved() { #[tokio::test] async fn get_repo_by_repo_did_400_on_invalid_did() { let server = MockServer::start().await; - let state = fresh_app(&server.uri()).await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; let app = router(state); let resp = app .oneshot(xrpc_request( diff --git a/crates/xrpc/tests/knot_proxy.rs b/crates/xrpc/tests/knot_proxy.rs index a093e94..b7d3821 100644 --- a/crates/xrpc/tests/knot_proxy.rs +++ b/crates/xrpc/tests/knot_proxy.rs @@ -11,6 +11,9 @@ use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader}; use bobbin_slingshot_client::SlingshotClient; use bobbin_xrpc::{AppState, router}; use http::{Request, StatusCode}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::recordkey::Rkey; use serde_json::{Value, json}; use tower::ServiceExt; use url::Url; @@ -20,6 +23,14 @@ use wiremock::{Mock, MockServer, ResponseTemplate}; const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; +fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() +} + +fn rkey(s: &str) -> Rkey { + Rkey::new_owned(s).unwrap() +} + fn test_config() -> KnotProxyConfig { KnotProxyConfig { failure_threshold: FailureThreshold::new(2).unwrap(), @@ -79,15 +90,24 @@ impl Harness { } } - async fn mount_repo_record(&self, did: &str, rkey: &str, name: &str) { + async fn mount_repo_record(&self, did: &Did, rkey: &Rkey, name: &str) { self.mount_repo_record_inner(did, rkey, Some(name)).await; } - async fn mount_repo_record_rkey_as_name(&self, did: &str, rkey: &str) { + async fn mount_repo_record_rkey_as_name( + &self, + did: &Did, + rkey: &Rkey, + ) { self.mount_repo_record_inner(did, rkey, None).await; } - async fn mount_repo_record_inner(&self, did: &str, rkey: &str, name: Option<&str>) { + async fn mount_repo_record_inner( + &self, + did: &Did, + rkey: &Rkey, + name: Option<&str>, + ) { let knot_value = self.knot.uri(); let mut record = json!({ "$type": "sh.tangled.repo", @@ -97,12 +117,16 @@ impl Harness { if let Some(n) = name { record["name"] = json!(n); } - let uri = format!("at://{did}/sh.tangled.repo/{rkey}"); + let uri = format!( + "at://{}/sh.tangled.repo/{}", + did.as_ref(), + rkey.as_ref() + ); Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) - .and(query_param("repo", did)) + .and(query_param("repo", did.as_ref())) .and(query_param("collection", "sh.tangled.repo")) - .and(query_param("rkey", rkey)) + .and(query_param("rkey", rkey.as_ref())) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "uri": uri, "cid": CID, @@ -151,7 +175,7 @@ async fn body_value(resp: http::Response) -> Value { async fn proxies_repo_blob_with_did_slash_name_repo_param() { let h = Harness::new().await; let tid = "3jzfcijpj2z2a"; - h.mount_repo_record("did:plc:abalone", tid, "barnacle") + h.mount_repo_record(&did("did:plc:abalone"), &rkey(tid), "barnacle") .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.blob")) @@ -183,7 +207,7 @@ async fn proxies_repo_blob_with_did_slash_name_repo_param() { #[tokio::test] async fn modern_rkey_as_name_uses_rkey_even_when_name_field_set() { let h = Harness::new().await; - h.mount_repo_record("did:plc:abalone", "core", "Tangled Core") + h.mount_repo_record(&did("did:plc:abalone"), &rkey("core"), "Tangled Core") .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.getDefaultBranch")) @@ -209,7 +233,7 @@ async fn modern_rkey_as_name_uses_rkey_even_when_name_field_set() { #[tokio::test] async fn modern_rkey_as_name_works_when_name_field_null() { let h = Harness::new().await; - h.mount_repo_record_rkey_as_name("did:plc:abalone", "core") + h.mount_repo_record_rkey_as_name(&did("did:plc:abalone"), &rkey("core")) .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.getDefaultBranch")) @@ -234,7 +258,7 @@ async fn modern_rkey_as_name_works_when_name_field_null() { async fn legacy_tid_rkey_falls_back_to_name_field() { let h = Harness::new().await; let tid_rkey = "3jzfcijpj2z2a"; - h.mount_repo_record("did:plc:abalone", tid_rkey, "dotfiles") + h.mount_repo_record(&did("did:plc:abalone"), &rkey(tid_rkey), "dotfiles") .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.getDefaultBranch")) @@ -259,7 +283,7 @@ async fn legacy_tid_rkey_falls_back_to_name_field() { async fn tid_rkey_without_name_falls_back_to_tid() { let h = Harness::new().await; let tid_rkey = "3jzfcijpj2z2a"; - h.mount_repo_record_rkey_as_name("did:plc:abalone", tid_rkey) + h.mount_repo_record_rkey_as_name(&did("did:plc:abalone"), &rkey(tid_rkey)) .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.getDefaultBranch")) @@ -283,7 +307,7 @@ async fn tid_rkey_without_name_falls_back_to_tid() { async fn streams_binary_archive_through_proxy() { let h = Harness::new().await; let tid = "3jzfcijpj2z2b"; - h.mount_repo_record("did:plc:limpet", tid, "kelp").await; + h.mount_repo_record(&did("did:plc:limpet"), &rkey(tid), "kelp").await; let payload: Vec = (0u8..=255).collect(); Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.archive")) @@ -341,7 +365,7 @@ async fn unknown_repo_propagates_404_from_slingshot() { #[tokio::test] async fn knot_5xx_routes_to_upstream_failed() { let h = Harness::new().await; - h.mount_repo_record("did:plc:abalone", "r1", "barnacle") + h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.blob")) @@ -361,7 +385,7 @@ async fn knot_5xx_routes_to_upstream_failed() { #[tokio::test] async fn knot_4xx_passes_through_unchanged() { let h = Harness::new().await; - h.mount_repo_record("did:plc:abalone", "r1", "barnacle") + h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.blob")) @@ -384,7 +408,7 @@ async fn knot_4xx_passes_through_unchanged() { #[tokio::test] async fn breaker_opens_after_threshold_then_short_circuits() { let h = Harness::new().await; - h.mount_repo_record("did:plc:abalone", "r1", "barnacle") + h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.blob")) @@ -479,7 +503,7 @@ async fn missing_knot_param_on_knot_route_returns_400() { async fn second_proxy_call_skips_slingshot_via_lru() { let h = Harness::new().await; let tid = "3jzfcijpj2z2c"; - h.mount_repo_record("did:plc:abalone", tid, "barnacle") + h.mount_repo_record(&did("did:plc:abalone"), &rkey(tid), "barnacle") .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.tree")) @@ -515,7 +539,7 @@ async fn second_proxy_call_skips_slingshot_via_lru() { #[tokio::test] async fn does_not_inject_auth_or_atproto_proxy_headers() { let h = Harness::new().await; - h.mount_repo_record("did:plc:abalone", "r1", "barnacle") + h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.blob")) @@ -554,7 +578,7 @@ async fn does_not_inject_auth_or_atproto_proxy_headers() { async fn forwards_range_and_conditional_request_headers() { let h = Harness::new().await; let tid = "3jzfcijpj2z2d"; - h.mount_repo_record("did:plc:limpet", tid, "kelp").await; + h.mount_repo_record(&did("did:plc:limpet"), &rkey(tid), "kelp").await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.archive")) .and(query_param("repo", "did:plc:limpet/kelp")) @@ -607,7 +631,7 @@ async fn forwards_range_and_conditional_request_headers() { #[tokio::test] async fn drops_disallowed_client_headers() { let h = Harness::new().await; - h.mount_repo_record("did:plc:abalone", "r1", "barnacle") + h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.blob")) @@ -689,20 +713,20 @@ async fn record_with_private_knot_returns_invalid_record() { ..test_config() }; let h = Harness::with_config(strict).await; - let did = "did:plc:abalone"; - let rkey = "r1"; + let owner = did("did:plc:abalone"); + let rk = rkey("r1"); let record = json!({ "$type": "sh.tangled.repo", "createdAt": "2026-05-01T00:00:00Z", "knot": "http://10.0.0.5:3000", "name": "barnacle", }); - let uri = format!("at://{did}/sh.tangled.repo/{rkey}"); + let uri = format!("at://{}/sh.tangled.repo/{}", owner.as_ref(), rk.as_ref()); Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) - .and(query_param("repo", did)) + .and(query_param("repo", owner.as_ref())) .and(query_param("collection", "sh.tangled.repo")) - .and(query_param("rkey", rkey)) + .and(query_param("rkey", rk.as_ref())) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "uri": uri, "cid": CID, @@ -731,20 +755,20 @@ async fn strips_basic_auth_from_credentialed_knot_url() { parsed.host_str().unwrap(), parsed.port().unwrap(), ); - let did = "did:plc:abalone"; - let rkey = "r1"; + let owner = did("did:plc:abalone"); + let rk = rkey("r1"); let record = json!({ "$type": "sh.tangled.repo", "createdAt": "2026-05-01T00:00:00Z", "knot": knot_with_creds, "name": "barnacle", }); - let uri = format!("at://{did}/sh.tangled.repo/{rkey}"); + let uri = format!("at://{}/sh.tangled.repo/{}", owner.as_ref(), rk.as_ref()); Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) - .and(query_param("repo", did)) + .and(query_param("repo", owner.as_ref())) .and(query_param("collection", "sh.tangled.repo")) - .and(query_param("rkey", rkey)) + .and(query_param("rkey", rk.as_ref())) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "uri": uri, "cid": CID, @@ -777,7 +801,7 @@ async fn strips_basic_auth_from_credentialed_knot_url() { async fn knot_redirect_surfaces_as_upstream_failed() { let h = Harness::new().await; let secondary = MockServer::start().await; - h.mount_repo_record("did:plc:abalone", "r1", "barnacle") + h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") .await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.blob")) @@ -808,7 +832,7 @@ async fn knot_redirect_surfaces_as_upstream_failed() { #[tokio::test] async fn forwards_repeated_query_params() { let h = Harness::new().await; - h.mount_repo_record("did:plc:limpet", "r4", "kelp").await; + h.mount_repo_record(&did("did:plc:limpet"), &rkey("r4"), "kelp").await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.tags")) .respond_with(ResponseTemplate::new(200).set_body_raw(r#"{"tags":[]}"#, "application/json")) @@ -902,20 +926,20 @@ async fn record_with_plaintext_knot_returns_invalid_record_when_https_required() ..test_config() }; let h = Harness::with_config(strict).await; - let did = "did:plc:abalone"; - let rkey = "r1"; + let owner = did("did:plc:abalone"); + let rk = rkey("r1"); let record = json!({ "$type": "sh.tangled.repo", "createdAt": "2026-05-01T00:00:00Z", "knot": "http://oyster.cafe", "name": "barnacle", }); - let uri = format!("at://{did}/sh.tangled.repo/{rkey}"); + let uri = format!("at://{}/sh.tangled.repo/{}", owner.as_ref(), rk.as_ref()); Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) - .and(query_param("repo", did)) + .and(query_param("repo", owner.as_ref())) .and(query_param("collection", "sh.tangled.repo")) - .and(query_param("rkey", rkey)) + .and(query_param("rkey", rk.as_ref())) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "uri": uri, "cid": CID, @@ -944,7 +968,7 @@ async fn record_with_plaintext_knot_returns_invalid_record_when_https_required() #[tokio::test] async fn knot_not_modified_passes_through() { let h = Harness::new().await; - h.mount_repo_record("did:plc:limpet", "r5", "kelp").await; + h.mount_repo_record(&did("did:plc:limpet"), &rkey("r5"), "kelp").await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.archive")) .respond_with(ResponseTemplate::new(304).insert_header("etag", "\"v1\"")) diff --git a/crates/xrpc/tests/search.rs b/crates/xrpc/tests/search.rs index b6916f6..86e3e8d 100644 --- a/crates/xrpc/tests/search.rs +++ b/crates/xrpc/tests/search.rs @@ -13,6 +13,7 @@ use bobbin_xrpc::{AppState, router}; use http::{Request, StatusCode}; use jacquard_common::DefaultStr; use jacquard_common::types::nsid::Nsid; +use jacquard_common::types::recordkey::Rkey; use jacquard_common::types::string::{AtUri, Did}; use serde_json::{Value, json}; use tower::ServiceExt; @@ -27,6 +28,14 @@ fn at(s: &str) -> AtUri { AtUri::new_owned(s).unwrap() } +fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() +} + +fn rkey(s: &str) -> Rkey { + Rkey::new_owned(s).unwrap() +} + fn nsid(s: &'static str) -> Nsid { Nsid::new_static(s).unwrap() } @@ -76,38 +85,58 @@ impl Harness { } } - async fn index_issue(&self, did: &str, rkey: &str, title: &str, body: &str) { + async fn index_issue( + &self, + did: &Did, + rkey: &Rkey, + title: &str, + body: &str, + ) { self.index_issue_at(did, rkey, title, body, None, None) .await; } async fn index_issue_at( &self, - did: &str, - rkey: &str, + did: &Did, + rkey: &Rkey, title: &str, body: &str, created_at: Option, - repo: Option<&str>, + repo: Option<&Did>, ) { - let uri = format!("at://{did}/sh.tangled.repo.issue/{rkey}"); + let uri = format!( + "at://{}/sh.tangled.repo.issue/{}", + did.as_ref(), + rkey.as_ref() + ); self.search .upsert(SearchDoc { uri: at(&uri), nsid: nsid("sh.tangled.repo.issue"), title: title.to_owned(), body: body.to_owned(), - author: Some(Did::::new_owned(did).unwrap()), + author: Some(did.clone()), created_at, - repo: repo.map(|d| Did::::new_owned(d).unwrap()), + repo: repo.cloned(), }) .await; self.search.flush().await; self.mount_issue(did, rkey, title, body).await; } - async fn index_string(&self, did: &str, rkey: &str, filename: &str, contents: &str) { - let uri = format!("at://{did}/sh.tangled.string/{rkey}"); + async fn index_string( + &self, + did: &Did, + rkey: &Rkey, + filename: &str, + contents: &str, + ) { + let uri = format!( + "at://{}/sh.tangled.string/{}", + did.as_ref(), + rkey.as_ref() + ); self.search .upsert(SearchDoc { uri: at(&uri), @@ -123,8 +152,18 @@ impl Harness { self.mount_string(did, rkey, filename, contents).await; } - async fn mount_issue(&self, did: &str, rkey: &str, title: &str, body: &str) { - let uri = format!("at://{did}/sh.tangled.repo.issue/{rkey}"); + async fn mount_issue( + &self, + did: &Did, + rkey: &Rkey, + title: &str, + body: &str, + ) { + let uri = format!( + "at://{}/sh.tangled.repo.issue/{}", + did.as_ref(), + rkey.as_ref() + ); let value = json!({ "$type": "sh.tangled.repo.issue", "repo": "did:plc:abalone", @@ -134,9 +173,9 @@ impl Harness { }); Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) - .and(query_param("repo", did)) + .and(query_param("repo", did.as_ref())) .and(query_param("collection", "sh.tangled.repo.issue")) - .and(query_param("rkey", rkey)) + .and(query_param("rkey", rkey.as_ref())) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "uri": uri, "cid": CID, @@ -146,8 +185,18 @@ impl Harness { .await; } - async fn mount_string(&self, did: &str, rkey: &str, filename: &str, contents: &str) { - let uri = format!("at://{did}/sh.tangled.string/{rkey}"); + async fn mount_string( + &self, + did: &Did, + rkey: &Rkey, + filename: &str, + contents: &str, + ) { + let uri = format!( + "at://{}/sh.tangled.string/{}", + did.as_ref(), + rkey.as_ref() + ); let value = json!({ "$type": "sh.tangled.string", "filename": filename, @@ -157,9 +206,9 @@ impl Harness { }); Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) - .and(query_param("repo", did)) + .and(query_param("repo", did.as_ref())) .and(query_param("collection", "sh.tangled.string")) - .and(query_param("rkey", rkey)) + .and(query_param("rkey", rkey.as_ref())) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "uri": uri, "cid": CID, @@ -221,8 +270,8 @@ async fn empty_query_returns_no_hits() { async fn indexed_issue_hydrates_typed_value() { let h = Harness::new().await; h.index_issue( - "did:plc:nel", - "abcabcabcabcz", + &did("did:plc:nel"), + &rkey("abcabcabcabcz"), "barnacle pagination overflow", "scrolling resets when the cursor wraps", ) @@ -258,9 +307,9 @@ async fn indexed_issue_hydrates_typed_value() { #[tokio::test] async fn nsid_filter_narrows_to_single_collection() { let h = Harness::new().await; - h.index_issue("did:plc:nel", "i1", "anemone tide", "high water") + h.index_issue(&did("did:plc:nel"), &rkey("i1"), "anemone tide", "high water") .await; - h.index_string("did:plc:teq", "k1", "anemone.md", "anemone notes") + h.index_string(&did("did:plc:teq"), &rkey("k1"), "anemone.md", "anemone notes") .await; let app = router(h.state.clone()); let resp = app @@ -283,7 +332,7 @@ async fn nsid_filter_narrows_to_single_collection() { #[tokio::test] async fn hit_set_stable_across_coverage_promotion() { let h = Harness::new().await; - h.index_issue("did:plc:nel", "i1", "limpet survey", "tidal pools") + h.index_issue(&did("did:plc:nel"), &rkey("i1"), "limpet survey", "tidal pools") .await; let app = router(h.state.clone()); h.warming(3, 5); @@ -315,8 +364,8 @@ async fn pagination_round_trips_via_cursor() { let names = ["nel", "olaren", "teq", "lyna", "bailey"]; for (i, owner) in names.iter().enumerate() { h.index_issue( - &format!("did:plc:{owner}"), - &format!("r{i}"), + &did(&format!("did:plc:{owner}")), + &rkey(&format!("r{i}")), "anemone tides", "shell", ) @@ -415,7 +464,7 @@ async fn tombstoned_hit_silently_dropped_from_results() { }) .await; h.search.flush().await; - h.index_issue("did:plc:teq", "i2", "kelp survives", "still here") + h.index_issue(&did("did:plc:teq"), &rkey("i2"), "kelp survives", "still here") .await; let app = router(h.state.clone()); let resp = app.oneshot(search_request(&[("q", "kelp")])).await.unwrap(); @@ -508,8 +557,8 @@ async fn second_query_short_circuits_via_lru_without_re_querying_slingshot() { #[tokio::test] async fn author_filter_narrows_to_matching_did() { let h = Harness::new().await; - h.index_issue("did:plc:nel", "i1", "kelp tide", "").await; - h.index_issue("did:plc:teq", "i2", "kelp wave", "").await; + h.index_issue(&did("did:plc:nel"), &rkey("i1"), "kelp tide", "").await; + h.index_issue(&did("did:plc:teq"), &rkey("i2"), "kelp wave", "").await; let app = router(h.state.clone()); let resp = app .oneshot(search_request(&[("q", "kelp"), ("author", "did:plc:nel")])) @@ -533,11 +582,11 @@ async fn since_until_window_filters_by_created_at() { let early = 1_700_000_000; let mid = 1_750_000_000; let late = 1_800_000_000; - h.index_issue_at("did:plc:nel", "i1", "kelp early", "", Some(early), None) + h.index_issue_at(&did("did:plc:nel"), &rkey("i1"), "kelp early", "", Some(early), None) .await; - h.index_issue_at("did:plc:nel", "i2", "kelp mid", "", Some(mid), None) + h.index_issue_at(&did("did:plc:nel"), &rkey("i2"), "kelp mid", "", Some(mid), None) .await; - h.index_issue_at("did:plc:nel", "i3", "kelp late", "", Some(late), None) + h.index_issue_at(&did("did:plc:nel"), &rkey("i3"), "kelp late", "", Some(late), None) .await; let app = router(h.state.clone()); let resp = app @@ -561,15 +610,29 @@ async fn since_until_window_filters_by_created_at() { #[tokio::test] async fn repo_filter_scopes_to_owning_repo() { let h = Harness::new().await; - let abalone = "did:plc:abalone"; - let limpet = "did:plc:limpet"; - h.index_issue_at("did:plc:nel", "i1", "kelp one", "", None, Some(abalone)) - .await; - h.index_issue_at("did:plc:teq", "i2", "kelp two", "", None, Some(limpet)) - .await; + let abalone = did("did:plc:abalone"); + let limpet = did("did:plc:limpet"); + h.index_issue_at( + &did("did:plc:nel"), + &rkey("i1"), + "kelp one", + "", + None, + Some(&abalone), + ) + .await; + h.index_issue_at( + &did("did:plc:teq"), + &rkey("i2"), + "kelp two", + "", + None, + Some(&limpet), + ) + .await; let app = router(h.state.clone()); let resp = app - .oneshot(search_request(&[("q", "kelp"), ("repo", abalone)])) + .oneshot(search_request(&[("q", "kelp"), ("repo", abalone.as_ref())])) .await .unwrap(); let (status, body) = json_response(resp).await; -- 2.51.2