diff --git a/crates/xrpc/src/lib.rs b/crates/xrpc/src/lib.rs index c11232e..37f210b 100644 --- a/crates/xrpc/src/lib.rs +++ b/crates/xrpc/src/lib.rs @@ -78,6 +78,7 @@ use futures::stream::{self, StreamExt, TryStreamExt}; use jacquard_common::types::did::Did; use jacquard_common::types::ident::AtIdentifier; use jacquard_common::types::nsid::Nsid; +use jacquard_common::types::recordkey::Rkey; use jacquard_common::types::string::{AtUri, Cid}; use jacquard_common::xrpc::XrpcResp; use jacquard_common::{DefaultStr, IntoStatic}; @@ -471,90 +472,114 @@ fn register_proxied( handler: H, ) -> Router where - H: Fn(AppState, HeaderMap, ProxyParams, &'static str) -> Fut + Clone + Send + Sync + 'static, + H: Fn(AppState, HeaderMap, ProxyParams, Nsid) -> Fut + + Clone + + Send + + Sync + + 'static, Fut: Future> + Send + 'static, { - nsids.iter().fold(router, |router, &nsid| { + nsids.iter().fold(router, |router, &nsid_lit| { let handler = handler.clone(); + let nsid = nsid_static(nsid_lit); router.route( - &format!("/xrpc/{nsid}"), + &format!("/xrpc/{nsid_lit}"), get( move |State(state): State, headers: HeaderMap, Query(params): Query| { - handler(state, headers, params, nsid) + handler(state, headers, params, nsid.clone()) }, ), ) }) } -#[derive(Clone, Debug, Deserialize)] -#[serde(transparent)] -pub struct RawAtUriParam(String); +#[derive(Clone, Debug)] +pub enum SubjectQuery { + Did(Did), + Uri(AtUri), +} -impl RawAtUriParam { - pub fn as_str(&self) -> &str { - &self.0 +impl<'de> Deserialize<'de> for SubjectQuery { + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'de>, + { + let raw = String::deserialize(deserializer)?; + if let Ok(did) = Did::::new_owned(&raw) { + return Ok(Self::Did(did)); + } + AtUri::::new_owned(&raw) + .map(Self::Uri) + .map_err(serde::de::Error::custom) } } -#[derive(Clone, Copy, Debug, Eq, PartialEq)] -pub struct ExpectedNsid<'a>(&'a str); +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ExpectedNsid(Nsid); -impl<'a> ExpectedNsid<'a> { - pub const fn new(nsid: &'a str) -> Self { +impl ExpectedNsid { + pub fn new(nsid: Nsid) -> Self { Self(nsid) } - pub const fn as_str(self) -> &'a str { - self.0 + pub fn from_static(s: &'static str) -> Self { + Self(nsid_static(s)) + } + + pub fn as_nsid(&self) -> &Nsid { + &self.0 + } + + pub fn as_str(&self) -> &str { + self.0.as_ref() } } #[derive(Debug, Deserialize)] struct GetRepoQuery { - repo: RawAtUriParam, + repo: AtUri, } #[derive(Debug, Deserialize)] struct GetRepoByRepoDidQuery { #[serde(rename = "repoDid")] - repo_did: String, + repo_did: Did, } #[derive(Debug, Deserialize)] struct GetProfileQuery { - actor: RawAtUriParam, + actor: AtUri, } #[derive(Debug, Deserialize)] struct GetIssueQuery { - issue: RawAtUriParam, + issue: AtUri, } #[derive(Debug, Deserialize)] struct GetPullQuery { - pull: RawAtUriParam, + pull: AtUri, } #[derive(Debug, Deserialize)] struct ListQuery { - subject: RawAtUriParam, + subject: SubjectQuery, cursor: Option, limit: Option, } #[derive(Debug, Deserialize)] struct IssueListQuery { - subject: RawAtUriParam, + subject: SubjectQuery, cursor: Option, limit: Option, - author: Option, + author: Option>, } impl IssueListQuery { - fn split(self) -> (ListQuery, Option) { + fn split(self) -> (ListQuery, Option>) { ( ListQuery { subject: self.subject, @@ -568,15 +593,15 @@ impl IssueListQuery { #[derive(Debug, Deserialize)] struct CountQuery { - subject: RawAtUriParam, + subject: SubjectQuery, } #[derive(Debug, Deserialize)] struct SearchQueryParams { q: String, - nsid: Option, - author: Option, - repo: Option, + nsid: Option>, + author: Option>, + repo: Option>, since: Option, until: Option, cursor: Option, @@ -706,12 +731,11 @@ fn format_micros(micros: u64) -> String { rfc.unwrap_or_else(|| micros.to_string()) } -fn source_did_prefix(source: &AtUri) -> Option<&str> { - source - .as_str() - .strip_prefix("at://") - .and_then(|rest| rest.split('/').next()) - .filter(|s| s.starts_with("did:")) +fn source_authority_did(source: &AtUri) -> Option> { + match source.authority() { + AtIdentifier::Did(d) => Some(d.clone().into_static()), + AtIdentifier::Handle(_) => None, + } } fn enrich_issue_list( @@ -722,14 +746,14 @@ fn enrich_issue_list( .items .into_iter() .map(|v| { - let issue_author = source_did_prefix(&v.uri).map(str::to_owned); - let repo_did = v.value.repo.as_str().to_owned(); + let issue_author = source_authority_did(&v.uri); + let repo_did = v.value.repo.clone(); enrich_view( &state.edges, - "sh.tangled.repo.issue.comment", + nsid_static("sh.tangled.repo.issue.comment"), &state.issue_states, v, - move |src| accept_state_source(src, issue_author.as_deref(), &repo_did), + move |src| accept_state_source(src, issue_author.as_ref(), &repo_did), ) }) .collect(); @@ -747,14 +771,14 @@ fn enrich_pull_list( .items .into_iter() .map(|v| { - let pull_author = source_did_prefix(&v.uri).map(str::to_owned); - let target_repo = v.value.target.repo.as_str().to_owned(); + let pull_author = source_authority_did(&v.uri); + let target_repo = v.value.target.repo.clone(); enrich_view( &state.edges, - "sh.tangled.repo.pull.comment", + nsid_static("sh.tangled.repo.pull.comment"), &state.pull_statuses, v, - move |src| accept_state_source(src, pull_author.as_deref(), &target_repo), + move |src| accept_state_source(src, pull_author.as_ref(), &target_repo), ) }) .collect(); @@ -766,18 +790,18 @@ fn enrich_pull_list( fn accept_state_source( source: &AtUri, - entity_author: Option<&str>, - repo_owner: &str, + entity_author: Option<&Did>, + repo_owner: &Did, ) -> bool { - let Some(src) = source_did_prefix(source) else { + let Some(src) = source_authority_did(source) else { return false; }; - Some(src) == entity_author || src == repo_owner + Some(&src) == entity_author || &src == repo_owner } fn enrich_view( edges: &EdgeStore, - comment_nsid: &'static str, + comment_nsid: Nsid, states: &StateIndex, view: RecordView, accept: F, @@ -787,7 +811,7 @@ where F: Fn(&AtUri) -> bool, { let comment_count = edges.count(&EdgeKey::new( - nsid_static(comment_nsid), + comment_nsid, SubjectRef::Uri(view.uri.clone()), )); let (state, state_updated_at) = states @@ -1051,25 +1075,27 @@ impl HasSubject for VouchRecord { const SHAPE: SubjectShape = SubjectShape::BareDid; } -fn parse_subject(raw: &RawAtUriParam, shape: SubjectShape) -> Result { - if let Ok(did) = Did::::new_owned(raw.as_str()) { - return match shape { - SubjectShape::BareDid | SubjectShape::BareDidOrOneOfCollections(_) => { - Ok(SubjectRef::Did(did)) - } - SubjectShape::Collection(expected) => Err(XrpcError::InvalidParams(format!( - "subject must be at:///{expected}/, got bare did" - ))), - SubjectShape::OneOfCollections(allowed) => Err(XrpcError::InvalidParams(format!( - "subject must be at://// with nsid in [{}], got bare did", - allowed.join(", "), - ))), - SubjectShape::AnyAtUri => Err(XrpcError::InvalidParams( - "subject must be at-uri form, got bare did".into(), - )), - }; - } - let uri = parse_uri(raw.as_str())?; +fn parse_subject(raw: &SubjectQuery, shape: SubjectShape) -> Result { + let uri = match raw { + SubjectQuery::Did(did) => { + return match shape { + SubjectShape::BareDid | SubjectShape::BareDidOrOneOfCollections(_) => { + Ok(SubjectRef::Did(did.clone())) + } + SubjectShape::Collection(expected) => Err(XrpcError::InvalidParams(format!( + "subject must be at:///{expected}/, got bare did" + ))), + SubjectShape::OneOfCollections(allowed) => Err(XrpcError::InvalidParams(format!( + "subject must be at://// with nsid in [{}], got bare did", + allowed.join(", "), + ))), + SubjectShape::AnyAtUri => Err(XrpcError::InvalidParams( + "subject must be at-uri form, got bare did".into(), + )), + }; + } + SubjectQuery::Uri(uri) => uri, + }; if matches!(uri.authority(), AtIdentifier::Handle(_)) { return Err(XrpcError::InvalidParams( "subject authority must be a did, not a handle".into(), @@ -1086,29 +1112,29 @@ fn parse_subject(raw: &RawAtUriParam, shape: SubjectShape) -> Result { - require_rkey(&uri, expected)?; - Ok(SubjectRef::Uri(uri)) + require_rkey(uri, expected)?; + Ok(SubjectRef::Uri(uri.clone())) } SubjectShape::Collection(expected) => Err(XrpcError::InvalidParams(format!( "subject must be at:///{expected}/, got collection {c}" ))), SubjectShape::OneOfCollections(allowed) if allowed.contains(&c) => { - require_rkey(&uri, c)?; - Ok(SubjectRef::Uri(uri)) + require_rkey(uri, c)?; + Ok(SubjectRef::Uri(uri.clone())) } SubjectShape::OneOfCollections(allowed) => Err(XrpcError::InvalidParams(format!( "subject must be at://// with nsid in [{}], got collection {c}", allowed.join(", "), ))), SubjectShape::BareDidOrOneOfCollections(allowed) if allowed.contains(&c) => { - require_rkey(&uri, c)?; - Ok(SubjectRef::Uri(uri)) + require_rkey(uri, c)?; + Ok(SubjectRef::Uri(uri.clone())) } SubjectShape::BareDidOrOneOfCollections(allowed) => Err(XrpcError::InvalidParams(format!( "subject must be a bare did or at://// with nsid in [{}], got collection {c}", allowed.join(", "), ))), - SubjectShape::AnyAtUri => Ok(SubjectRef::Uri(uri)), + SubjectShape::AnyAtUri => Ok(SubjectRef::Uri(uri.clone())), } } @@ -1130,22 +1156,20 @@ fn parse_limit(raw: Option) -> Result { .map_err(|e| XrpcError::InvalidParams(format!("limit: {e}"))) } -fn parse_author_prefix(raw: Option<&str>) -> Result, XrpcError> { - raw.map(|s| { - Did::::new_owned(s) - .map(|did| format!("at://{}/", did.as_ref())) - .map_err(|e| XrpcError::InvalidParams(format!("author: {e}"))) - }) - .transpose() +fn at_uri_owned_by(uri: &AtUri, author: &Did) -> bool { + match uri.authority() { + AtIdentifier::Did(d) => d.as_ref() == author.as_ref(), + AtIdentifier::Handle(_) => false, + } } async fn resolve_for_view( state: &AppState, - expected_nsid: &str, + expected_nsid: &Nsid, uri: AtUri, ) -> Result, XrpcError> { let raw = uri.as_ref().to_owned(); - resolve(state, ExpectedNsid::new(expected_nsid), uri) + resolve(state, ExpectedNsid::new(expected_nsid.clone()), uri) .await .map(|(body, _did)| body) .map_err(|e| match e { @@ -1156,7 +1180,7 @@ async fn resolve_for_view( async fn resolve( state: &AppState, - expected: ExpectedNsid<'_>, + expected: ExpectedNsid, uri: AtUri, ) -> Result<(Arc, Did), XrpcError> { let collection = uri @@ -1190,7 +1214,7 @@ async fn resolve( .get_record(&did_ref, &collection, &rkey) .await .map_err(map_slingshot)?; - verify_type_tag(&body, expected)?; + verify_type_tag(&body, &expected)?; state.records.put(uri, body.clone()); Ok((body, did)) } @@ -1201,7 +1225,7 @@ struct TypeTag<'a> { ty: &'a str, } -fn verify_type_tag(body: &RecordBody, expected: ExpectedNsid<'_>) -> Result<(), XrpcError> { +fn verify_type_tag(body: &RecordBody, expected: &ExpectedNsid) -> Result<(), XrpcError> { let tag: TypeTag<'_> = serde_json::from_slice(&body.value) .map_err(|e| XrpcError::InvalidRecord(format!("$type peek: {e}")))?; if tag.ty != expected.as_str() { @@ -1216,7 +1240,7 @@ fn verify_type_tag(body: &RecordBody, expected: ExpectedNsid<'_>) -> Result<(), async fn deserialize_or_upgrade( state: &AppState, - nsid: &str, + nsid: &Nsid, bytes: &[u8], ) -> Result where @@ -1241,8 +1265,9 @@ where V: serde::de::DeserializeOwned + NormalizeRepoRefs, { let raw = uri.as_str().to_owned(); - let (body, _did) = resolve(state, ExpectedNsid::new(R::NSID), uri).await?; - let value: V = deserialize_or_upgrade(state, R::NSID, &body.value).await?; + let nsid = nsid_static(R::NSID); + let (body, _did) = resolve(state, ExpectedNsid::new(nsid.clone()), uri).await?; + let value: V = deserialize_or_upgrade(state, &nsid, &body.value).await?; let value = value .normalize(&state.resolver) .await @@ -1252,13 +1277,13 @@ where async fn fetch( state: &AppState, - uri: &RawAtUriParam, + uri: &AtUri, ) -> Result<(Arc, V), XrpcError> where R: XrpcResp, V: serde::de::DeserializeOwned + NormalizeRepoRefs, { - fetch_from_uri::(state, parse_uri(uri.as_str())?).await + fetch_from_uri::(state, uri.clone()).await } async fn get_repo( @@ -1277,11 +1302,9 @@ async fn get_repo_by_repo_did( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result>, XrpcError> { - let repo_did = Did::::new_owned(&q.repo_did) - .map_err(|e| XrpcError::InvalidParams(format!("repoDid: {e}")))?; let ident = state .resolver - .lookup_by_repo_did(&repo_did) + .lookup_by_repo_did(&q.repo_did) .await .ok_or(XrpcError::NotFound)?; let uri = AtUri::::from_parts_owned( @@ -1406,28 +1429,34 @@ where .iter() .map(|s| parse_uri(s)) .collect::>()?; + let nsid = nsid_static(R::NSID); let items: Vec> = stream::iter(parsed) - .map(|uri| async move { - match resolve(state, ExpectedNsid::new(R::NSID), uri).await { - Ok((body, _)) => { - match deserialize_or_upgrade::(state, R::NSID, &body.value).await { - Ok(value) => { - let Some(value) = value.normalize(&state.resolver).await else { - return Ok(None); - }; - Ok(Some(RecordView { - uri: body.uri.clone(), - cid: Some(body.cid.clone()), - value, - })) + .map(|uri| { + let nsid = nsid.clone(); + async move { + match resolve(state, ExpectedNsid::new(nsid.clone()), uri).await { + Ok((body, _)) => { + match deserialize_or_upgrade::(state, &nsid, &body.value).await { + Ok(value) => { + let Some(value) = value.normalize(&state.resolver).await else { + return Ok(None); + }; + Ok(Some(RecordView { + uri: body.uri.clone(), + cid: Some(body.cid.clone()), + value, + })) + } + Err(_) => Ok(None), } - Err(_) => Ok(None), } + Err( + XrpcError::NotFound + | XrpcError::UpstreamGone(_) + | XrpcError::InvalidRecord(_), + ) => Ok(None), + Err(other) => Err(other), } - Err( - XrpcError::NotFound | XrpcError::UpstreamGone(_) | XrpcError::InvalidRecord(_), - ) => Ok(None), - Err(other) => Err(other), } }) .buffered(FETCH_CONCURRENCY) @@ -1448,7 +1477,7 @@ where async fn list_records_with_author( state: &AppState, q: ListQuery, - author: Option<&str>, + author: Option<&Did>, ) -> Result, XrpcError> where R: XrpcResp + HasSubject, @@ -1457,29 +1486,32 @@ where let subject = parse_subject(&q.subject, R::SHAPE)?; let cursor = parse_cursor(q.cursor.as_deref())?; let limit = parse_limit(q.limit)?; - let author_prefix = parse_author_prefix(author)?; - let key = EdgeKey::new(nsid_static(R::NSID), subject); - let EdgePage { items, next } = match author_prefix.as_deref() { - Some(prefix) => state + let nsid = nsid_static(R::NSID); + let key = EdgeKey::new(nsid.clone(), subject); + let EdgePage { items, next } = match author { + Some(did) => state .edges - .list_filtered(&key, cursor, limit, |uri| uri.as_ref().starts_with(prefix)), + .list_filtered(&key, cursor, limit, |uri| at_uri_owned_by(uri, did)), None => state.edges.list(&key, cursor, limit), }; let items = stream::iter(items) - .map(|uri| async move { - let body = resolve_for_view(state, R::NSID, uri).await?; - let value: V = match deserialize_or_upgrade::(state, R::NSID, &body.value).await { - Ok(v) => v, - Err(_) => return Ok::<_, XrpcError>(None), - }; - let Some(value) = value.normalize(&state.resolver).await else { - return Ok::<_, XrpcError>(None); - }; - Ok::<_, XrpcError>(Some(RecordView { - uri: body.uri.clone(), - cid: Some(body.cid.clone()), - value, - })) + .map(|uri| { + let nsid = nsid.clone(); + async move { + let body = resolve_for_view(state, &nsid, uri).await?; + let value: V = match deserialize_or_upgrade::(state, &nsid, &body.value).await { + Ok(v) => v, + Err(_) => return Ok::<_, XrpcError>(None), + }; + let Some(value) = value.normalize(&state.resolver).await else { + return Ok::<_, XrpcError>(None); + }; + Ok::<_, XrpcError>(Some(RecordView { + uri: body.uri.clone(), + cid: Some(body.cid.clone()), + value, + })) + } }) .buffered(FETCH_CONCURRENCY) .try_filter_map(|opt| async move { Ok(opt) }) @@ -1513,22 +1545,26 @@ where let limit = parse_limit(q.limit)?; let key = EdgeKey::new(nsid_static(M::EDGE_KIND), subject); let EdgePage { items, next } = state.edges.list(&key, cursor, limit); + let record_nsid = nsid_static(::NSID); let items = stream::iter(items) - .map(|uri| async move { - let nsid = ::NSID; - let body = resolve_for_view(state, nsid, uri).await?; - let value: V = match deserialize_or_upgrade::(state, nsid, &body.value).await { - Ok(v) => v, - Err(_) => return Ok::<_, XrpcError>(None), - }; - let Some(value) = value.normalize(&state.resolver).await else { - return Ok::<_, XrpcError>(None); - }; - Ok::<_, XrpcError>(Some(RecordView { - uri: body.uri.clone(), - cid: Some(body.cid.clone()), - value, - })) + .map(|uri| { + let record_nsid = record_nsid.clone(); + async move { + let body = resolve_for_view(state, &record_nsid, uri).await?; + let value: V = + match deserialize_or_upgrade::(state, &record_nsid, &body.value).await { + Ok(v) => v, + Err(_) => return Ok::<_, XrpcError>(None), + }; + let Some(value) = value.normalize(&state.resolver).await else { + return Ok::<_, XrpcError>(None); + }; + Ok::<_, XrpcError>(Some(RecordView { + uri: body.uri.clone(), + cid: Some(body.cid.clone()), + value, + })) + } }) .buffered(FETCH_CONCURRENCY) .try_filter_map(|opt| async move { Ok(opt) }) @@ -1582,7 +1618,7 @@ async fn list_issues( XrpcQuery(q): XrpcQuery, ) -> Result>>, XrpcError> { let (base, author) = q.split(); - let list = list_records_with_author::(&state, base, author.as_deref()).await?; + let list = list_records_with_author::(&state, base, author.as_ref()).await?; Ok(Json(enrich_issue_list(&state, list))) } @@ -1598,7 +1634,7 @@ async fn list_pulls( XrpcQuery(q): XrpcQuery, ) -> Result>>, XrpcError> { let (base, author) = q.split(); - let list = list_records_with_author::(&state, base, author.as_deref()).await?; + let list = list_records_with_author::(&state, base, author.as_ref()).await?; Ok(Json(enrich_pull_list(&state, list))) } @@ -2153,18 +2189,16 @@ async fn count_strings( #[derive(Deserialize)] struct ResolveMiniDocParams { - identifier: String, + identifier: AtIdentifier, } async fn resolve_mini_doc( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result { - let identifier = AtIdentifier::::new_owned(q.identifier.trim()) - .map_err(|e| XrpcError::InvalidParams(format!("identifier: {e}")))?; let body = state .slingshot - .resolve_mini_doc(&identifier) + .resolve_mini_doc(&q.identifier) .await .map_err(map_slingshot)?; Ok((StatusCode::OK, [(CONTENT_TYPE, "application/json")], body).into_response()) @@ -2204,30 +2238,6 @@ async fn search_query( } fn build_search_filters(q: &SearchQueryParams) -> Result { - let nsid = q - .nsid - .as_deref() - .map(|s| { - Nsid::::new_owned(s) - .map_err(|e| XrpcError::InvalidParams(format!("nsid: {e}"))) - }) - .transpose()?; - let author = q - .author - .as_deref() - .map(|s| { - Did::::new_owned(s) - .map_err(|e| XrpcError::InvalidParams(format!("author: {e}"))) - }) - .transpose()?; - let repo = q - .repo - .as_deref() - .map(|s| { - Did::::new_owned(s) - .map_err(|e| XrpcError::InvalidParams(format!("repo: {e}"))) - }) - .transpose()?; let since = q .since .as_deref() @@ -2246,9 +2256,9 @@ fn build_search_filters(q: &SearchQueryParams) -> Result Result, XrpcError> { let SearchHit { uri, nsid, score } = hit; - let body = match resolve_for_view(state, nsid.as_ref(), uri).await { + let body = match resolve_for_view(state, &nsid, uri).await { Ok(b) => b, Err(XrpcError::UpstreamGone(_)) => return Ok(None), Err(other) => return Err(other), }; - let value = match SearchableRecord::from_json_bytes(nsid.as_ref(), &body.value) { + let value = match SearchableRecord::from_json_bytes(&nsid, &body.value) { Ok(v) => v, Err(canon_err) => { - let canon_bytes = upgrade_wire_bytes(nsid.as_ref(), &body.value, &state.resolver) + let canon_bytes = upgrade_wire_bytes(&nsid, &body.value, &state.resolver) .await .map_err(|_| XrpcError::InvalidRecord(canon_err.to_string()))?; - SearchableRecord::from_json_bytes(nsid.as_ref(), &canon_bytes) + SearchableRecord::from_json_bytes(&nsid, &canon_bytes) .map_err(|e| XrpcError::InvalidRecord(e.to_string()))? } }; @@ -2348,29 +2358,28 @@ fn validate_client_supplied_knot(state: &AppState, host: &KnotHost) -> Result<() async fn resolve_knot_target( state: &AppState, - repo_uri_raw: &str, + repo_uri: AtUri, ) -> Result<(KnotHost, RepoSlug), XrpcError> { - let repo_uri = parse_uri(repo_uri_raw)?; - let rkey: Option = repo_uri.rkey().map(|r| r.as_str().to_owned()); - let (body, did) = resolve(state, ExpectedNsid::new(RepoRecord::NSID), repo_uri).await?; + let rkey: Option> = repo_uri.rkey().map(|r| r.clone().into_static()); + let (body, did) = resolve(state, ExpectedNsid::from_static(RepoRecord::NSID), repo_uri).await?; let value: Repo = serde_json::from_slice(&body.value) .map_err(|e| XrpcError::InvalidRecord(format!("decode repo record: {e}")))?; let host = KnotHost::parse(value.knot.as_ref()) .map_err(|e| XrpcError::InvalidRecord(format!("knot field: {e}")))?; - let name = pick_human_slug(rkey.as_deref(), value.name.as_deref()).ok_or_else(|| { + let name = pick_human_slug(rkey.as_ref(), value.name.as_deref()).ok_or_else(|| { XrpcError::InvalidRecord("at-uri missing rkey and record missing name".to_string()) })?; - let slug = RepoSlug::new(did.as_ref(), &name) + let slug = RepoSlug::new(&did, &name) .map_err(|e| XrpcError::InvalidRecord(format!("repo slug: {e}")))?; Ok((host, slug)) } -fn pick_human_slug(rkey: Option<&str>, name: Option<&str>) -> Option { +fn pick_human_slug(rkey: Option<&Rkey>, name: Option<&str>) -> Option { match rkey { - Some(r) if jacquard_common::types::tid::Tid::new(r).is_ok() => { - Some(name.unwrap_or(r).to_owned()) + Some(r) if jacquard_common::types::tid::Tid::new(r.as_ref()).is_ok() => { + Some(name.unwrap_or(r.as_ref()).to_owned()) } - Some(r) => Some(r.to_owned()), + Some(r) => Some(r.as_ref().to_owned()), None => name.map(str::to_owned), } } @@ -2406,7 +2415,7 @@ fn upstream_to_axum(resp: ProxyResponse) -> Response { async fn dispatch_proxy( state: AppState, headers: HeaderMap, - nsid: &'static str, + nsid: Nsid, host: KnotHost, params: ProxyParams, ) -> Result { @@ -2417,7 +2426,7 @@ async fn dispatch_proxy( let allowed = filter_request_headers(&headers); let upstream = state .knots - .forward(&host, nsid, &forward, allowed) + .forward(&host, &nsid, &forward, allowed) .await .map_err(map_proxy_error)?; Ok(upstream_to_axum(upstream)) @@ -2443,11 +2452,12 @@ async fn proxy_repo_handler( state: AppState, headers: HeaderMap, params: ProxyParams, - nsid: &'static str, + nsid: Nsid, ) -> Result { let (repo_raw, rest) = extract_param(params, REPO_PARAM)? .ok_or_else(|| XrpcError::InvalidParams("missing repo".into()))?; - let (host, slug) = resolve_knot_target(&state, &repo_raw).await?; + let repo_uri = parse_uri(&repo_raw)?; + let (host, slug) = resolve_knot_target(&state, repo_uri).await?; let forward = rest .into_iter() .chain(std::iter::once(( @@ -2462,7 +2472,7 @@ async fn proxy_knot_handler( state: AppState, headers: HeaderMap, params: ProxyParams, - nsid: &'static str, + nsid: Nsid, ) -> Result { let (knot_raw, forward) = extract_param(params, KNOT_HOST_PARAM)? .ok_or_else(|| XrpcError::InvalidParams("missing knot".into()))?; diff --git a/crates/xrpc/tests/extended.rs b/crates/xrpc/tests/extended.rs index d0cdf22..245b0f2 100644 --- a/crates/xrpc/tests/extended.rs +++ b/crates/xrpc/tests/extended.rs @@ -15,6 +15,7 @@ use http::{Request, StatusCode}; use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; +use jacquard_common::types::recordkey::Rkey; use jacquard_common::types::string::AtUri; use serde_json::{Value, json}; use tower::ServiceExt; @@ -31,6 +32,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() } @@ -80,23 +89,34 @@ impl Harness { } } - fn add_edge(&self, kind: &'static str, subject: &str, source: &str) { + fn add_edge(&self, kind: &Nsid, subject: &AtUri, source: &AtUri) { static EDGE_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1); self.edges.add(Edge { - kind: nsid(kind), - subject: subj(subject), - source: at(source), + kind: kind.clone(), + subject: subj(subject.as_ref()), + source: source.clone(), sort_micros: EDGE_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed), }); } - 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, @@ -142,17 +162,17 @@ fn label_definition_body(name: &str) -> Value { }) } -fn label_op_body(subject: &str, def_uri: &str, value: &str) -> Value { +fn label_op_body(subject: &AtUri, def_uri: &AtUri, value: &str) -> Value { json!({ "$type": "sh.tangled.label.op", "performedAt": "2026-05-01T00:00:00Z", - "subject": subject, - "add": [{"key": def_uri, "value": value}], + "subject": subject.as_ref(), + "add": [{"key": def_uri.as_ref(), "value": value}], "delete": [] }) } -fn pipeline_body(repo_did: &str) -> Value { +fn pipeline_body(repo_did: &Did) -> Value { json!({ "$type": "sh.tangled.pipeline", "workflows": [], @@ -160,7 +180,7 @@ fn pipeline_body(repo_did: &str) -> Value { "kind": "manual", "repo": { "did": "did:plc:teq", - "repoDid": repo_did, + "repoDid": repo_did.as_ref(), "knot": "nel.pet", "defaultBranch": "main" } @@ -168,14 +188,14 @@ fn pipeline_body(repo_did: &str) -> Value { }) } -fn pipeline_body_owner_only(owner_did: &str) -> Value { +fn pipeline_body_owner_only(owner_did: &Did) -> Value { json!({ "$type": "sh.tangled.pipeline", "workflows": [], "triggerMetadata": { "kind": "manual", "repo": { - "did": owner_did, + "did": owner_did.as_ref(), "knot": "nel.pet", "defaultBranch": "main" } @@ -183,22 +203,22 @@ fn pipeline_body_owner_only(owner_did: &str) -> Value { }) } -fn pipeline_status_body(pipeline_uri: &str) -> Value { +fn pipeline_status_body(pipeline_uri: &AtUri) -> Value { json!({ "$type": "sh.tangled.pipeline.status", "createdAt": "2026-05-01T00:00:00Z", - "pipeline": pipeline_uri, - "workflow": pipeline_uri, + "pipeline": pipeline_uri.as_ref(), + "workflow": pipeline_uri.as_ref(), "status": "success" }) } -fn artifact_body(repo_did: &str, name: &str) -> Value { +fn artifact_body(repo_did: &Did, name: &str) -> Value { json!({ "$type": "sh.tangled.repo.artifact", "createdAt": "2026-05-01T00:00:00Z", "name": name, - "repoDid": repo_did, + "repoDid": repo_did.as_ref(), "tag": {"$bytes": TAG_BYTES}, "artifact": { "$type": "blob", @@ -209,20 +229,20 @@ fn artifact_body(repo_did: &str, name: &str) -> Value { }) } -fn knot_member_body(subject_did: &str) -> Value { +fn knot_member_body(subject_did: &Did) -> Value { json!({ "$type": "sh.tangled.knot.member", "createdAt": "2026-05-01T00:00:00Z", - "subject": subject_did, + "subject": subject_did.as_ref(), "domain": "oyster.cafe" }) } -fn spindle_member_body(subject_did: &str) -> Value { +fn spindle_member_body(subject_did: &Did) -> Value { json!({ "$type": "sh.tangled.spindle.member", "createdAt": "2026-05-01T00:00:00Z", - "subject": subject_did, + "subject": subject_did.as_ref(), "instance": "spin.nel.pet" }) } @@ -240,18 +260,22 @@ fn string_body(filename: &str, contents: &str) -> Value { #[tokio::test] async fn list_label_definitions_keys_on_owner_did() { let h = Harness::new().await; - let owner = "did:plc:abalone"; - let rkey = "bug"; - let source = format!("at://{owner}/sh.tangled.label.definition/{rkey}"); + let owner = did("did:plc:abalone"); + let rk = rkey("bug"); + let source = at(&format!( + "at://{}/sh.tangled.label.definition/{}", + owner.as_ref(), + rk.as_ref() + )); h.add_edge( - "sh.tangled.label.definition", - &format!("at://{owner}"), + &nsid("sh.tangled.label.definition"), + &at(&format!("at://{}", owner.as_ref())), &source, ); h.mount( - owner, - "sh.tangled.label.definition", - rkey, + &owner, + &nsid("sh.tangled.label.definition"), + &rk, label_definition_body("bug"), ) .await; @@ -260,7 +284,7 @@ async fn list_label_definitions_keys_on_owner_did() { let (status, body) = json_response( app.oneshot(list_request( "sh.tangled.label.listDefinitions", - &format!("at://{owner}"), + &format!("at://{}", owner.as_ref()), &[], )) .await @@ -280,24 +304,30 @@ async fn list_label_definitions_keys_on_owner_did() { #[tokio::test] async fn count_label_definitions_dedupes_per_author() { let h = Harness::new().await; - let owner = "did:plc:abalone"; - let subject = format!("at://{owner}"); + let owner = did("did:plc:abalone"); + let subject = at(&format!("at://{}", owner.as_ref())); h.add_edge( - "sh.tangled.label.definition", + &nsid("sh.tangled.label.definition"), &subject, - &format!("at://{owner}/sh.tangled.label.definition/bug"), + &at(&format!( + "at://{}/sh.tangled.label.definition/bug", + owner.as_ref() + )), ); h.add_edge( - "sh.tangled.label.definition", + &nsid("sh.tangled.label.definition"), &subject, - &format!("at://{owner}/sh.tangled.label.definition/wontfix"), + &at(&format!( + "at://{}/sh.tangled.label.definition/wontfix", + owner.as_ref() + )), ); let app = router(h.state.clone()); let (_, body) = json_response( app.oneshot(list_request( "sh.tangled.label.countDefinitions", - &subject, + subject.as_ref(), &[], )) .await @@ -311,26 +341,30 @@ async fn count_label_definitions_dedupes_per_author() { #[tokio::test] async fn list_label_ops_accepts_issue_subject() { let h = Harness::new().await; - let issue_uri = "at://did:plc:abalone/sh.tangled.repo.issue/i1"; - let author = "did:plc:nel"; - let rkey = "op1"; - let def_uri = "at://did:plc:abalone/sh.tangled.label.definition/bug"; + let issue_uri = at("at://did:plc:abalone/sh.tangled.repo.issue/i1"); + let author = did("did:plc:nel"); + let rk = rkey("op1"); + let def_uri = at("at://did:plc:abalone/sh.tangled.label.definition/bug"); h.add_edge( - "sh.tangled.label.op", - issue_uri, - &format!("at://{author}/sh.tangled.label.op/{rkey}"), + &nsid("sh.tangled.label.op"), + &issue_uri, + &at(&format!( + "at://{}/sh.tangled.label.op/{}", + author.as_ref(), + rk.as_ref() + )), ); h.mount( - author, - "sh.tangled.label.op", - rkey, - label_op_body(issue_uri, def_uri, "true"), + &author, + &nsid("sh.tangled.label.op"), + &rk, + label_op_body(&issue_uri, &def_uri, "true"), ) .await; let app = router(h.state.clone()); let (status, body) = json_response( - app.oneshot(list_request("sh.tangled.label.listOps", issue_uri, &[])) + app.oneshot(list_request("sh.tangled.label.listOps", issue_uri.as_ref(), &[])) .await .unwrap(), ) @@ -338,32 +372,36 @@ async fn list_label_ops_accepts_issue_subject() { assert_eq!(status, StatusCode::OK); let items = body["items"].as_array().unwrap(); assert_eq!(items.len(), 1); - assert_eq!(items[0]["value"]["subject"], json!(issue_uri)); - assert_eq!(items[0]["value"]["add"][0]["key"], json!(def_uri)); + assert_eq!(items[0]["value"]["subject"], json!(issue_uri.as_ref())); + assert_eq!(items[0]["value"]["add"][0]["key"], json!(def_uri.as_ref())); } #[tokio::test] async fn list_label_ops_pull_subject_round_trip() { let h = Harness::new().await; - let pull_uri = "at://did:plc:abalone/sh.tangled.repo.pull/p1"; - let author = "did:plc:bailey"; - let rkey = "op1"; - let def_uri = "at://did:plc:abalone/sh.tangled.label.definition/wontfix"; - let source = format!("at://{author}/sh.tangled.label.op/{rkey}"); - let body = label_op_body(pull_uri, def_uri, "true"); + let pull_uri = at("at://did:plc:abalone/sh.tangled.repo.pull/p1"); + let author = did("did:plc:bailey"); + let rk = rkey("op1"); + let def_uri = at("at://did:plc:abalone/sh.tangled.label.definition/wontfix"); + let source = at(&format!( + "at://{}/sh.tangled.label.op/{}", + author.as_ref(), + rk.as_ref() + )); + let body = label_op_body(&pull_uri, &def_uri, "true"); let parsed = bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.label.op"), body.clone()) .expect("parse label.op record"); parsed - .extract_edges(&at(&source)) + .extract_edges(&source) .expect("extract") .into_iter() .for_each(|e| h.edges.add(e)); - h.mount(author, "sh.tangled.label.op", rkey, body).await; + h.mount(&author, &nsid("sh.tangled.label.op"), &rk, body).await; let app = router(h.state.clone()); let (status, json) = json_response( - app.oneshot(list_request("sh.tangled.label.listOps", pull_uri, &[])) + app.oneshot(list_request("sh.tangled.label.listOps", pull_uri.as_ref(), &[])) .await .unwrap(), ) @@ -371,8 +409,8 @@ async fn list_label_ops_pull_subject_round_trip() { assert_eq!(status, StatusCode::OK); let items = json["items"].as_array().unwrap(); assert_eq!(items.len(), 1); - assert_eq!(items[0]["value"]["subject"], json!(pull_uri)); - assert_eq!(items[0]["value"]["add"][0]["key"], json!(def_uri)); + assert_eq!(items[0]["value"]["subject"], json!(pull_uri.as_ref())); + assert_eq!(items[0]["value"]["add"][0]["key"], json!(def_uri.as_ref())); } #[tokio::test] @@ -417,20 +455,24 @@ async fn list_label_ops_rejects_unrelated_collection() { #[tokio::test] async fn list_pipelines_keys_on_repo_did() { let h = Harness::new().await; - let repo_did = "did:plc:abalone"; - let subject = format!("at://{repo_did}"); - let spindle_did = "did:plc:lyna"; - let rkey = "pl1"; + let repo_did = did("did:plc:abalone"); + let subject = at(&format!("at://{}", repo_did.as_ref())); + let spindle_did = did("did:plc:lyna"); + let rk = rkey("pl1"); h.add_edge( - "sh.tangled.pipeline", + &nsid("sh.tangled.pipeline"), &subject, - &format!("at://{spindle_did}/sh.tangled.pipeline/{rkey}"), + &at(&format!( + "at://{}/sh.tangled.pipeline/{}", + spindle_did.as_ref(), + rk.as_ref() + )), ); h.mount( - spindle_did, - "sh.tangled.pipeline", - rkey, - pipeline_body(repo_did), + &spindle_did, + &nsid("sh.tangled.pipeline"), + &rk, + pipeline_body(&repo_did), ) .await; @@ -438,7 +480,7 @@ async fn list_pipelines_keys_on_repo_did() { let (status, body) = json_response( app.oneshot(list_request( "sh.tangled.pipeline.listPipelines", - &subject, + subject.as_ref(), &[], )) .await @@ -450,7 +492,7 @@ async fn list_pipelines_keys_on_repo_did() { assert_eq!(items.len(), 1); assert_eq!( items[0]["value"]["triggerMetadata"]["repo"]["repoDid"], - json!(repo_did) + json!(repo_did.as_ref()) ); } @@ -474,19 +516,23 @@ async fn count_pipelines_returns_zero_when_no_edges() { #[tokio::test] async fn list_pipeline_statuses_keys_on_pipeline_uri() { let h = Harness::new().await; - let pipeline_uri = "at://did:plc:lyna/sh.tangled.pipeline/pl1"; - let author = "did:plc:bailey"; - let rkey = "s1"; + let pipeline_uri = at("at://did:plc:lyna/sh.tangled.pipeline/pl1"); + let author = did("did:plc:bailey"); + let rk = rkey("s1"); h.add_edge( - "sh.tangled.pipeline.status", - pipeline_uri, - &format!("at://{author}/sh.tangled.pipeline.status/{rkey}"), + &nsid("sh.tangled.pipeline.status"), + &pipeline_uri, + &at(&format!( + "at://{}/sh.tangled.pipeline.status/{}", + author.as_ref(), + rk.as_ref() + )), ); h.mount( - author, - "sh.tangled.pipeline.status", - rkey, - pipeline_status_body(pipeline_uri), + &author, + &nsid("sh.tangled.pipeline.status"), + &rk, + pipeline_status_body(&pipeline_uri), ) .await; @@ -494,7 +540,7 @@ async fn list_pipeline_statuses_keys_on_pipeline_uri() { let (status, body) = json_response( app.oneshot(list_request( "sh.tangled.pipeline.listStatuses", - pipeline_uri, + pipeline_uri.as_ref(), &[], )) .await @@ -504,7 +550,7 @@ async fn list_pipeline_statuses_keys_on_pipeline_uri() { assert_eq!(status, StatusCode::OK); let items = body["items"].as_array().unwrap(); assert_eq!(items.len(), 1); - assert_eq!(items[0]["value"]["pipeline"], json!(pipeline_uri)); + assert_eq!(items[0]["value"]["pipeline"], json!(pipeline_uri.as_ref())); assert_eq!(items[0]["value"]["status"], json!("success")); } @@ -535,88 +581,108 @@ async fn pipeline_status_endpoint_rejects_bare_did_subject() { #[tokio::test] async fn list_artifacts_keys_on_repo_did() { let h = Harness::new().await; - let repo_did = "did:plc:abalone"; - let subject = format!("at://{repo_did}"); - let owner = "did:plc:nel"; - let rkey = "a1"; + let repo_did = did("did:plc:abalone"); + let subject = at(&format!("at://{}", repo_did.as_ref())); + let owner = did("did:plc:nel"); + let rk = rkey("a1"); h.add_edge( - "sh.tangled.repo.artifact", + &nsid("sh.tangled.repo.artifact"), &subject, - &format!("at://{owner}/sh.tangled.repo.artifact/{rkey}"), + &at(&format!( + "at://{}/sh.tangled.repo.artifact/{}", + owner.as_ref(), + rk.as_ref() + )), ); h.mount( - owner, - "sh.tangled.repo.artifact", - rkey, - artifact_body(repo_did, "out.bin"), + &owner, + &nsid("sh.tangled.repo.artifact"), + &rk, + artifact_body(&repo_did, "out.bin"), ) .await; let app = router(h.state.clone()); let (status, body) = json_response( - app.oneshot(list_request("sh.tangled.repo.listArtifacts", &subject, &[])) - .await - .unwrap(), + app.oneshot(list_request( + "sh.tangled.repo.listArtifacts", + subject.as_ref(), + &[], + )) + .await + .unwrap(), ) .await; assert_eq!(status, StatusCode::OK); let items = body["items"].as_array().unwrap(); assert_eq!(items.len(), 1); assert_eq!(items[0]["value"]["name"], json!("out.bin")); - assert_eq!(items[0]["value"]["repoDid"], json!(repo_did)); + assert_eq!(items[0]["value"]["repoDid"], json!(repo_did.as_ref())); } #[tokio::test] async fn list_knot_members_keys_on_subject_did() { let h = Harness::new().await; - let subject_did = "did:plc:nel"; - let subject = format!("at://{subject_did}"); - let admin = "did:plc:teq"; - let rkey = "m1"; + let subject_did = did("did:plc:nel"); + let subject = at(&format!("at://{}", subject_did.as_ref())); + let admin = did("did:plc:teq"); + let rk = rkey("m1"); h.add_edge( - "sh.tangled.knot.member", + &nsid("sh.tangled.knot.member"), &subject, - &format!("at://{admin}/sh.tangled.knot.member/{rkey}"), + &at(&format!( + "at://{}/sh.tangled.knot.member/{}", + admin.as_ref(), + rk.as_ref() + )), ); h.mount( - admin, - "sh.tangled.knot.member", - rkey, - knot_member_body(subject_did), + &admin, + &nsid("sh.tangled.knot.member"), + &rk, + knot_member_body(&subject_did), ) .await; let app = router(h.state.clone()); let (status, body) = json_response( - app.oneshot(list_request("sh.tangled.knot.listMembers", &subject, &[])) - .await - .unwrap(), + app.oneshot(list_request( + "sh.tangled.knot.listMembers", + subject.as_ref(), + &[], + )) + .await + .unwrap(), ) .await; assert_eq!(status, StatusCode::OK); let items = body["items"].as_array().unwrap(); assert_eq!(items.len(), 1); - assert_eq!(items[0]["value"]["subject"], json!(subject_did)); + assert_eq!(items[0]["value"]["subject"], json!(subject_did.as_ref())); assert_eq!(items[0]["value"]["domain"], json!("oyster.cafe")); } #[tokio::test] async fn list_spindle_members_keys_on_subject_did() { let h = Harness::new().await; - let subject_did = "did:plc:olaren"; - let subject = format!("at://{subject_did}"); - let admin = "did:plc:teq"; - let rkey = "m1"; + let subject_did = did("did:plc:olaren"); + let subject = at(&format!("at://{}", subject_did.as_ref())); + let admin = did("did:plc:teq"); + let rk = rkey("m1"); h.add_edge( - "sh.tangled.spindle.member", + &nsid("sh.tangled.spindle.member"), &subject, - &format!("at://{admin}/sh.tangled.spindle.member/{rkey}"), + &at(&format!( + "at://{}/sh.tangled.spindle.member/{}", + admin.as_ref(), + rk.as_ref() + )), ); h.mount( - admin, - "sh.tangled.spindle.member", - rkey, - spindle_member_body(subject_did), + &admin, + &nsid("sh.tangled.spindle.member"), + &rk, + spindle_member_body(&subject_did), ) .await; @@ -624,7 +690,7 @@ async fn list_spindle_members_keys_on_subject_did() { let (status, body) = json_response( app.oneshot(list_request( "sh.tangled.spindle.listMembers", - &subject, + subject.as_ref(), &[], )) .await @@ -634,34 +700,42 @@ async fn list_spindle_members_keys_on_subject_did() { assert_eq!(status, StatusCode::OK); let items = body["items"].as_array().unwrap(); assert_eq!(items.len(), 1); - assert_eq!(items[0]["value"]["subject"], json!(subject_did)); + assert_eq!(items[0]["value"]["subject"], json!(subject_did.as_ref())); assert_eq!(items[0]["value"]["instance"], json!("spin.nel.pet")); } #[tokio::test] async fn list_strings_keys_on_owner_did() { let h = Harness::new().await; - let owner = "did:plc:abalone"; - let subject = format!("at://{owner}"); - let rkey = "k1"; + let owner = did("did:plc:abalone"); + let subject = at(&format!("at://{}", owner.as_ref())); + let rk = rkey("k1"); h.add_edge( - "sh.tangled.string", + &nsid("sh.tangled.string"), &subject, - &format!("at://{owner}/sh.tangled.string/{rkey}"), + &at(&format!( + "at://{}/sh.tangled.string/{}", + owner.as_ref(), + rk.as_ref() + )), ); h.mount( - owner, - "sh.tangled.string", - rkey, + &owner, + &nsid("sh.tangled.string"), + &rk, string_body("snippet.rs", "fn main() {}"), ) .await; let app = router(h.state.clone()); let (status, body) = json_response( - app.oneshot(list_request("sh.tangled.string.listStrings", &subject, &[])) - .await - .unwrap(), + app.oneshot(list_request( + "sh.tangled.string.listStrings", + subject.as_ref(), + &[], + )) + .await + .unwrap(), ) .await; assert_eq!(status, StatusCode::OK); @@ -674,13 +748,17 @@ async fn list_strings_keys_on_owner_did() { #[tokio::test] async fn count_strings_dedupes_per_owner() { let h = Harness::new().await; - let owner = "did:plc:abalone"; - let subject = format!("at://{owner}"); - ["k1", "k2", "k3"].iter().for_each(|rkey| { + let owner = did("did:plc:abalone"); + let subject = at(&format!("at://{}", owner.as_ref())); + ["k1", "k2", "k3"].iter().for_each(|r| { h.add_edge( - "sh.tangled.string", + &nsid("sh.tangled.string"), &subject, - &format!("at://{owner}/sh.tangled.string/{rkey}"), + &at(&format!( + "at://{}/sh.tangled.string/{}", + owner.as_ref(), + r + )), ); }); @@ -688,7 +766,7 @@ async fn count_strings_dedupes_per_owner() { let (_, body) = json_response( app.oneshot(list_request( "sh.tangled.string.countStrings", - &subject, + subject.as_ref(), &[], )) .await @@ -702,27 +780,31 @@ async fn count_strings_dedupes_per_owner() { #[tokio::test] async fn extractor_to_xrpc_round_trip_for_pipeline() { let h = Harness::new().await; - let repo_did = "did:plc:abalone"; - let spindle_did = "did:plc:lyna"; - let rkey = "pl1"; - let source = format!("at://{spindle_did}/sh.tangled.pipeline/{rkey}"); - let body = pipeline_body(repo_did); + let repo_did = did("did:plc:abalone"); + let spindle_did = did("did:plc:lyna"); + let rk = rkey("pl1"); + let source = at(&format!( + "at://{}/sh.tangled.pipeline/{}", + spindle_did.as_ref(), + rk.as_ref() + )); + let body = pipeline_body(&repo_did); let parsed = bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.pipeline"), body.clone()) .expect("parse pipeline record"); parsed - .extract_edges(&at(&source)) + .extract_edges(&source) .expect("extract") .into_iter() .for_each(|e| h.edges.add(e)); - h.mount(spindle_did, "sh.tangled.pipeline", rkey, body) + h.mount(&spindle_did, &nsid("sh.tangled.pipeline"), &rk, body) .await; let app = router(h.state.clone()); let (status, json) = json_response( app.oneshot(list_request( "sh.tangled.pipeline.listPipelines", - &format!("at://{repo_did}"), + &format!("at://{}", repo_did.as_ref()), &[], )) .await @@ -741,27 +823,31 @@ async fn extractor_to_xrpc_round_trip_for_pipeline() { #[tokio::test] async fn list_pipelines_drops_records_without_resolvable_repo_did() { let h = Harness::new().await; - let owner_did = "did:plc:nel"; - let spindle_did = "did:plc:lyna"; - let rkey = "pl1"; - let source = format!("at://{spindle_did}/sh.tangled.pipeline/{rkey}"); - let body = pipeline_body_owner_only(owner_did); + let owner_did = did("did:plc:nel"); + let spindle_did = did("did:plc:lyna"); + let rk = rkey("pl1"); + let source = at(&format!( + "at://{}/sh.tangled.pipeline/{}", + spindle_did.as_ref(), + rk.as_ref() + )); + let body = pipeline_body_owner_only(&owner_did); let parsed = bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.pipeline"), body.clone()) .expect("parse pipeline record"); parsed - .extract_edges(&at(&source)) + .extract_edges(&source) .expect("extract") .into_iter() .for_each(|e| h.edges.add(e)); - h.mount(spindle_did, "sh.tangled.pipeline", rkey, body) + h.mount(&spindle_did, &nsid("sh.tangled.pipeline"), &rk, body) .await; let app = router(h.state.clone()); let (status, json) = json_response( app.oneshot(list_request( "sh.tangled.pipeline.listPipelines", - &format!("at://{owner_did}"), + &format!("at://{}", owner_did.as_ref()), &[], )) .await