diff --git a/Cargo.lock b/Cargo.lock index 862bd51..92de3e5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -370,6 +370,7 @@ dependencies = [ "bytes", "futures", "http", + "jacquard-common", "reqwest", "scc", "thiserror 2.0.18", diff --git a/crates/edge-index/src/lib.rs b/crates/edge-index/src/lib.rs index 7fd5a6c..fe35bcd 100644 --- a/crates/edge-index/src/lib.rs +++ b/crates/edge-index/src/lib.rs @@ -266,7 +266,7 @@ impl EdgeStore { fn add_locked(&self, edge: Edge) { let source_spur = self.source_interner.get_or_intern(edge.source.as_ref()); let id = SourceId::from_spur(source_spur); - let author = self.intern_author(edge.source.as_ref()); + let author = self.intern_author(&edge.source); let sort_micros = edge.sort_micros; let key = EdgeKey::new(edge.kind, edge.subject); @@ -292,7 +292,7 @@ impl EdgeStore { return; }; let id = SourceId::from_spur(source_spur); - let author = author_str(source.as_ref()) + let author = source_authority_did(source) .and_then(|s| self.did_interner.get(s)) .map(AuthorId::from_spur); let Some((_, entries)) = self.reverse.remove_sync(&id) else { @@ -312,8 +312,8 @@ impl EdgeStore { }); } - fn intern_author(&self, source: &str) -> Option { - let did = author_str(source)?; + fn intern_author(&self, source: &AtUri) -> Option { + let did = source_authority_did(source)?; Some(AuthorId::from_spur(self.did_interner.get_or_intern(did))) } @@ -437,9 +437,11 @@ impl EdgeStore { } } -fn author_str(source: &str) -> Option<&str> { - let rest = source.strip_prefix("at://")?; - rest.split('/').next() +fn source_authority_did(source: &AtUri) -> Option<&str> { + let rest = source.as_ref().strip_prefix("at://")?; + let end = rest.find('/').unwrap_or(rest.len()); + let candidate = &rest[..end]; + candidate.starts_with("did:").then_some(candidate) } fn bump_ref(refs: &mut HashMap, author: AuthorId) { @@ -482,8 +484,12 @@ mod tests { AtUri::new_owned(s).unwrap() } + fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() + } + fn did_subj(s: &str) -> SubjectRef { - SubjectRef::Did(Did::new_owned(s).unwrap()) + SubjectRef::Did(did(s)) } fn limit(n: u32) -> PageLimit { @@ -492,15 +498,15 @@ mod tests { const NAMES: [&str; 5] = ["nel", "olaren", "teq", "lyna", "bailey"]; - fn star_edge(source: &str, subject_did: &str) -> Edge { - star_edge_at(source, subject_did, 0) + fn star_edge(source: AtUri, subject: Did) -> Edge { + star_edge_at(source, subject, 0) } - fn star_edge_at(source: &str, subject_did: &str, sort_micros: u64) -> Edge { + fn star_edge_at(source: AtUri, subject: Did, sort_micros: u64) -> Edge { Edge { kind: nsid("sh.tangled.feed.star"), - subject: did_subj(subject_did), - source: at(source), + subject: SubjectRef::Did(subject), + source, sort_micros, } } @@ -511,16 +517,16 @@ mod tests { let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); store.add(star_edge( - "at://did:plc:nel/sh.tangled.feed.star/r1", - "did:plc:abalone", + at("at://did:plc:nel/sh.tangled.feed.star/r1"), + did("did:plc:abalone"), )); store.add(star_edge( - "at://did:plc:olaren/sh.tangled.feed.star/r2", - "did:plc:abalone", + at("at://did:plc:olaren/sh.tangled.feed.star/r2"), + did("did:plc:abalone"), )); store.add(star_edge( - "at://did:plc:nel/sh.tangled.feed.star/r3", - "did:plc:abalone", + at("at://did:plc:nel/sh.tangled.feed.star/r3"), + did("did:plc:abalone"), )); assert_eq!(store.count(&key), 3); @@ -532,8 +538,8 @@ mod tests { let store = store(); let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); let edge = star_edge( - "at://did:plc:nel/sh.tangled.feed.star/r1", - "did:plc:abalone", + at("at://did:plc:nel/sh.tangled.feed.star/r1"), + did("did:plc:abalone"), ); store.add(edge.clone()); store.add(edge); @@ -610,8 +616,8 @@ mod tests { let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); (0..5).for_each(|i| { store.add(star_edge_at( - &format!("at://did:plc:{}/sh.tangled.feed.star/r{i}", NAMES[i]), - "did:plc:abalone", + at(&format!("at://did:plc:{}/sh.tangled.feed.star/r{i}", NAMES[i])), + did("did:plc:abalone"), 1_000_000 + i as u64 * 1_000_000, )); }); @@ -638,8 +644,8 @@ mod tests { let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); (0..2).for_each(|i| { store.add(star_edge( - &format!("at://did:plc:{}/sh.tangled.feed.star/r{i}", NAMES[i]), - "did:plc:abalone", + at(&format!("at://did:plc:{}/sh.tangled.feed.star/r{i}", NAMES[i])), + did("did:plc:abalone"), )); }); @@ -664,19 +670,19 @@ mod tests { fn distinct_authors_decreases_when_last_source_from_author_removed() { let store = store(); let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); - let s1 = "at://did:plc:nel/sh.tangled.feed.star/r1"; - let s2 = "at://did:plc:nel/sh.tangled.feed.star/r2"; - let s3 = "at://did:plc:olaren/sh.tangled.feed.star/r3"; + let s1 = at("at://did:plc:nel/sh.tangled.feed.star/r1"); + let s2 = at("at://did:plc:nel/sh.tangled.feed.star/r2"); + let s3 = at("at://did:plc:olaren/sh.tangled.feed.star/r3"); - store.add(star_edge(s1, "did:plc:abalone")); - store.add(star_edge(s2, "did:plc:abalone")); - store.add(star_edge(s3, "did:plc:abalone")); + store.add(star_edge(s1.clone(), did("did:plc:abalone"))); + store.add(star_edge(s2.clone(), did("did:plc:abalone"))); + store.add(star_edge(s3, did("did:plc:abalone"))); assert_eq!(store.count_distinct_authors(&key), 2); - store.remove_source(&at(s1)); + store.remove_source(&s1); assert_eq!(store.count_distinct_authors(&key), 2, "user1 still has s2"); - store.remove_source(&at(s2)); + store.remove_source(&s2); assert_eq!(store.count_distinct_authors(&key), 1, "user1 fully gone"); } @@ -743,12 +749,12 @@ mod tests { #[test] fn list_does_not_drop_entries_at_same_sort_micros_across_pages() { let store = store(); - let subject = "did:plc:limpet"; - let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj(subject)); + let subject_did = did("did:plc:limpet"); + let key = EdgeKey::new(nsid("sh.tangled.feed.star"), SubjectRef::Did(subject_did.clone())); (0..4).for_each(|i| { store.add(star_edge_at( - &format!("at://did:plc:{}/sh.tangled.feed.star/r{i}", NAMES[i]), - subject, + at(&format!("at://did:plc:{}/sh.tangled.feed.star/r{i}", NAMES[i])), + subject_did.clone(), 42, )); }); @@ -776,19 +782,19 @@ mod tests { #[test] fn list_filtered_narrows_by_predicate_and_paginates() { let store = store(); - let subject = "did:plc:limpet"; - let key = EdgeKey::new(nsid("sh.tangled.repo.issue"), did_subj(subject)); + let subject_did = did("did:plc:limpet"); + let key = EdgeKey::new(nsid("sh.tangled.repo.issue"), SubjectRef::Did(subject_did.clone())); let by_nel = (0..3).map(|i| { star_edge_at( - &format!("at://did:plc:nel/sh.tangled.repo.issue/r{i}"), - subject, + at(&format!("at://did:plc:nel/sh.tangled.repo.issue/r{i}")), + subject_did.clone(), 100 + i as u64, ) }); let by_olaren = (0..2).map(|i| { star_edge_at( - &format!("at://did:plc:olaren/sh.tangled.repo.issue/o{i}"), - subject, + at(&format!("at://did:plc:olaren/sh.tangled.repo.issue/o{i}")), + subject_did.clone(), 500 + i as u64, ) }); diff --git a/crates/ingest/benches/json_decode.rs b/crates/ingest/benches/json_decode.rs index 40717ed..978e12f 100644 --- a/crates/ingest/benches/json_decode.rs +++ b/crates/ingest/benches/json_decode.rs @@ -2,6 +2,8 @@ use std::hint::black_box; use bobbin_types::edges::Record; use criterion::{Criterion, criterion_group, criterion_main}; +use jacquard_common::DefaultStr; +use jacquard_common::types::nsid::Nsid; use serde::Deserialize; use serde_json::value::RawValue; @@ -12,7 +14,7 @@ struct ValueFrame { #[derive(Deserialize)] struct ValueRecordFrame { - collection: String, + collection: Nsid, #[serde(default)] record: Option, } @@ -25,7 +27,7 @@ struct BorrowedRawFrame<'a> { #[derive(Deserialize)] struct BorrowedRawRecordFrame<'a> { - collection: String, + collection: Nsid, #[serde(borrow, default)] record: Option<&'a RawValue>, } @@ -37,7 +39,7 @@ struct OwnedRawFrame { #[derive(Deserialize)] struct OwnedRawRecordFrame { - collection: String, + collection: Nsid, #[serde(default)] record: Option>, } @@ -46,24 +48,22 @@ fn current_path(text: &str) -> Record { let frame: ValueFrame = serde_json::from_str(text).expect("frame"); let value = frame.record.record.expect("upsert body"); let bytes = serde_json::to_vec(&value).expect("re-serialize Value"); - Record::from_json_bytes(frame.record.collection.as_str(), &bytes).expect("decode record") + Record::from_json_bytes(&frame.record.collection, &bytes).expect("decode record") } fn raw_value_borrowed(text: &str) -> Record { let frame: BorrowedRawFrame<'_> = serde_json::from_str(text).expect("frame"); let raw = frame.record.record.expect("upsert body"); - Record::from_json_bytes(frame.record.collection.as_str(), raw.get().as_bytes()) - .expect("decode record") + Record::from_json_bytes(&frame.record.collection, raw.get().as_bytes()).expect("decode record") } fn raw_value_owned(text: &str) -> Record { let frame: OwnedRawFrame = serde_json::from_str(text).expect("frame"); let raw = frame.record.record.expect("upsert body"); - Record::from_json_bytes(frame.record.collection.as_str(), raw.get().as_bytes()) - .expect("decode record") + Record::from_json_bytes(&frame.record.collection, raw.get().as_bytes()).expect("decode record") } -fn floor_record_only(collection: &str, record_bytes: &[u8]) -> Record { +fn floor_record_only(collection: &Nsid, record_bytes: &[u8]) -> Record { Record::from_json_bytes(collection, record_bytes).expect("decode record") } @@ -189,7 +189,7 @@ fn follow_text() -> &'static str { }"# } -fn extract_record_slice(text: &str) -> (String, Vec) { +fn extract_record_slice(text: &str) -> (Nsid, Vec) { let frame: ValueFrame = serde_json::from_str(text).expect("frame"); let value = frame.record.record.expect("upsert body"); let bytes = serde_json::to_vec(&value).expect("re-serialize Value"); @@ -230,10 +230,9 @@ fn bench_decode(c: &mut Criterion) { }); group.bench_function("floor_record_only", |b| { - let collection = collection.as_str(); let bytes = record_bytes.as_slice(); b.iter(|| { - let r = floor_record_only(black_box(collection), black_box(bytes)); + let r = floor_record_only(black_box(&collection), black_box(bytes)); black_box(r); }); }); diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index 404675f..bc0102b 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -1003,7 +1003,7 @@ async fn prepare_record( }; let wire_bytes = Bytes::copy_from_slice(raw.get().as_bytes()); let (parsed, bytes) = match decode_canon_or_upgrade_bytes( - record.collection.as_ref(), + &record.collection, &wire_bytes, ctx.resolver, ) @@ -1525,6 +1525,10 @@ mod tests { SubjectRef::Uri(AtUri::new_owned(s).unwrap()) } + fn rkey(s: &str) -> Rkey { + Rkey::new_owned(s).unwrap() + } + #[allow(clippy::type_complexity)] fn fresh() -> ( Arc, @@ -1663,7 +1667,7 @@ mod tests { #[tokio::test] async fn update_replaces_prior_edges() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); - let mk = |subject_did: &str, id: u64| -> HydrantFrame { + let mk = |subject_did: &Did, id: u64| -> HydrantFrame { parse_frame(json!({ "id": id, "type": "record", @@ -1677,13 +1681,13 @@ mod tests { "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", - "subjectDid": subject_did + "subjectDid": subject_did.as_ref() } } })) }; handle_frame( - mk("did:plc:abalone", 1), + mk(&Did::new_owned("did:plc:abalone").unwrap(), 1), &store, &issue_states, &pull_statuses, @@ -1696,7 +1700,7 @@ mod tests { ) .await; handle_frame( - mk("did:plc:uni", 2), + mk(&Did::new_owned("did:plc:uni").unwrap(), 2), &store, &issue_states, &pull_statuses, @@ -2101,7 +2105,7 @@ mod tests { } } - fn star_frame_text(id: u64, rkey: &str) -> String { + fn star_frame_text(id: u64, rkey: &Rkey) -> String { json!({ "id": id, "type": "record", @@ -2110,7 +2114,7 @@ mod tests { "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", - "rkey": rkey, + "rkey": rkey.as_ref(), "action": "create", "record": { "$type": "sh.tangled.feed.star", @@ -2160,7 +2164,7 @@ mod tests { )); ws_tx - .send(Ok(WsMessage::Text(star_frame_text(1, "starrkeyaa001")))) + .send(Ok(WsMessage::Text(star_frame_text(1, &rkey("starrkeyaa001"))))) .await .unwrap(); ws_tx.send(Ok(WsMessage::Pong(Bytes::new()))).await.unwrap(); @@ -2227,11 +2231,11 @@ mod tests { )); ws_tx - .send(Ok(WsMessage::Text(star_frame_text(1, "heldrkeyaa001")))) + .send(Ok(WsMessage::Text(star_frame_text(1, &rkey("heldrkeyaa001"))))) .await .unwrap(); ws_tx - .send(Ok(WsMessage::Text(star_frame_text(2, "heldrkeyaa002")))) + .send(Ok(WsMessage::Text(star_frame_text(2, &rkey("heldrkeyaa002"))))) .await .unwrap(); @@ -2978,15 +2982,11 @@ mod tests { let owner = Did::new_owned("did:plc:nel").unwrap(); let abalone = Did::new_owned("did:plc:abalone").unwrap(); - let fast_rkeys = ["fastrkeyaa01", "fastrkeyaa02"]; - let slow_rkeys = ["slowrkeyaa01", "slowrkeyaa02"]; + let fast_rkeys: [Rkey; 2] = [rkey("fastrkeyaa01"), rkey("fastrkeyaa02")]; + let slow_rkeys: [Rkey; 2] = [rkey("slowrkeyaa01"), rkey("slowrkeyaa02")]; for r in &fast_rkeys { resolver - .observe( - owner.clone(), - Rkey::new_owned(r).unwrap(), - Some(abalone.clone()), - ) + .observe(owner.clone(), r.clone(), Some(abalone.clone())) .await; } @@ -3037,7 +3037,7 @@ mod tests { .await; }); - let mk = |id: u64, idx: u64, repo_rkey: &str| -> HydrantFrame { + let mk = |id: u64, idx: u64, repo_rkey: &Rkey| -> HydrantFrame { parse_frame(json!({ "id": id, "type": "record", @@ -3052,16 +3052,16 @@ mod tests { "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", - "subject": format!("at://did:plc:nel/sh.tangled.repo/{repo_rkey}"), + "subject": format!("at://did:plc:nel/sh.tangled.repo/{}", repo_rkey.as_ref()), } } })) }; - frame_tx.send(mk(1, 1, slow_rkeys[0])).await.unwrap(); - frame_tx.send(mk(2, 2, fast_rkeys[0])).await.unwrap(); - frame_tx.send(mk(3, 3, slow_rkeys[1])).await.unwrap(); - frame_tx.send(mk(4, 4, fast_rkeys[1])).await.unwrap(); + frame_tx.send(mk(1, 1, &slow_rkeys[0])).await.unwrap(); + frame_tx.send(mk(2, 2, &fast_rkeys[0])).await.unwrap(); + frame_tx.send(mk(3, 3, &slow_rkeys[1])).await.unwrap(); + frame_tx.send(mk(4, 4, &fast_rkeys[1])).await.unwrap(); drop(frame_tx); pipeline.await.unwrap(); diff --git a/crates/ingest/src/warming.rs b/crates/ingest/src/warming.rs index 63aa454..1ea6210 100644 --- a/crates/ingest/src/warming.rs +++ b/crates/ingest/src/warming.rs @@ -328,7 +328,7 @@ mod tests { fn make_upsert(source: &str) -> ParkedUpsert { let source_uri = at(source); let star_json = br#"{"$type":"sh.tangled.feed.star","createdAt":"2026-05-01T00:00:00Z","subject":{"$type":"sh.tangled.feed.star#repo","did":"did:plc:abalone"}}"#; - let parsed = Record::from_json_bytes("sh.tangled.feed.star", star_json) + let parsed = Record::from_json_bytes(&nsid_static("sh.tangled.feed.star"), star_json) .expect("star fixture parses"); ParkedUpsert { cursor: HydrantCursor::new(1), diff --git a/crates/knot-proxy/Cargo.toml b/crates/knot-proxy/Cargo.toml index c946974..f2dc6f6 100644 --- a/crates/knot-proxy/Cargo.toml +++ b/crates/knot-proxy/Cargo.toml @@ -10,6 +10,7 @@ bobbin-runtime = { workspace = true } bytes = { workspace = true } futures = { workspace = true } http = { workspace = true } +jacquard-common = { workspace = true } reqwest = { workspace = true } scc = { workspace = true } thiserror = { workspace = true } diff --git a/crates/knot-proxy/src/host.rs b/crates/knot-proxy/src/host.rs index 16c009f..7374c6a 100644 --- a/crates/knot-proxy/src/host.rs +++ b/crates/knot-proxy/src/host.rs @@ -1,5 +1,8 @@ use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; +use jacquard_common::BosStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::nsid::Nsid; use thiserror::Error; use url::{Host, Url}; @@ -49,9 +52,9 @@ impl KnotHost { &self.0 } - pub fn xrpc_url(&self, nsid: &str) -> Url { + pub fn xrpc_url>(&self, nsid: &Nsid) -> Url { let mut url = self.0.clone(); - url.set_path(&format!("/xrpc/{nsid}")); + url.set_path(&format!("/xrpc/{}", nsid.as_ref())); url } @@ -164,12 +167,6 @@ pub struct RepoSlug(String); #[derive(Clone, Debug, Error)] pub enum RepoSlugError { - #[error("empty did")] - EmptyDid, - #[error("did missing did: prefix: {0}")] - DidPrefixMissing(String), - #[error("did contains slash: {0}")] - DidHasSlash(String), #[error("empty repo name")] EmptyName, #[error("repo name contains slash: {0}")] @@ -177,23 +174,14 @@ pub enum RepoSlugError { } impl RepoSlug { - pub fn new(did: &str, name: &str) -> Result { - if did.is_empty() { - return Err(RepoSlugError::EmptyDid); - } - if !did.starts_with("did:") { - return Err(RepoSlugError::DidPrefixMissing(did.to_owned())); - } - if did.contains('/') { - return Err(RepoSlugError::DidHasSlash(did.to_owned())); - } + pub fn new>(did: &Did, name: &str) -> Result { if name.is_empty() { return Err(RepoSlugError::EmptyName); } if name.contains('/') { return Err(RepoSlugError::NameHasSlash(name.to_owned())); } - Ok(Self(format!("{did}/{name}"))) + Ok(Self(format!("{}/{}", did.as_ref(), name))) } pub fn as_str(&self) -> &str { @@ -204,6 +192,15 @@ impl RepoSlug { #[cfg(test)] mod tests { use super::*; + use jacquard_common::DefaultStr; + + fn did(s: &'static str) -> Did { + Did::new_static(s).unwrap() + } + + fn nsid(s: &'static str) -> Nsid { + Nsid::new_static(s).unwrap() + } #[test] fn parses_bare_host_as_https() { @@ -255,7 +252,7 @@ mod tests { #[test] fn xrpc_url_appends_nsid() { let host = KnotHost::parse("oyster.cafe").unwrap(); - let url = host.xrpc_url("sh.tangled.repo.blob"); + let url = host.xrpc_url(&nsid("sh.tangled.repo.blob")); assert_eq!( url.as_str(), "https://oyster.cafe/xrpc/sh.tangled.repo.blob", @@ -264,46 +261,22 @@ mod tests { #[test] fn slug_joins_did_and_name() { - let slug = RepoSlug::new("did:plc:abalone", "barnacle").unwrap(); - assert_eq!(slug.as_str(), "did:plc:abalone/barnacle"); + let slug = RepoSlug::new(&did("did:plc:squid"), "barnacle").unwrap(); + assert_eq!(slug.as_str(), "did:plc:squid/barnacle"); } #[test] fn slug_rejects_slash_in_name() { assert!(matches!( - RepoSlug::new("did:plc:abalone", "bad/name").unwrap_err(), + RepoSlug::new(&did("did:plc:squid"), "bad/name").unwrap_err(), RepoSlugError::NameHasSlash(s) if s == "bad/name", )); } - #[test] - fn slug_rejects_did_without_prefix() { - assert!(matches!( - RepoSlug::new("plc:abalone", "barnacle").unwrap_err(), - RepoSlugError::DidPrefixMissing(s) if s == "plc:abalone", - )); - } - - #[test] - fn slug_rejects_did_with_slash() { - assert!(matches!( - RepoSlug::new("did:plc:abalone/extra", "barnacle").unwrap_err(), - RepoSlugError::DidHasSlash(s) if s == "did:plc:abalone/extra", - )); - } - - #[test] - fn slug_rejects_empty_did() { - assert!(matches!( - RepoSlug::new("", "barnacle").unwrap_err(), - RepoSlugError::EmptyDid, - )); - } - #[test] fn slug_rejects_empty_name() { assert!(matches!( - RepoSlug::new("did:plc:abalone", "").unwrap_err(), + RepoSlug::new(&did("did:plc:squid"), "").unwrap_err(), RepoSlugError::EmptyName, )); } diff --git a/crates/knot-proxy/src/lib.rs b/crates/knot-proxy/src/lib.rs index 097a653..a9119f6 100644 --- a/crates/knot-proxy/src/lib.rs +++ b/crates/knot-proxy/src/lib.rs @@ -10,6 +10,8 @@ use bobbin_runtime::{ use bytes::Bytes; use futures::Stream; use http::{HeaderMap, StatusCode}; +use jacquard_common::BosStr; +use jacquard_common::types::nsid::Nsid; use reqwest::{Client, redirect::Policy}; use scc::HashMap as SccMap; use thiserror::Error; @@ -143,10 +145,10 @@ impl KnotProxy { self.require_https } - pub async fn forward( + pub async fn forward>( &self, host: &KnotHost, - nsid: &str, + nsid: &Nsid, query: &[(&str, &str)], headers: HeaderMap, ) -> Result { @@ -267,7 +269,11 @@ impl Stream for BodyStream { } } -fn build_xrpc_url(host: &KnotHost, nsid: &str, query: &[(&str, &str)]) -> Url { +fn build_xrpc_url>( + host: &KnotHost, + nsid: &Nsid, + query: &[(&str, &str)], +) -> Url { let mut url = host.xrpc_url(nsid); { let mut pairs = url.query_pairs_mut(); @@ -325,10 +331,15 @@ mod tests { use super::*; use bobbin_runtime::SystemClock; use futures::stream::TryStreamExt; + use jacquard_common::DefaultStr; use tokio::io::AsyncWriteExt; use wiremock::matchers::{method, path, query_param}; use wiremock::{Mock, MockServer, ResponseTemplate}; + fn nsid(s: &'static str) -> Nsid { + Nsid::new_static(s).unwrap() + } + pub(crate) fn config_for_test() -> KnotProxyConfig { KnotProxyConfig { failure_threshold: FailureThreshold::new(2).unwrap(), @@ -376,7 +387,7 @@ mod tests { let server = server().await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.repo.blob")) - .and(query_param("repo", "did:plc:abalone/barnacle")) + .and(query_param("repo", "did:plc:squid/barnacle")) .and(query_param("ref", "main")) .and(query_param("path", "README.md")) .respond_with( @@ -391,9 +402,9 @@ mod tests { let resp = proxy .forward( &host_of(&server), - "sh.tangled.repo.blob", + &nsid("sh.tangled.repo.blob"), &[ - ("repo", "did:plc:abalone/barnacle"), + ("repo", "did:plc:squid/barnacle"), ("ref", "main"), ("path", "README.md"), ], @@ -418,15 +429,15 @@ mod tests { let proxy = proxy_for_test(); let host = host_of(&server); let r1 = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert!(matches!(r1, Err(KnotProxyError::Upstream(_)))); let r2 = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert!(matches!(r2, Err(KnotProxyError::Upstream(_)))); let r3 = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert!(matches!(r3, Err(KnotProxyError::CircuitOpen))); } @@ -443,15 +454,15 @@ mod tests { let proxy = proxy_for_test(); let host = host_of(&server); let r1 = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert_eq!(r1.unwrap().status(), 404); let r2 = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert_eq!(r2.unwrap().status(), 404); let r3 = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert_eq!( r3.unwrap().status(), @@ -482,20 +493,20 @@ mod tests { let proxy = proxy_for_test(); let host = host_of(&server); let _ = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; let _ = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert!(matches!( proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await, Err(KnotProxyError::CircuitOpen), )); tokio::time::sleep(Duration::from_millis(120)).await; let recovered = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await .expect("must recover after cooldown"); assert_eq!(recovered.status(), 200); @@ -526,19 +537,19 @@ mod tests { let bad_host = host_of(&bad); let good_host = host_of(&good); let _ = proxy - .forward(&bad_host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&bad_host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; let _ = proxy - .forward(&bad_host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&bad_host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert!(matches!( proxy - .forward(&bad_host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&bad_host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await, Err(KnotProxyError::CircuitOpen), )); let resp = proxy - .forward(&good_host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&good_host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await .expect("healthy host stays open"); assert_eq!(resp.status(), 200); @@ -549,12 +560,12 @@ mod tests { let host = KnotHost::parse("https://oyster.cafe").unwrap(); let url = build_xrpc_url( &host, - "sh.tangled.repo.tree", - &[("repo", "did:plc:abalone/barnacle"), ("ref", "main")], + &nsid("sh.tangled.repo.tree"), + &[("repo", "did:plc:squid/barnacle"), ("ref", "main")], ); assert_eq!( url.as_str(), - "https://oyster.cafe/xrpc/sh.tangled.repo.tree?repo=did%3Aplc%3Aabalone%2Fbarnacle&ref=main", + "https://oyster.cafe/xrpc/sh.tangled.repo.tree?repo=did%3Aplc%3Asquid%2Fbarnacle&ref=main", ); } @@ -575,7 +586,7 @@ mod tests { let err = proxy .forward( &host_of(&server), - "sh.tangled.repo.blob", + &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new(), ) @@ -601,7 +612,7 @@ mod tests { let err = proxy .forward( &host_of(&server), - "sh.tangled.repo.blob", + &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new(), ) @@ -622,18 +633,18 @@ mod tests { let proxy = proxy_for_test(); let r1 = proxy - .forward(&dead, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&dead, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert!( r1.is_err(), "transport must fail against closed port: {r1:?}" ); let r2 = proxy - .forward(&dead, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&dead, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert!(r2.is_err(), "second transport must fail: {r2:?}"); let r3 = proxy - .forward(&dead, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&dead, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert!( matches!(r3, Err(KnotProxyError::CircuitOpen)), @@ -663,7 +674,7 @@ mod tests { let err = proxy .forward( &host_of(&primary), - "sh.tangled.repo.blob", + &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new(), ) @@ -689,7 +700,7 @@ mod tests { let resp = proxy .forward( &host_of(&server), - "sh.tangled.repo.blob", + &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new(), ) @@ -724,20 +735,20 @@ mod tests { let proxy = proxy_for_test(); let r1 = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await .expect("headers arrive even when body is truncated"); assert_eq!(r1.status(), 200); let _ = drain(r1.into_body_stream()).await; let r2 = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await .expect("second call still gets headers"); let _ = drain(r2.into_body_stream()).await; let r3 = proxy - .forward(&host, "sh.tangled.repo.blob", &[], HeaderMap::new()) + .forward(&host, &nsid("sh.tangled.repo.blob"), &[], HeaderMap::new()) .await; assert!( matches!(r3, Err(KnotProxyError::CircuitOpen)), diff --git a/crates/record-lru/src/lib.rs b/crates/record-lru/src/lib.rs index 0141387..ea9b41b 100644 --- a/crates/record-lru/src/lib.rs +++ b/crates/record-lru/src/lib.rs @@ -82,12 +82,12 @@ mod tests { AtUri::new_owned(s).unwrap() } - fn body(uri: &str, payload: &str) -> Arc { + fn body(uri: AtUri, payload: &str) -> Arc { let cid: Cid = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i" .parse() .unwrap(); Arc::new(RecordBody { - uri: at(uri), + uri, cid, value: Bytes::copy_from_slice(payload.as_bytes()), }) @@ -97,7 +97,7 @@ mod tests { fn put_then_get_round_trips() { let store = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); let uri = at("at://did:plc:nel/sh.tangled.repo/r1"); - let b = body(uri.as_ref(), "{\"hi\":1}"); + let b = body(uri.clone(), "{\"hi\":1}"); store.put(uri.clone(), b.clone()); let got = store.get(&uri).expect("hit"); assert_eq!(got.value, b.value); @@ -107,7 +107,7 @@ mod tests { fn remove_evicts_entry() { let store = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); let uri = at("at://did:plc:nel/sh.tangled.repo/r1"); - store.put(uri.clone(), body(uri.as_ref(), "{\"v\":1}")); + store.put(uri.clone(), body(uri.clone(), "{\"v\":1}")); assert!(store.get(&uri).is_some()); store.remove(&uri); assert!(store.get(&uri).is_none()); @@ -134,14 +134,14 @@ mod tests { let store = LruRecordStore::new(CacheCapacity::from_bytes(2_048)); let big = "x".repeat(900); let kept = (0..16).fold(0usize, |kept, i| { - let uri = format!("at://did:plc:olaren/sh.tangled.string/r{i}"); - store.put(at(&uri), body(&uri, &big)); - kept + usize::from(store.get(&at(&uri)).is_some()) + let uri = at(&format!("at://did:plc:olaren/sh.tangled.string/r{i}")); + store.put(uri.clone(), body(uri.clone(), &big)); + kept + usize::from(store.get(&uri).is_some()) }); let resident = (0..16) .filter(|i| { - let uri = format!("at://did:plc:olaren/sh.tangled.string/r{i}"); - store.get(&at(&uri)).is_some() + let uri = at(&format!("at://did:plc:olaren/sh.tangled.string/r{i}")); + store.get(&uri).is_some() }) .count(); assert!( @@ -162,7 +162,7 @@ mod tests { #[test] fn weight_includes_uri_value_and_cid() { - let b = body("at://did:plc:nel/sh.tangled.string/k", "{\"a\":1}"); + let b = body(at("at://did:plc:nel/sh.tangled.string/k"), "{\"a\":1}"); let baseline = b.value.len() as u64 + b.uri.as_ref().len() as u64 + b.cid.as_ref().len() as u64; assert!( diff --git a/crates/resolver/src/legacy_upgrade.rs b/crates/resolver/src/legacy_upgrade.rs index e958ebf..17e363b 100644 --- a/crates/resolver/src/legacy_upgrade.rs +++ b/crates/resolver/src/legacy_upgrade.rs @@ -8,9 +8,10 @@ use bobbin_types::sh_tangled::git::ref_update::RefUpdate; use bobbin_types::sh_tangled::repo::collaborator::Collaborator; use bobbin_types::sh_tangled::repo::issue::Issue; use bobbin_types::sh_tangled::repo::pull::{Pull, Source, Target}; -use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; +use jacquard_common::types::nsid::Nsid; use jacquard_common::types::string::AtUri; +use jacquard_common::{BosStr, DefaultStr}; use crate::normalize::{is_repo_at_uri, resolve_repo_uri}; use crate::{RepoIdResolver, Resolution}; @@ -25,7 +26,10 @@ pub enum DecodedRecord { } impl DecodedRecord { - pub fn try_decode(nsid: &str, bytes: &[u8]) -> Result { + pub fn try_decode>( + nsid: &Nsid, + bytes: &[u8], + ) -> Result { match Record::from_json_bytes(nsid, bytes) { Ok(record) => Ok(Self::Canon(record)), Err(canon_err) => match LegacyRecord::from_json_bytes(nsid, bytes) { @@ -48,22 +52,23 @@ async fn upgrade_repo_did( resolve_repo_uri(resolver, &uri).await } -pub async fn upgrade_wire_bytes( - nsid: &str, +pub async fn upgrade_wire_bytes>( + nsid: &Nsid, bytes: &[u8], resolver: &RepoIdResolver, ) -> Result, ExtractError> { let legacy = LegacyRecord::from_json_bytes(nsid, bytes)?; let Some(canon) = upgrade(legacy, resolver).await else { return Err(ExtractError::UnknownCollection(alloc::format!( - "{nsid}: legacy upgrade failed" + "{}: legacy upgrade failed", + nsid.as_ref() ))); }; serialize_canon_variant(&canon).map_err(ExtractError::DecodeJson) } -pub async fn decode_canon_or_upgrade_bytes<'a>( - nsid: &str, +pub async fn decode_canon_or_upgrade_bytes<'a, S: BosStr + AsRef>( + nsid: &Nsid, bytes: &'a [u8], resolver: &RepoIdResolver, ) -> Result<(Record, alloc::borrow::Cow<'a, [u8]>), ExtractError> { @@ -73,7 +78,8 @@ pub async fn decode_canon_or_upgrade_bytes<'a>( DecodedRecord::Legacy(legacy) => { let Some(canon) = upgrade(legacy, resolver).await else { return Err(ExtractError::UnknownCollection(alloc::format!( - "{nsid}: legacy upgrade failed" + "{}: legacy upgrade failed", + nsid.as_ref() ))); }; let canon_bytes = serialize_canon_variant(&canon).map_err(ExtractError::DecodeJson)?; @@ -263,10 +269,14 @@ mod tests { Rkey::new_owned(s).unwrap() } + fn nsid(s: &'static str) -> Nsid { + Nsid::new_static(s).unwrap() + } + #[test] fn legacy_decode_routes_through_try_decode_for_known_nsids() { - let json = br#"{"$type":"sh.tangled.repo.issue","repoDid":"did:plc:abalone","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; - let decoded = DecodedRecord::try_decode("sh.tangled.repo.issue", json) + let json = br#"{"$type":"sh.tangled.repo.issue","repoDid":"did:plc:squid","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; + let decoded = DecodedRecord::try_decode(&nsid("sh.tangled.repo.issue"), json) .expect("legacy issue must decode"); assert!(matches!( decoded, @@ -276,12 +286,12 @@ mod tests { #[test] fn canon_decode_wins_when_wire_matches_new_shape() { - let json = br#"{"$type":"sh.tangled.repo.issue","repo":"did:plc:abalone","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; - let decoded = DecodedRecord::try_decode("sh.tangled.repo.issue", json) + let json = br#"{"$type":"sh.tangled.repo.issue","repo":"did:plc:squid","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; + let decoded = DecodedRecord::try_decode(&nsid("sh.tangled.repo.issue"), json) .expect("canon issue must decode"); match decoded { DecodedRecord::Canon(Record::Issue(i)) => { - assert_eq!(i.repo.as_ref(), "did:plc:abalone") + assert_eq!(i.repo.as_ref(), "did:plc:squid") } other => panic!("expected canon issue, got {other:?}"), } @@ -290,7 +300,7 @@ mod tests { #[test] fn legacy_decode_passes_through_for_unaffected_nsids() { let json = br#"{"$type":"sh.tangled.graph.follow","subject":"did:plc:bailey","createdAt":"2026-05-01T00:00:00Z"}"#; - let decoded = DecodedRecord::try_decode("sh.tangled.graph.follow", json) + let decoded = DecodedRecord::try_decode(&nsid("sh.tangled.graph.follow"), json) .expect("follow has no legacy form, must decode canon"); assert!(matches!(decoded, DecodedRecord::Canon(Record::Follow(_)))); } @@ -299,7 +309,7 @@ mod tests { async fn upgrade_issue_uses_repo_did_directly() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let json = br#"{"$type":"sh.tangled.repo.issue","repoDid":"did:plc:abalone","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; - let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.issue"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Issue(i) => assert_eq!(i.repo, did("did:plc:abalone")), @@ -316,7 +326,7 @@ mod tests { .observe(owner.clone(), key.clone(), Some(did("did:plc:abalone"))) .await; let json = br#"{"$type":"sh.tangled.repo.issue","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; - let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.issue"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Issue(i) => assert_eq!(i.repo, did("did:plc:abalone")), @@ -328,7 +338,7 @@ mod tests { async fn upgrade_issue_drops_when_resolver_cannot_map() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let json = br#"{"$type":"sh.tangled.repo.issue","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; - let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.issue"), json).expect("decode"); assert!( upgrade(legacy, &resolver).await.is_none(), "no resolver entry and no repoDid means the canon Did cannot be constructed", @@ -339,7 +349,7 @@ mod tests { async fn upgrade_pull_propagates_target_resolution() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repoDid":"did:plc:abalone"}}"#; - let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.pull", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.pull"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Pull(p) => { @@ -354,7 +364,7 @@ mod tests { async fn upgrade_pull_source_repo_resolution_is_independent_of_target() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repoDid":"did:plc:abalone"},"source":{"branch":"feat","repo":"at://did:plc:nel/sh.tangled.repo/missing"}}"#; - let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.pull", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.pull"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Pull(p) => { @@ -375,7 +385,7 @@ mod tests { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let json = br#"{"$type":"sh.tangled.git.refUpdate","ref":"refs/heads/main","committerDid":"did:plc:olaren","repoDid":"did:plc:abalone","oldSha":"0000000000000000000000000000000000000000","newSha":"1111111111111111111111111111111111111111","meta":{"isDefaultRef":true,"commitCount":{}}}"#; let legacy = - LegacyRecord::from_json_bytes("sh.tangled.git.refUpdate", json).expect("decode"); + LegacyRecord::from_json_bytes(&nsid("sh.tangled.git.refUpdate"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::RefUpdate(r) => assert_eq!(r.repo, did("did:plc:abalone")), @@ -387,7 +397,7 @@ mod tests { async fn upgrade_star_prefers_subject_did_over_subject_uri() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let json = br#"{"$type":"sh.tangled.feed.star","createdAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:nel/sh.tangled.string/k1","subjectDid":"did:plc:abalone"}"#; - let legacy = LegacyRecord::from_json_bytes("sh.tangled.feed.star", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.feed.star"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Star(s) => match s.subject { @@ -402,7 +412,7 @@ mod tests { async fn upgrade_star_falls_back_to_string_when_repo_uri_not_in_cache() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let json = br#"{"$type":"sh.tangled.feed.star","createdAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"}"#; - let legacy = LegacyRecord::from_json_bytes("sh.tangled.feed.star", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.feed.star"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Star(s) => match s.subject { @@ -426,7 +436,7 @@ mod tests { .observe(owner.clone(), key.clone(), Some(did("did:plc:abalone"))) .await; let json = br#"{"$type":"sh.tangled.feed.star","createdAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"}"#; - let legacy = LegacyRecord::from_json_bytes("sh.tangled.feed.star", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.feed.star"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Star(s) => match s.subject { @@ -442,7 +452,7 @@ mod tests { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let with_did = br#"{"$type":"sh.tangled.repo.collaborator","createdAt":"2026-05-01T00:00:00Z","subject":"did:plc:lyna","repoDid":"did:plc:abalone"}"#; let canon = upgrade( - LegacyRecord::from_json_bytes("sh.tangled.repo.collaborator", with_did) + LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.collaborator"), with_did) .expect("decode"), &resolver, ) @@ -454,8 +464,9 @@ mod tests { } let no_resolution = br#"{"$type":"sh.tangled.repo.collaborator","createdAt":"2026-05-01T00:00:00Z","subject":"did:plc:lyna","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"}"#; - let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.collaborator", no_resolution) - .expect("decode"); + let legacy = + LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.collaborator"), no_resolution) + .expect("decode"); assert!(upgrade(legacy, &resolver).await.is_none()); } } diff --git a/crates/resolver/src/lib.rs b/crates/resolver/src/lib.rs index 2f77003..537c438 100644 --- a/crates/resolver/src/lib.rs +++ b/crates/resolver/src/lib.rs @@ -417,7 +417,7 @@ fn is_garbage_response(err: &SlingshotError) -> bool { } fn repo_did_from_body(body: &[u8]) -> Result>, ExtractError> { - match Record::from_json_bytes(REPO_COLLECTION, body)? { + match Record::from_json_bytes(&nsid_static(REPO_COLLECTION), body)? { Record::Repo(repo) => Ok(repo.repo_did), _ => Ok(None), } diff --git a/crates/search/src/lib.rs b/crates/search/src/lib.rs index 64d50ad..842c22b 100644 --- a/crates/search/src/lib.rs +++ b/crates/search/src/lib.rs @@ -497,10 +497,15 @@ mod tests { NsidType::new_static(s).unwrap() } - fn doc(uri: &str, nsid_s: &'static str, title: &str, body: &str) -> SearchDoc { + fn doc( + uri: AtUri, + nsid: NsidType, + title: &str, + body: &str, + ) -> SearchDoc { SearchDoc { - uri: at(uri), - nsid: nsid(nsid_s), + uri, + nsid, title: title.to_owned(), body: body.to_owned(), author: None, @@ -517,15 +522,15 @@ mod tests { async fn upsert_then_query_returns_matching_doc() { let idx = build(); idx.upsert(doc( - "at://did:plc:nel/sh.tangled.repo.issue/r1", - "sh.tangled.repo.issue", + at("at://did:plc:nel/sh.tangled.repo.issue/r1"), + nsid("sh.tangled.repo.issue"), "barnacle pagination", "scroll resets", )) .await; idx.upsert(doc( - "at://did:plc:teq/sh.tangled.repo.issue/r2", - "sh.tangled.repo.issue", + at("at://did:plc:teq/sh.tangled.repo.issue/r2"), + nsid("sh.tangled.repo.issue"), "kelp grew sideways", "kelp", )) @@ -553,15 +558,15 @@ mod tests { async fn query_filters_by_nsid() { let idx = build(); idx.upsert(doc( - "at://did:plc:nel/sh.tangled.repo.issue/r1", - "sh.tangled.repo.issue", + at("at://did:plc:nel/sh.tangled.repo.issue/r1"), + nsid("sh.tangled.repo.issue"), "anemone tide", "", )) .await; idx.upsert(doc( - "at://did:plc:teq/sh.tangled.string/k1", - "sh.tangled.string", + at("at://did:plc:teq/sh.tangled.string/k1"), + nsid("sh.tangled.string"), "anemone.md", "anemone again", )) @@ -587,10 +592,15 @@ mod tests { #[tokio::test] async fn upsert_replaces_prior_document_for_same_uri() { let idx = build(); - let uri = "at://did:plc:nel/sh.tangled.repo.issue/r1"; - idx.upsert(doc(uri, "sh.tangled.repo.issue", "abalone", "")) - .await; - idx.upsert(doc(uri, "sh.tangled.repo.issue", "limpet", "")) + let uri = at("at://did:plc:nel/sh.tangled.repo.issue/r1"); + idx.upsert(doc( + uri.clone(), + nsid("sh.tangled.repo.issue"), + "abalone", + "", + )) + .await; + idx.upsert(doc(uri, nsid("sh.tangled.repo.issue"), "limpet", "")) .await; idx.flush().await; @@ -611,8 +621,13 @@ mod tests { async fn remove_drops_document_from_index() { let idx = build(); let uri = at("at://did:plc:nel/sh.tangled.repo.issue/r1"); - idx.upsert(doc(uri.as_ref(), "sh.tangled.repo.issue", "whelk", "shell")) - .await; + idx.upsert(doc( + uri.clone(), + nsid("sh.tangled.repo.issue"), + "whelk", + "shell", + )) + .await; idx.remove(&uri).await; idx.flush().await; let hits = idx @@ -628,8 +643,8 @@ mod tests { let names = ["nel", "olaren", "teq", "lyna", "bailey"]; for (i, owner) in names.iter().enumerate() { idx.upsert(doc( - &format!("at://did:plc:{owner}/sh.tangled.repo.issue/r{i}"), - "sh.tangled.repo.issue", + at(&format!("at://did:plc:{owner}/sh.tangled.repo.issue/r{i}")), + nsid("sh.tangled.repo.issue"), "anemone", "tides", )) @@ -671,8 +686,8 @@ mod tests { async fn empty_query_short_circuits_to_empty_page() { let idx = build(); idx.upsert(doc( - "at://did:plc:nel/sh.tangled.repo.issue/r1", - "sh.tangled.repo.issue", + at("at://did:plc:nel/sh.tangled.repo.issue/r1"), + nsid("sh.tangled.repo.issue"), "abalone", "", )) @@ -696,8 +711,8 @@ mod tests { async fn writer_auto_commits_on_idle_interval() { let idx = build(); idx.upsert(doc( - "at://did:plc:nel/sh.tangled.repo.issue/r1", - "sh.tangled.repo.issue", + at("at://did:plc:nel/sh.tangled.repo.issue/r1"), + nsid("sh.tangled.repo.issue"), "auto", "", )) diff --git a/crates/slingshot-client/src/lib.rs b/crates/slingshot-client/src/lib.rs index ece71bd..5aa2526 100644 --- a/crates/slingshot-client/src/lib.rs +++ b/crates/slingshot-client/src/lib.rs @@ -146,7 +146,7 @@ impl SlingshotClient { StatusCode::OK => { let bytes = read_bounded(resp).await?; let body = decode(&bytes)?; - verify_addresses(&body, repo.as_ref(), collection.as_ref(), rkey.as_ref())?; + verify_addresses(&body, repo, collection, rkey)?; Ok(Arc::new(body)) } StatusCode::NOT_FOUND => Err(SlingshotError::NotFound), @@ -196,13 +196,18 @@ fn decode(bytes: &[u8]) -> Result { }) } -fn verify_addresses( +fn verify_addresses( body: &RecordBody, - repo: &str, - collection: &str, - rkey: &str, + repo: &Did, + collection: &Nsid, + rkey: &Rkey, ) -> Result<(), SlingshotError> { - let expected = format!("at://{repo}/{collection}/{rkey}"); + let expected = format!( + "at://{}/{}/{}", + repo.as_ref(), + collection.as_ref(), + rkey.as_ref() + ); if body.uri.as_ref() == expected { Ok(()) } else { diff --git a/crates/types/src/edges.rs b/crates/types/src/edges.rs index cd5c855..089d6c1 100644 --- a/crates/types/src/edges.rs +++ b/crates/types/src/edges.rs @@ -87,16 +87,19 @@ impl Record { value: serde_json::Value, ) -> Result { let bytes = serde_json::to_vec(&value)?; - Self::from_json_bytes(collection.as_ref(), &bytes) + Self::from_json_bytes(collection, &bytes) } - pub fn from_json_bytes(nsid_str: &str, bytes: &[u8]) -> Result { + pub fn from_json_bytes>( + nsid: &Nsid, + bytes: &[u8], + ) -> Result { macro_rules! parse { ($variant:ident) => { Ok(Self::$variant(serde_json::from_slice(bytes)?)) }; } - match nsid_str { + match nsid.as_ref() { "sh.tangled.actor.profile" => parse!(Profile), "sh.tangled.feed.reaction" => parse!(Reaction), "sh.tangled.feed.star" => parse!(Star), @@ -251,7 +254,7 @@ fn repo_subject( } fn uri_subject_for_record(uri: &AtUri) -> Option { - crate::ids::owner_did_from_aturi(uri.as_ref())?; + crate::ids::owner_did_from_aturi(uri)?; Some(SubjectRef::Uri(uri.clone())) } @@ -311,7 +314,7 @@ fn mirror_kind_for(kind: &str) -> Option<&'static str> { } fn append_mirror_edges(primary: Vec, source: &AtUri) -> Vec { - let Some(author) = crate::ids::owner_did_from_aturi(source.as_ref()) else { + let Some(author) = crate::ids::owner_did_from_aturi(source) else { return primary; }; let author_subject = SubjectRef::Did(author); @@ -539,7 +542,7 @@ fn pipeline_edges( } fn owner_self_edges(kind: &'static str, source: &AtUri) -> Vec { - crate::ids::owner_did_from_aturi(source.as_ref()) + crate::ids::owner_did_from_aturi(source) .map(|did| one_edge(kind, SubjectRef::Did(did), source)) .unwrap_or_default() } @@ -706,7 +709,10 @@ mod tests { assert_eq!(edges[0].subject, did_subj("did:plc:abalone")); } - fn artifact_body(repo: Option<&str>, repo_did: Option<&str>) -> serde_json::Value { + fn artifact_body( + repo: Option<&AtUri>, + repo_did: Option<&Did>, + ) -> serde_json::Value { let mut body = json!({ "$type": "sh.tangled.repo.artifact", "createdAt": "2026-05-01T00:00:00Z", @@ -721,10 +727,10 @@ mod tests { }); let obj = body.as_object_mut().expect("artifact_body returns object"); if let Some(r) = repo { - obj.insert("repo".into(), json!(r)); + obj.insert("repo".into(), json!(r.as_ref())); } if let Some(d) = repo_did { - obj.insert("repoDid".into(), json!(d)); + obj.insert("repoDid".into(), json!(d.as_ref())); } body } @@ -735,8 +741,8 @@ mod tests { "sh.tangled.repo.artifact", "at://did:plc:nel/sh.tangled.repo.artifact/abcabcabcabcz", artifact_body( - Some("at://did:plc:abalone/sh.tangled.repo/r1"), - Some("did:plc:lyna"), + Some(&at("at://did:plc:abalone/sh.tangled.repo/r1")), + Some(&did("did:plc:lyna")), ), ); assert_eq!(edges.len(), 1); @@ -748,7 +754,7 @@ mod tests { let edges = extract( "sh.tangled.repo.artifact", "at://did:plc:nel/sh.tangled.repo.artifact/abcabcabcabcz", - artifact_body(None, Some("did:plc:abalone")), + artifact_body(None, Some(&did("did:plc:abalone"))), ); assert_eq!(edges.len(), 1); assert_eq!(edges[0].subject, did_subj("did:plc:abalone")); @@ -769,7 +775,7 @@ mod tests { let edges = extract( "sh.tangled.repo.artifact", "at://did:plc:nel/sh.tangled.repo.artifact/abcabcabcabcz", - artifact_body(Some("at://did:plc:abalone/sh.tangled.repo/r1"), None), + artifact_body(Some(&at("at://did:plc:abalone/sh.tangled.repo/r1")), None), ); assert_eq!(edges.len(), 1); assert_eq!( @@ -783,7 +789,7 @@ mod tests { let edges = extract( "sh.tangled.repo.artifact", "at://did:plc:nel/sh.tangled.repo.artifact/abcabcabcabcz", - artifact_body(Some("at://oyster.cafe/sh.tangled.repo/r1"), None), + artifact_body(Some(&at("at://oyster.cafe/sh.tangled.repo/r1")), None), ); assert!(edges.is_empty()); } diff --git a/crates/types/src/ids.rs b/crates/types/src/ids.rs index 3146e52..0a0564e 100644 --- a/crates/types/src/ids.rs +++ b/crates/types/src/ids.rs @@ -81,10 +81,13 @@ pub fn nsid_static(s: &'static str) -> Nsid { Nsid::new_static(s).expect("compile-time NSID literal must validate") } -pub fn owner_did_from_aturi(uri: &str) -> Option> { - let rest = uri.strip_prefix("at://")?; - let end = rest.find('/').unwrap_or(rest.len()); - Did::new_owned(&rest[..end]).ok() +pub fn owner_did_from_aturi(uri: &AtUri) -> Option> { + use jacquard_common::IntoStatic; + use jacquard_common::types::ident::AtIdentifier; + match uri.authority() { + AtIdentifier::Did(d) => Some(d.into_static()), + AtIdentifier::Handle(_) => None, + } } #[cfg(test)] @@ -95,10 +98,14 @@ mod tests { Did::new_static(s).unwrap() } + fn at(s: &str) -> AtUri { + AtUri::new_owned(s).unwrap() + } + #[test] fn owner_did_from_full_uri() { assert_eq!( - owner_did_from_aturi("at://did:plc:olaren/sh.tangled.feed.star/3lk1"), + owner_did_from_aturi(&at("at://did:plc:olaren/sh.tangled.feed.star/3lk1")), Some(did("did:plc:olaren")), ); } @@ -106,7 +113,7 @@ mod tests { #[test] fn owner_did_from_did_only_uri() { assert_eq!( - owner_did_from_aturi("at://did:plc:olaren"), + owner_did_from_aturi(&at("at://did:plc:olaren")), Some(did("did:plc:olaren")) ); } @@ -114,23 +121,18 @@ mod tests { #[test] fn owner_did_rejects_handle_authority() { assert_eq!( - owner_did_from_aturi("at://oyster.cafe/sh.tangled.repo/r1"), + owner_did_from_aturi(&at("at://oyster.cafe/sh.tangled.repo/r1")), None ); } - #[test] - fn owner_did_rejects_missing_scheme() { - assert_eq!(owner_did_from_aturi("did:plc:nel"), None); - } - #[test] fn subject_ref_did_and_uri_hash_distinct() { use std::collections::HashSet; let mut s = HashSet::new(); - s.insert(SubjectRef::Did(did("did:plc:abalone"))); + s.insert(SubjectRef::Did(did("did:plc:squid"))); s.insert(SubjectRef::Uri( - AtUri::new_owned("at://did:plc:abalone").unwrap(), + AtUri::new_owned("at://did:plc:squid").unwrap(), )); assert_eq!( s.len(), diff --git a/crates/types/src/legacy.rs b/crates/types/src/legacy.rs index ebafaa6..440e210 100644 --- a/crates/types/src/legacy.rs +++ b/crates/types/src/legacy.rs @@ -2,6 +2,7 @@ use alloc::collections::BTreeMap; use alloc::vec::Vec; use jacquard_common::deps::smol_str::SmolStr; +use jacquard_common::types::nsid::Nsid; use jacquard_common::types::string::{AtUri, Datetime, Did}; use jacquard_common::types::value::Data; use jacquard_common::{BosStr, DefaultStr}; @@ -153,8 +154,11 @@ pub enum LegacyRecord { } impl LegacyRecord { - pub fn from_json_bytes(nsid: &str, bytes: &[u8]) -> Result { - match nsid { + pub fn from_json_bytes>( + nsid: &Nsid, + bytes: &[u8], + ) -> Result { + match nsid.as_ref() { "sh.tangled.repo.issue" => Ok(Self::Issue(serde_json::from_slice(bytes)?)), "sh.tangled.repo.pull" => Ok(Self::Pull(serde_json::from_slice(bytes)?)), "sh.tangled.repo.collaborator" => { diff --git a/crates/types/src/search.rs b/crates/types/src/search.rs index 604d4a6..6c0a219 100644 --- a/crates/types/src/search.rs +++ b/crates/types/src/search.rs @@ -5,7 +5,7 @@ use core::future::Future; use jacquard_common::types::ident::AtIdentifier; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::string::{AtUri, Did}; -use jacquard_common::{DefaultStr, IntoStatic}; +use jacquard_common::{BosStr, DefaultStr, IntoStatic}; use serde::Serialize; use crate::edges::{ExtractError, Record}; @@ -86,10 +86,13 @@ impl SearchableRecord { } } - pub fn from_json_bytes(nsid_str: &str, bytes: &[u8]) -> Result { - let parsed = Record::from_json_bytes(nsid_str, bytes)?; + pub fn from_json_bytes>( + nsid: &Nsid, + bytes: &[u8], + ) -> Result { + let parsed = Record::from_json_bytes(nsid, bytes)?; Self::try_from_record(parsed) - .ok_or_else(|| ExtractError::UnknownCollection(nsid_str.into())) + .ok_or_else(|| ExtractError::UnknownCollection(nsid.as_ref().into())) } pub fn nsid(&self) -> Nsid { @@ -292,20 +295,24 @@ mod tests { Nsid::new_static(s).unwrap() } - fn extract(collection: &'static str, source: &str, body: serde_json::Value) -> SearchDoc { - let parsed = Record::from_json_value(&nsid(collection), body).expect("parse"); + fn extract( + collection: &Nsid, + source: &AtUri, + body: serde_json::Value, + ) -> SearchDoc { + let parsed = Record::from_json_value(collection, body).expect("parse"); let searchable = SearchableRecord::try_from_record(parsed).expect("indexable record"); - searchable.to_search_doc(&at(source)) + searchable.to_search_doc(source) } #[test] fn issue_titles_and_body_indexable() { let doc = extract( - "sh.tangled.repo.issue", - "at://did:plc:nel/sh.tangled.repo.issue/abcabcabcabcz", + &nsid("sh.tangled.repo.issue"), + &at("at://did:plc:nel/sh.tangled.repo.issue/abcabcabcabcz"), json!({ "$type": "sh.tangled.repo.issue", - "repo": "did:plc:abalone", + "repo": "did:plc:squid", "title": "barnacle pagination overflow", "body": "scrolling resets when the cursor wraps", "createdAt": "2026-05-01T00:00:00Z" @@ -319,11 +326,11 @@ mod tests { #[test] fn issue_without_body_yields_empty_body() { let doc = extract( - "sh.tangled.repo.issue", - "at://did:plc:nel/sh.tangled.repo.issue/abcabcabcabcz", + &nsid("sh.tangled.repo.issue"), + &at("at://did:plc:nel/sh.tangled.repo.issue/abcabcabcabcz"), json!({ "$type": "sh.tangled.repo.issue", - "repo": "did:plc:abalone", + "repo": "did:plc:squid", "title": "kelp ate my newline", "createdAt": "2026-05-01T00:00:00Z" }), @@ -335,18 +342,18 @@ mod tests { #[test] fn repo_indexes_name_description_topics() { let doc = extract( - "sh.tangled.repo", - "at://did:plc:teq/sh.tangled.repo/r1", + &nsid("sh.tangled.repo"), + &at("at://did:plc:teq/sh.tangled.repo/r1"), json!({ "$type": "sh.tangled.repo", - "name": "abalone", + "name": "scallop", "knot": "oyster.cafe", "description": "shell index for tide pools", "topics": ["intertidal", "molluscs"], "createdAt": "2026-05-01T00:00:00Z" }), ); - assert_eq!(doc.title, "abalone"); + assert_eq!(doc.title, "scallop"); assert!(doc.body.contains("shell index")); assert!(doc.body.contains("intertidal")); assert!(doc.body.contains("molluscs")); @@ -355,8 +362,8 @@ mod tests { #[test] fn repo_without_name_field_falls_back_to_rkey() { let doc = extract( - "sh.tangled.repo", - "at://did:plc:teq/sh.tangled.repo/limpet", + &nsid("sh.tangled.repo"), + &at("at://did:plc:teq/sh.tangled.repo/limpet"), json!({ "$type": "sh.tangled.repo", "knot": "oyster.cafe", @@ -375,7 +382,7 @@ mod tests { "createdAt": "2026-05-01T00:00:00Z", "subject": { "$type": "sh.tangled.feed.star#repo", - "did": "did:plc:abalone" + "did": "did:plc:squid" } }), ) @@ -386,8 +393,8 @@ mod tests { #[test] fn string_indexes_filename_description_contents() { let doc = extract( - "sh.tangled.string", - "at://did:plc:teq/sh.tangled.string/k1", + &nsid("sh.tangled.string"), + &at("at://did:plc:teq/sh.tangled.string/k1"), json!({ "$type": "sh.tangled.string", "filename": "anemone.md", @@ -405,12 +412,13 @@ mod tests { fn from_json_bytes_round_trips_via_serialize() { let body = json!({ "$type": "sh.tangled.repo.issue", - "repo": "did:plc:abalone", + "repo": "did:plc:squid", "title": "uni shell", "createdAt": "2026-05-01T00:00:00Z" }); let bytes = serde_json::to_vec(&body).unwrap(); - let view = SearchableRecord::from_json_bytes("sh.tangled.repo.issue", &bytes).unwrap(); + let view = + SearchableRecord::from_json_bytes(&nsid("sh.tangled.repo.issue"), &bytes).unwrap(); assert_eq!(view.nsid(), nsid("sh.tangled.repo.issue")); let serialized = serde_json::to_value(&view).unwrap(); assert_eq!(serialized["title"], json!("uni shell")); @@ -424,11 +432,11 @@ mod tests { "createdAt": "2026-05-01T00:00:00Z", "subject": { "$type": "sh.tangled.feed.star#repo", - "did": "did:plc:abalone" + "did": "did:plc:squid" } }); let bytes = serde_json::to_vec(&body).unwrap(); - let err = SearchableRecord::from_json_bytes("sh.tangled.feed.star", &bytes) + let err = SearchableRecord::from_json_bytes(&nsid("sh.tangled.feed.star"), &bytes) .expect_err("non-searchable nsid must be rejected at search hydration"); assert!(matches!(err, ExtractError::UnknownCollection(_))); }