From 5c53ee7bc9a86ea5a2d0469b1dc6bc4ec5564eee Mon Sep 17 00:00:00 2001 From: dawn Date: Mon, 5 Oct 2026 21:01:09 +0300 Subject: [PATCH] lexicons,bobbin: put a repo's knot on its repo views Signed-off-by: dawn --- bobbin/crates/resolver/src/hydrant.rs | 102 +++++++++++++++++- bobbin/crates/resolver/src/identity.rs | 81 +++++++++++++- bobbin/crates/xrpc/src/lib.rs | 39 +++++-- lexicons/repo/getRepo.json | 4 + lexicons/repo/getRepoByName.json | 4 + lexicons/repo/getRepoByRepoDid.json | 4 + .../lexicons/types/sh/tangled/repo/getRepo.ts | 4 + .../types/sh/tangled/repo/getRepoByName.ts | 4 + .../types/sh/tangled/repo/getRepoByRepoDid.ts | 4 + 9 files changed, 235 insertions(+), 11 deletions(-) diff --git a/bobbin/crates/resolver/src/hydrant.rs b/bobbin/crates/resolver/src/hydrant.rs index ac301b2cf..a601b64aa 100644 --- a/bobbin/crates/resolver/src/hydrant.rs +++ b/bobbin/crates/resolver/src/hydrant.rs @@ -38,6 +38,34 @@ pub enum RepoIdentity { Dead, } +#[derive(Debug, Deserialize)] +struct DescribedRepo { + #[serde(rename = "didDoc")] + did_doc: DescribedDoc, +} + +#[derive(Debug, Deserialize)] +struct DescribedDoc { + #[serde(default)] + service: Vec, +} + +#[derive(Debug, Deserialize)] +struct DescribedService { + id: String, + #[serde(rename = "serviceEndpoint")] + endpoint: serde_json::Value, +} + +fn knot_host(did: &Did, services: &[DescribedService]) -> Option { + let service = services.iter().find(|service| { + service.id == "#tangled_knot" + || service.id.strip_prefix(did.as_str()) == Some("#tangled_knot") + })?; + let url = Url::parse(service.endpoint.as_str()?).ok()?; + Some(url[url::Position::BeforeHost..url::Position::AfterPort].to_owned()) +} + #[derive(Debug, Deserialize)] struct RepoInfo { #[serde(default)] @@ -158,6 +186,26 @@ impl HydrantClient { Ok(serde_json::from_slice::(&bytes)?.id) } + /// the host in a did's `#tangled_knot` service, `None` when its doc names no knot + pub async fn repo_knot(&self, did: &Did) -> Result, HydrantError> { + // the only hydrant endpoint with the raw doc, the rest are minidocs without services + let mut url = self.endpoint(["xrpc", "com.atproto.repo.describeRepo"])?; + url.query_pairs_mut().append_pair("repo", did.as_str()); + let resp = self + .http + .execute(HttpRequest { + url, + headers: HeaderMap::new(), + }) + .await?; + if resp.status != StatusCode::OK { + return Err(HydrantError::Upstream(resp.status)); + } + let bytes = read_bounded(resp).await?; + let described = serde_json::from_slice::(&bytes)?; + Ok(knot_host(did, &described.did_doc.service)) + } + /// `None` when hydrant does not track the repo or has no handle for it pub async fn repo_identity( &self, @@ -222,7 +270,7 @@ mod tests { use super::*; use bobbin_runtime::ReqwestHttp; use bobbin_slingshot_client::default_http_client; - use wiremock::matchers::{method, path}; + use wiremock::matchers::{method, path, query_param}; use wiremock::{Mock, MockServer, ResponseTemplate}; fn did(value: &str) -> Did { @@ -533,4 +581,56 @@ mod tests { "got {err:?}" ); } + + fn described(services: serde_json::Value) -> ResponseTemplate { + ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:kelp", + "handle": "handle.invalid", + "handleIsCorrect": false, + "collections": [], + "didDoc": { "id": "did:plc:kelp", "service": services }, + })) + } + + #[tokio::test] + async fn repo_knot_takes_the_host_of_the_tangled_knot_service() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/com.atproto.repo.describeRepo")) + .and(query_param("repo", "did:plc:kelp")) + .respond_with(described(serde_json::json!([ + { + "id": "#atproto_pds", + "type": "AtprotoPersonalDataServer", + "serviceEndpoint": "https://knot.example.com", + }, + { + "id": "did:plc:kelp#tangled_knot", + "type": "TangledKnot", + "serviceEndpoint": "https://knot.example.com:8443/repo/abc", + }, + ]))) + .mount(&server) + .await; + + let knot = client(&server).await.repo_knot(&did("did:plc:kelp")).await; + assert_eq!(knot.unwrap().as_deref(), Some("knot.example.com:8443")); + } + + #[tokio::test] + async fn repo_knot_is_none_for_a_doc_with_only_a_pds() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/com.atproto.repo.describeRepo")) + .respond_with(described(serde_json::json!([{ + "id": "#atproto_pds", + "type": "AtprotoPersonalDataServer", + "serviceEndpoint": "https://pds.example.com", + }]))) + .mount(&server) + .await; + + let knot = client(&server).await.repo_knot(&did("did:plc:kelp")).await; + assert_eq!(knot.unwrap(), None); + } } diff --git a/bobbin/crates/resolver/src/identity.rs b/bobbin/crates/resolver/src/identity.rs index c4b53ebc9..738dad33f 100644 --- a/bobbin/crates/resolver/src/identity.rs +++ b/bobbin/crates/resolver/src/identity.rs @@ -27,6 +27,9 @@ use crate::hydrant::{HydrantClient, RepoIdentity}; const WARM_CONCURRENCY: usize = 32; // how long a failed lookup suppresses another attempt for that did const WARM_RETRY_AFTER: Duration = Duration::from_secs(300); +// identity frames are best effort and can miss a did doc change, so cached knots +// expire on their own too +const KNOT_TTL: Duration = Duration::from_secs(3600); #[derive(Clone, Copy)] enum WarmKind { @@ -172,6 +175,12 @@ struct IdentityResolverStats { warm_failed: AtomicU64, } +#[derive(Clone)] +struct KnotEntry { + knot: Option, + expires: Option, +} + // in the future this would spill to disk probably pub struct IdentityResolver { by_did: SccMap, IdentityState, RuntimeHasher>, @@ -181,6 +190,7 @@ pub struct IdentityResolver { queued: SccSet, RuntimeHasher>, // dids whose last warm failed, mapped to when we may try again cooldown: SccMap, Instant, RuntimeHasher>, + knots: SccMap, KnotEntry, RuntimeHasher>, warm_tx: mpsc::UnboundedSender, warm_rx: Mutex>>, hydrant: Option, @@ -222,7 +232,8 @@ impl IdentityResolver { by_handle: SccMap::with_hasher(hasher.clone()), in_flight: SccMap::with_hasher(hasher.clone()), queued: SccSet::with_hasher(hasher.clone()), - cooldown: SccMap::with_hasher(hasher), + cooldown: SccMap::with_hasher(hasher.clone()), + knots: SccMap::with_hasher(hasher), warm_tx, warm_rx: Mutex::new(Some(warm_rx)), hydrant: None, @@ -259,9 +270,40 @@ impl IdentityResolver { } pub fn refresh(&self, did: &Did) { + self.knots.remove_sync(did); self.enqueue(did, WarmKind::Refresh); } + /// the knot host a repo did names, `None` when it names none or hydrant can't say + pub async fn repo_knot(&self, did: &Did) -> Option { + let now = self.clock.as_ref().map(|clock| clock.now_instant()); + let cached = self.knots.read_sync(did, |_, entry| entry.clone()); + if let Some(entry) = cached + && entry + .expires + .zip(now) + .is_none_or(|(expires, now)| now < expires) + { + return entry.knot; + } + let knot = self + .hydrant + .as_ref()? + .repo_knot(did) + .await + .inspect_err(|error| debug!(did = did.as_str(), %error, "knot lookup failed")) + .ok()?; + let expires = now.map(|now| now + KNOT_TTL); + self.knots.upsert_sync( + did.clone(), + KnotEntry { + knot: knot.clone(), + expires, + }, + ); + knot + } + fn enqueue(&self, did: &Did, kind: WarmKind) { if self.cooling_down(did) { return; @@ -379,6 +421,7 @@ impl IdentityResolver { } pub fn observe(&self, did: Did, handle: Handle) { + self.knots.remove_sync(&did); let observed = MiniDoc { did: did.clone(), handle, @@ -438,6 +481,7 @@ impl IdentityResolver { } pub fn deactivate(&self, did: Did) { + self.knots.remove_sync(&did); let _mutation = self .mutations .lock() @@ -1109,6 +1153,41 @@ mod tests { assert_eq!(resolver.stats().upstream_requests, 0); } + #[tokio::test] + async fn repo_knot_is_cached_until_an_identity_change() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/com.atproto.repo.describeRepo")) + .and(query_param("repo", "did:plc:kelp")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "didDoc": { + "id": "did:plc:kelp", + "service": [{ + "id": "#tangled_knot", + "type": "TangledKnot", + "serviceEndpoint": "https://knot.example.com/repo/abc", + }], + }, + }))) + .expect(2) + .mount(&server) + .await; + + let resolver = hydrant_backed(&server); + let repo = did("did:plc:kelp"); + for _ in 0..2 { + assert_eq!( + resolver.repo_knot(&repo).await.as_deref(), + Some("knot.example.com") + ); + } + resolver.refresh(&repo); + assert_eq!( + resolver.repo_knot(&repo).await.as_deref(), + Some("knot.example.com") + ); + } + #[tokio::test] async fn warming_falls_back_to_slingshot_for_dids_hydrant_lacks() { let server = MockServer::start().await; diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index 970a74970..2961203f1 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -1437,21 +1437,21 @@ where async fn get_repo( State(state): State, XrpcQuery(q): XrpcQuery, -) -> Result>>, XrpcError> { +) -> Result>, XrpcError> { get_manifest(&state, q.repo).await } async fn get_repo_by_repo_did( State(state): State, XrpcQuery(q): XrpcQuery, -) -> Result>>, XrpcError> { +) -> Result>, XrpcError> { get_manifest(&state, view::manifest_uri(&q.repo_did)).await } async fn get_repo_by_name( State(state): State, XrpcQuery(q): XrpcQuery, -) -> Result>>, XrpcError> { +) -> Result>, XrpcError> { let repo_did = state .claims .repo_by_name(q.owner.as_str(), &q.name) @@ -1460,15 +1460,36 @@ async fn get_repo_by_name( get_manifest(&state, view::manifest_uri(&repo_did)).await } +#[derive(Serialize)] +struct RepoOutput { + #[serde(flatten)] + record: ManifestGetRecordOutput, + #[serde(skip_serializing_if = "Option::is_none")] + knot: Option, +} + async fn get_manifest( state: &AppState, uri: AtUri, -) -> Result>>, XrpcError> { - let (body, value) = fetch_from_uri::>(state, uri).await?; - Ok(Json(Deduped(ManifestGetRecordOutput { - cid: Some(body.cid.clone()), - uri: body.uri.clone(), - value, +) -> Result>, XrpcError> { + let repo_did = owner_did_from_aturi(&uri); + let (fetched, knot) = tokio::join!( + fetch_from_uri::>(state, uri), + async { + match repo_did { + Some(did) => state.identity.repo_knot(&did).await, + None => None, + } + } + ); + let (body, value) = fetched?; + Ok(Json(Deduped(RepoOutput { + record: ManifestGetRecordOutput { + cid: Some(body.cid.clone()), + uri: body.uri.clone(), + value, + }, + knot, }))) } diff --git a/lexicons/repo/getRepo.json b/lexicons/repo/getRepo.json index bf715ab01..64af9e580 100644 --- a/lexicons/repo/getRepo.json +++ b/lexicons/repo/getRepo.json @@ -32,6 +32,10 @@ "value": { "type": "unknown", "description": "Embedded sh.tangled.repo record." + }, + "knot": { + "type": "string", + "description": "Host of the knot named by the #tangled_knot service in the repo's DID document. Absent when it names none." } } } diff --git a/lexicons/repo/getRepoByName.json b/lexicons/repo/getRepoByName.json index 44a90d3ea..35df41fe4 100644 --- a/lexicons/repo/getRepoByName.json +++ b/lexicons/repo/getRepoByName.json @@ -36,6 +36,10 @@ "value": { "type": "unknown", "description": "Embedded sh.tangled.repo record." + }, + "knot": { + "type": "string", + "description": "Host of the knot named by the #tangled_knot service in the repo's DID document. Absent when it names none." } } } diff --git a/lexicons/repo/getRepoByRepoDid.json b/lexicons/repo/getRepoByRepoDid.json index a78387043..d1d6b2550 100644 --- a/lexicons/repo/getRepoByRepoDid.json +++ b/lexicons/repo/getRepoByRepoDid.json @@ -32,6 +32,10 @@ "value": { "type": "unknown", "description": "Embedded sh.tangled.repo record." + }, + "knot": { + "type": "string", + "description": "Host of the knot named by the #tangled_knot service in the repo's DID document. Absent when it names none." } } } diff --git a/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepo.ts b/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepo.ts index ae2ba692c..3bd50fbd6 100644 --- a/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepo.ts +++ b/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepo.ts @@ -13,6 +13,10 @@ const _mainSchema = /*#__PURE__*/ v.query("sh.tangled.repo.getRepo", { type: "lex", schema: /*#__PURE__*/ v.object({ cid: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.cidString()), + /** + * Host of the knot named by the #tangled_knot service in the repo's DID document. Absent when it names none. + */ + knot: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.string()), uri: /*#__PURE__*/ v.resourceUriString(), /** * Embedded sh.tangled.repo record. diff --git a/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepoByName.ts b/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepoByName.ts index f8c505758..9972520a0 100644 --- a/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepoByName.ts +++ b/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepoByName.ts @@ -17,6 +17,10 @@ const _mainSchema = /*#__PURE__*/ v.query("sh.tangled.repo.getRepoByName", { type: "lex", schema: /*#__PURE__*/ v.object({ cid: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.cidString()), + /** + * Host of the knot named by the #tangled_knot service in the repo's DID document. Absent when it names none. + */ + knot: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.string()), uri: /*#__PURE__*/ v.resourceUriString(), /** * Embedded sh.tangled.repo record. diff --git a/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepoByRepoDid.ts b/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepoByRepoDid.ts index f9e4f3aff..d4ecf12d9 100644 --- a/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepoByRepoDid.ts +++ b/web/src/lib/api/lexicons/types/sh/tangled/repo/getRepoByRepoDid.ts @@ -13,6 +13,10 @@ const _mainSchema = /*#__PURE__*/ v.query("sh.tangled.repo.getRepoByRepoDid", { type: "lex", schema: /*#__PURE__*/ v.object({ cid: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.cidString()), + /** + * Host of the knot named by the #tangled_knot service in the repo's DID document. Absent when it names none. + */ + knot: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.string()), uri: /*#__PURE__*/ v.resourceUriString(), /** * Embedded sh.tangled.repo record. -- 2.51.2