From 69a0a7d2dd321ac7af243ec26fced542fa94f8dc Mon Sep 17 00:00:00 2001 From: dawn Date: Fri, 24 Jul 2026 21:56:04 +0300 Subject: [PATCH] bobbin: add sh.tangled.repo.getRepoByName Signed-off-by: dawn --- bobbin/crates/ingest/src/lib.rs | 105 ++++++++++++++++ bobbin/crates/resolver/src/lib.rs | 158 +++++++++++++++++++++++++ bobbin/crates/xrpc/src/lib.rs | 30 +++++ bobbin/crates/xrpc/tests/cold_start.rs | 106 +++++++++++++++++ lexicons/repo/getRepoByName.json | 45 +++++++ 5 files changed, 444 insertions(+) create mode 100644 lexicons/repo/getRepoByName.json diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index 5ec1278c..28517c65 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -1053,6 +1053,15 @@ async fn prepare_record( repo.repo_did.clone(), ) .await; + // urls carry the cosmetic name, older repos only have their rkey + let name = repo + .name + .as_ref() + .map(AsRef::as_ref) + .unwrap_or_else(|| record.rkey.as_ref()); + ctx.resolver + .observe_name(record.did.clone(), record.rkey.clone(), name) + .await; if let Some(prior) = superseded { let prior_uri = format!( "at://{}/sh.tangled.repo/{}", @@ -2985,6 +2994,94 @@ mod tests { ); } + #[tokio::test] + async fn repo_record_indexes_its_name() { + let (store, issue_states, pull_statuses, cov, resolver) = fresh(); + let repo: HydrantFrame = parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:nel", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.repo", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.repo", + "createdAt": "2026-05-01T00:00:00Z", + "knot": "oyster.cafe", + "name": "abalone" + } + } + })); + handle_frame( + repo, + &store, + &issue_states, + &pull_statuses, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + &sys_clock(), + now(), + ) + .await; + let owner = Did::new_owned("did:plc:nel").unwrap(); + assert_eq!( + resolver.lookup_by_name(&owner, "abalone").await, + Some(bobbin_types::ids::RepoIdent::new( + owner, + Rkey::new_owned("abcabcabcabcz").unwrap() + )), + ); + } + + #[tokio::test] + async fn repo_record_without_a_name_indexes_its_rkey() { + let (store, issue_states, pull_statuses, cov, resolver) = fresh(); + let repo: HydrantFrame = parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:nel", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.repo", + "rkey": "abalone", + "action": "create", + "record": { + "$type": "sh.tangled.repo", + "createdAt": "2026-05-01T00:00:00Z", + "knot": "oyster.cafe" + } + } + })); + handle_frame( + repo, + &store, + &issue_states, + &pull_statuses, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + &sys_clock(), + now(), + ) + .await; + let owner = Did::new_owned("did:plc:nel").unwrap(); + assert_eq!( + resolver.lookup_by_name(&owner, "abalone").await, + Some(bobbin_types::ids::RepoIdent::new( + owner, + Rkey::new_owned("abalone").unwrap() + )), + "repos made before the name field are only reachable by rkey", + ); + } + #[tokio::test] async fn delete_repo_record_evicts_resolver_cache() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); @@ -2997,6 +3094,9 @@ mod tests { Some(Did::new_owned("did:plc:abalone").unwrap()), ) .await; + resolver + .observe_name(owner.clone(), rkey.clone(), "abcabcabcabcz") + .await; assert!( resolver.cached_resolution(&owner, &rkey).await.is_some(), "observe must seed the cache", @@ -3030,6 +3130,11 @@ mod tests { resolver.cached_resolution(&owner, &rkey).await.is_none(), "deleting the repo record must clear the resolver cache so future observes are not blocked by a stale Authoritative entry", ); + assert_eq!( + resolver.lookup_by_name(&owner, "abcabcabcabcz").await, + None, + "a deleted repo must stop answering on its url", + ); } #[tokio::test] diff --git a/bobbin/crates/resolver/src/lib.rs b/bobbin/crates/resolver/src/lib.rs index 3d1674a0..eb8cc8cd 100644 --- a/bobbin/crates/resolver/src/lib.rs +++ b/bobbin/crates/resolver/src/lib.rs @@ -72,6 +72,22 @@ impl CacheEntry { } } +// repos are addressed by name in urls, so the name has to map back to an ident +#[derive(Clone, Debug, Eq, Hash, PartialEq)] +struct RepoNameKey { + owner: Did, + name: String, +} + +impl RepoNameKey { + fn new(owner: Did, name: &str) -> Self { + Self { + owner, + name: name.to_owned(), + } + } +} + #[derive(Default)] pub struct ResolverStats { hits: AtomicU64, @@ -169,6 +185,8 @@ struct SlingshotProbe { pub struct RepoIdResolver { cache: SccMap, by_repo_did: SccMap, RepoIdent, RuntimeHasher>, + by_name: SccMap, + name_of: SccMap, in_flight: SccMap>, RuntimeHasher>, probe: Option, stats: ResolverStats, @@ -183,6 +201,8 @@ impl RepoIdResolver { Self { cache: SccMap::with_hasher(hasher.clone()), by_repo_did: SccMap::with_hasher(hasher.clone()), + by_name: SccMap::with_hasher(hasher.clone()), + name_of: SccMap::with_hasher(hasher.clone()), in_flight: SccMap::with_hasher(hasher), probe: Some(SlingshotProbe { client, clock }), stats: ResolverStats::default(), @@ -193,6 +213,8 @@ impl RepoIdResolver { Self { cache: SccMap::with_hasher(hasher.clone()), by_repo_did: SccMap::with_hasher(hasher.clone()), + by_name: SccMap::with_hasher(hasher.clone()), + name_of: SccMap::with_hasher(hasher.clone()), in_flight: SccMap::with_hasher(hasher), probe: None, stats: ResolverStats::default(), @@ -226,6 +248,43 @@ impl RepoIdResolver { .map(|e| e.get().clone()) } + pub async fn lookup_by_name(&self, owner: &Did, name: &str) -> Option { + let key = RepoNameKey::new(owner.clone(), name); + self.by_name.get_async(&key).await.map(|e| e.get().clone()) + } + + pub async fn observe_name(&self, owner: Did, rkey: Rkey, name: &str) { + let ident = RepoIdent::new(owner.clone(), rkey); + let key = RepoNameKey::new(owner, name); + + // a rename leaves the old name pointing here, so give it up first + if let Some(prior) = self + .name_of + .get_async(&ident) + .await + .map(|e| e.get().clone()) + && prior != key + { + self.by_name + .remove_if_async(&prior, |existing| *existing == ident) + .await; + } + + self.name_of + .entry_async(ident.clone()) + .await + .and_modify(|existing| *existing = key.clone()) + .or_insert(key.clone()); + + // two repos claiming one name is the owner's problem + // we just resolve to the last one + self.by_name + .entry_async(key) + .await + .and_modify(|existing| *existing = ident.clone()) + .or_insert(ident); + } + pub async fn observe( &self, owner: Did, @@ -268,6 +327,11 @@ impl RepoIdResolver { .remove_if_async(&repo_did, |existing| *existing == ident) .await; } + if let Some((_, name)) = self.name_of.remove_async(&ident).await { + self.by_name + .remove_if_async(&name, |existing| *existing == ident) + .await; + } } async fn fill_provisional(&self, key: RepoIdent, resolution: Resolution) { @@ -514,6 +578,100 @@ mod tests { ); } + #[tokio::test] + async fn lookup_by_name_finds_observed_ident() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + resolver + .observe_name(did("did:plc:nel"), rkey("3liuighjy2h22"), "core") + .await; + let got = resolver.lookup_by_name(&did("did:plc:nel"), "core").await; + assert_eq!( + got, + Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), + ); + } + + #[tokio::test] + async fn lookup_by_name_is_scoped_to_the_owner() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + resolver + .observe_name(did("did:plc:nel"), rkey("3liuighjy2h22"), "core") + .await; + let got = resolver + .lookup_by_name(&did("did:plc:olaren"), "core") + .await; + assert_eq!(got, None, "one owner's name must not answer for another's"); + } + + #[tokio::test] + async fn rename_drops_the_old_name() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let owner = did("did:plc:nel"); + resolver + .observe_name(owner.clone(), rkey("3liuighjy2h22"), "core") + .await; + resolver + .observe_name(owner.clone(), rkey("3liuighjy2h22"), "tangled") + .await; + + assert_eq!( + resolver.lookup_by_name(&owner, "core").await, + None, + "the old name must stop resolving or a renamed repo answers on two urls", + ); + assert_eq!( + resolver.lookup_by_name(&owner, "tangled").await, + Some(RepoIdent::new(owner, rkey("3liuighjy2h22"))), + ); + } + + #[tokio::test] + async fn a_taken_name_goes_to_the_last_observed_repo() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let owner = did("did:plc:nel"); + resolver + .observe_name(owner.clone(), rkey("3liuighjy2h22"), "core") + .await; + resolver + .observe_name(owner.clone(), rkey("xyzxyzxyzxyzx"), "core") + .await; + assert_eq!( + resolver.lookup_by_name(&owner, "core").await, + Some(RepoIdent::new(owner, rkey("xyzxyzxyzxyzx"))), + ); + } + + #[tokio::test] + async fn forget_clears_the_name() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let owner = did("did:plc:nel"); + resolver + .observe_name(owner.clone(), rkey("3liuighjy2h22"), "core") + .await; + resolver.forget(&owner, &rkey("3liuighjy2h22")).await; + assert_eq!(resolver.lookup_by_name(&owner, "core").await, None); + } + + #[tokio::test] + async fn forget_leaves_a_name_that_moved_on() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let owner = did("did:plc:nel"); + resolver + .observe_name(owner.clone(), rkey("3liuighjy2h22"), "core") + .await; + resolver + .observe_name(owner.clone(), rkey("xyzxyzxyzxyzx"), "core") + .await; + + resolver.forget(&owner, &rkey("3liuighjy2h22")).await; + + assert_eq!( + resolver.lookup_by_name(&owner, "core").await, + Some(RepoIdent::new(owner, rkey("xyzxyzxyzxyzx"))), + "deleting the repo that lost the name must not unhook the one that has it", + ); + } + #[tokio::test] async fn observation_with_repo_did_resolves_mapped() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index 6a70a2f0..338dff3c 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -181,6 +181,7 @@ pub fn router(state: AppState) -> Router { "/xrpc/sh.tangled.repo.getReposByRepoDids", get(get_repos_by_repo_dids), ) + .route("/xrpc/sh.tangled.repo.getRepoByName", get(get_repo_by_name)) .route("/xrpc/sh.tangled.actor.getProfile", get(get_profile)) .route("/xrpc/sh.tangled.actor.getProfiles", get(get_profiles)) .route("/xrpc/sh.tangled.repo.getIssue", get(get_issue)) @@ -607,6 +608,12 @@ struct GetRepoByRepoDidQuery { repo_did: Did, } +#[derive(Debug, Deserialize)] +struct GetRepoByNameQuery { + owner: Did, + name: String, +} + #[derive(Debug, Deserialize)] struct GetProfileQuery { actor: AtUri, @@ -1291,6 +1298,29 @@ async fn get_repo_by_repo_did( }))) } +async fn get_repo_by_name( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + let ident = state + .resolver + .lookup_by_name(&q.owner, &q.name) + .await + .ok_or(XrpcError::NotFound)?; + let uri = AtUri::::from_parts_owned( + ident.owner.as_str(), + RepoRecord::NSID, + ident.rkey.as_str(), + ) + .expect("Did and Rkey newtypes already validated, at-uri assembly cannot fail"); + let (body, value) = fetch_from_uri::>(&state, uri).await?; + Ok(Json(Deduped(RepoGetRecordOutput { + cid: Some(body.cid.clone()), + uri: body.uri.clone(), + value, + }))) +} + async fn get_profile( State(state): State, XrpcQuery(q): XrpcQuery, diff --git a/bobbin/crates/xrpc/tests/cold_start.rs b/bobbin/crates/xrpc/tests/cold_start.rs index 30bc3bea..9de747eb 100644 --- a/bobbin/crates/xrpc/tests/cold_start.rs +++ b/bobbin/crates/xrpc/tests/cold_start.rs @@ -90,6 +90,13 @@ fn xrpc_request(endpoint: &str, param: &str, value: &str) -> Request { .unwrap() } +fn xrpc_request2(endpoint: &str, a: (&str, &str), b: (&str, &str)) -> Request { + Request::builder() + .uri(format!("/xrpc/{endpoint}?{}={}&{}={}", a.0, a.1, b.0, b.1)) + .body(Body::empty()) + .unwrap() +} + fn xrpc_request_escaped(endpoint: &str, param: &str, value: &str) -> Request { let encoded: String = byte_serialize(value.as_bytes()).collect(); Request::builder() @@ -666,6 +673,105 @@ async fn get_repo_by_repo_did_returns_observed_record() { assert_eq!(body["value"]["repoDid"], repo_did.as_ref()); } +#[tokio::test] +async fn get_repo_by_name_returns_observed_record() { + let server = MockServer::start().await; + let owner_did = did("did:plc:scallop"); + let rk = rkey("3liuighjy2h22"); + mount_record( + &server, + &owner_did, + &nsid("sh.tangled.repo"), + &rk, + json!({ + "$type": "sh.tangled.repo", + "name": "core", + "knot": "oyster.cafe", + "createdAt": "2026-05-01T00:00:00Z", + }), + ) + .await; + + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; + state + .resolver + .observe_name(owner_did.clone(), rk.clone(), "core") + .await; + + let app = router(state); + let resp = app + .oneshot(xrpc_request2( + "sh.tangled.repo.getRepoByName", + ("owner", owner_did.as_ref()), + ("name", "core"), + )) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + body["uri"], + format!( + "at://{}/sh.tangled.repo/{}", + owner_did.as_ref(), + rk.as_ref() + ) + ); + assert_eq!(body["value"]["name"], "core"); +} + +#[tokio::test] +async fn get_repo_by_name_404_when_unobserved() { + let server = MockServer::start().await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; + let app = router(state); + let resp = app + .oneshot(xrpc_request2( + "sh.tangled.repo.getRepoByName", + ("owner", "did:plc:scallop"), + ("name", "core"), + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::NOT_FOUND); +} + +#[tokio::test] +async fn get_repo_by_name_404_for_another_owners_name() { + let server = MockServer::start().await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; + state + .resolver + .observe_name(did("did:plc:scallop"), rkey("3liuighjy2h22"), "core") + .await; + let app = router(state); + let resp = app + .oneshot(xrpc_request2( + "sh.tangled.repo.getRepoByName", + ("owner", "did:plc:whelk"), + ("name", "core"), + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::NOT_FOUND); +} + +#[tokio::test] +async fn get_repo_by_name_400_on_invalid_owner() { + let server = MockServer::start().await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; + let app = router(state); + let resp = app + .oneshot(xrpc_request2( + "sh.tangled.repo.getRepoByName", + ("owner", "not-a-did"), + ("name", "core"), + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::BAD_REQUEST); +} + #[tokio::test] async fn get_repo_by_repo_did_404_when_unobserved() { let server = MockServer::start().await; diff --git a/lexicons/repo/getRepoByName.json b/lexicons/repo/getRepoByName.json new file mode 100644 index 00000000..44a90d3e --- /dev/null +++ b/lexicons/repo/getRepoByName.json @@ -0,0 +1,45 @@ +{ + "lexicon": 1, + "id": "sh.tangled.repo.getRepoByName", + "defs": { + "main": { + "type": "query", + "parameters": { + "type": "params", + "required": ["owner", "name"], + "properties": { + "owner": { + "type": "string", + "format": "did", + "description": "DID of the account that owns the repo." + }, + "name": { + "type": "string", + "description": "Name of the repo as it appears in its url." + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["uri", "value"], + "properties": { + "uri": { + "type": "string", + "format": "at-uri" + }, + "cid": { + "type": "string", + "format": "cid" + }, + "value": { + "type": "unknown", + "description": "Embedded sh.tangled.repo record." + } + } + } + } + } + } +} -- 2.51.2