From da359d65e90564fbc1d9a4a6ebaed2b335a49476 Mon Sep 17 00:00:00 2001 From: dawn Date: Fri, 7 Aug 2026 00:42:53 +0300 Subject: [PATCH] bobbin/crates/{ingest,resolver,xrpc}: warm id using hydrant /repos, fallback to slingshot, add pds and signing key Signed-off-by: dawn --- Cargo.lock | 4 + bobbin/crates/bobbin/src/main.rs | 14 +- bobbin/crates/bobbin/src/mem/report.rs | 8 + bobbin/crates/ingest/src/frame.rs | 13 +- bobbin/crates/ingest/src/lib.rs | 135 +++-- bobbin/crates/resolver/Cargo.toml | 6 +- bobbin/crates/resolver/src/hydrant.rs | 459 +++++++++++++++ bobbin/crates/resolver/src/identity.rs | 768 ++++++++++++++++++++----- bobbin/crates/resolver/src/lib.rs | 8 +- bobbin/crates/xrpc/src/lib.rs | 1 + 10 files changed, 1222 insertions(+), 194 deletions(-) create mode 100644 bobbin/crates/resolver/src/hydrant.rs diff --git a/Cargo.lock b/Cargo.lock index f2b3bdf44..d55447c4f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -811,12 +811,16 @@ dependencies = [ "bobbin-runtime", "bobbin-slingshot-client", "bobbin-types", + "bytes", + "futures", + "http", "jacquard-common", "scc", "serde", "serde_json", "thiserror 2.0.18", "tokio", + "tokio-util", "tracing", "url", "wiremock", diff --git a/bobbin/crates/bobbin/src/main.rs b/bobbin/crates/bobbin/src/main.rs index 9bad0238c..f56a1792f 100644 --- a/bobbin/crates/bobbin/src/main.rs +++ b/bobbin/crates/bobbin/src/main.rs @@ -13,7 +13,7 @@ use bobbin_ingest::{ use bobbin_knot_ingest::{CapabilityGate, KnotClient, KnotRegistry, Orchestrator}; use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig, MirrorProxy, classify_ip}; use bobbin_record_lru::{CacheCapacity, LruRecordStore, RecordStore}; -use bobbin_resolver::IdentityResolver; +use bobbin_resolver::{HydrantClient, IdentityResolver}; use bobbin_runtime::{ Clock, GuardedWs, MemoryBudget, NetworkError, OsEntropy, ReqwestHttp, RuntimeHasher, SystemClock, TungsteniteWs, WsTransport, @@ -187,10 +187,13 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { let records: Arc = Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(lru_cap))); let slingshot = SlingshotClient::with_default_http(cfg.slingshot.url.clone())?; - let identity = Arc::new(IdentityResolver::with_slingshot( - slingshot.clone(), - hasher.clone(), - )); + let identity = Arc::new( + IdentityResolver::with_slingshot(slingshot.clone(), clock.clone(), hasher.clone()) + .with_hydrant(HydrantClient::new( + cfg.hydrant.url.clone(), + ReqwestHttp::shared(default_http_client()?), + )?), + ); let mut resolver_opts = ResolverOptions::default(); // NOTE: see https://tangled.org/nonbinary.computer/jacquard/issues/39. resolver_opts.did_order = vec![ @@ -304,6 +307,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { knot_gate: Some(knot_gate.clone()), }; let mut ingest_handle = tokio::spawn(run_ingest(ingest_cfg, ingest_runtime)); + let _identity_warmer = tokio::spawn(identity.clone().run_warming(cancel.clone())); let knot_orchestrator = Orchestrator { client: Arc::new(knot_client), diff --git a/bobbin/crates/bobbin/src/mem/report.rs b/bobbin/crates/bobbin/src/mem/report.rs index 5428e3938..3814372c0 100644 --- a/bobbin/crates/bobbin/src/mem/report.rs +++ b/bobbin/crates/bobbin/src/mem/report.rs @@ -90,6 +90,10 @@ struct Identity { hits: u64, misses: u64, upstream_requests: u64, + warm_queued: u64, + warm_dropped: u64, + warm_resolved: u64, + warm_failed: u64, } #[derive(Serialize)] @@ -180,6 +184,10 @@ async fn mem_report(State(p): State) -> Json { hits: identity.hits, misses: identity.misses, upstream_requests: identity.upstream_requests, + warm_queued: identity.warm_queued, + warm_dropped: identity.warm_dropped, + warm_resolved: identity.warm_resolved, + warm_failed: identity.warm_failed, }, derived: Derived { known_bytes, diff --git a/bobbin/crates/ingest/src/frame.rs b/bobbin/crates/ingest/src/frame.rs index b96d7f3bf..711092d5c 100644 --- a/bobbin/crates/ingest/src/frame.rs +++ b/bobbin/crates/ingest/src/frame.rs @@ -16,13 +16,22 @@ pub struct HydrantFrame { pub record: Option, #[serde(default)] pub identity: Option, + #[serde(default)] + pub account: Option, } #[derive(Clone, Debug, Deserialize)] pub struct IdentityFrame { pub did: Did, - pub handle: Handle, - pub is_active: bool, + // omitted when the did has no currently resolving handle + #[serde(default)] + pub handle: Option>, +} + +#[derive(Clone, Debug, Deserialize)] +pub struct AccountFrame { + pub did: Did, + pub active: bool, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index b2eb3aac9..7bd015aa2 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -42,7 +42,7 @@ mod resolver; mod shadow; mod warming; use frame::HydrantStreamErrorFrame; -pub use frame::{FrameKind, HydrantFrame, IdentityFrame, RecordAction, RecordFrame}; +pub use frame::{AccountFrame, FrameKind, HydrantFrame, IdentityFrame, RecordAction, RecordFrame}; pub use resolver::{RepoIdResolver, Resolution}; pub use shadow::{WarmingShadowBuffer, WarmingShadowSnapshot}; pub use warming::{ParkedUpsert, WarmingBuffer, WarmingBufferSnapshot}; @@ -837,6 +837,7 @@ struct Pending { enum PendingOp { Noop, Identity(Option), + Account(Option), ClearCache { source: AtUri, }, @@ -869,7 +870,10 @@ fn pending_nsid(op: &PendingOp) -> Option<&Nsid> { PendingOp::Upsert(pieces) => Some(&pieces.nsid), PendingOp::Delete { nsid, .. } => Some(nsid), PendingOp::Parked { nsid, .. } => Some(nsid), - PendingOp::Noop | PendingOp::Identity(_) | PendingOp::ClearCache { .. } => None, + PendingOp::Noop + | PendingOp::Identity(_) + | PendingOp::Account(_) + | PendingOp::ClearCache { .. } => None, } } @@ -932,7 +936,10 @@ async fn claim_pending( ctx.resolver.forget(&ident.owner, &ident.rkey).await; } } - PendingOp::Noop | PendingOp::Identity(_) | PendingOp::Parked { .. } => {} + PendingOp::Noop + | PendingOp::Identity(_) + | PendingOp::Account(_) + | PendingOp::Parked { .. } => {} } pending } @@ -1051,7 +1058,7 @@ async fn prepare_frame( let op = match frame.kind { FrameKind::Record => prepare_record(frame.record, ctx).await, FrameKind::Identity => PendingOp::Identity(frame.identity), - FrameKind::Account => PendingOp::Noop, + FrameKind::Account => PendingOp::Account(frame.account), FrameKind::Other => { debug!(id = frame.id, "ignoring unknown hydrant frame kind"); PendingOp::Noop @@ -1424,13 +1431,19 @@ async fn commit_pending( PendingOp::Noop | PendingOp::Parked { .. } => {} PendingOp::Identity(observed) => { if let Some(observed) = observed { - if observed.is_active { - identity.observe(observed.did, observed.handle); - } else { - identity.deactivate(observed.did); + match observed.handle { + Some(handle) => identity.observe(observed.did, handle), + None => identity.refresh(&observed.did), } } } + PendingOp::Account(account) => { + if let Some(account) = account + && !account.active + { + identity.deactivate(account.did); + } + } PendingOp::ClearCache { source } => records.remove(&source), PendingOp::Upsert(pieces) => { let UpsertPieces { @@ -2259,6 +2272,20 @@ mod tests { assert_eq!(frame.kind, FrameKind::Record); } + #[test] + fn classify_hydrant_identity_and_account_frames() { + for text in [ + r#"{"id":1,"type":"account","account":{"did":"did:web:guestbook.gaze.systems","active":true,"status":"desynchronized"}}"#, + r#"{"id":2,"type":"identity","identity":{"did":"did:web:guestbook.gaze.systems","handle":"guestbook.gaze.systems"}}"#, + r#"{"id":43,"type":"identity","identity":{"did":"did:plc:olaren"}}"#, + r#"{"id":44,"type":"account","account":{"did":"did:plc:olaren","active":false}}"#, + ] { + let frame = classify_text_frame(text) + .unwrap_or_else(|e| panic!("hydrant frame must parse: {text} -> {e:?}")); + assert!(frame.identity.is_some() || frame.account.is_some()); + } + } + #[test] fn classify_garbage_object_returns_decode_error() { let text = r#"{"random":"object","without":"required fields"}"#; @@ -2573,7 +2600,7 @@ mod tests { } #[tokio::test] - async fn identity_frames_update_identity_resolver_lifecycle() { + async fn identity_and_account_frames_drive_the_identity_resolver() { let (store, issue_states, pull_statuses, coverage, resolver) = fresh(); let identity = IdentityResolver::detached(RuntimeHasher::default()); let records = NoopRecordStore; @@ -2591,64 +2618,62 @@ mod tests { knot_registry: None, knot_gate: None, }; - let frame: HydrantFrame = parse_frame(json!({ + let olaren = Did::new_static("did:plc:olaren").unwrap(); + let apply = async |frame: HydrantFrame| { + let pending = prepare_frame(frame, &ctx, now()).await; + let pending = resolve_pending(pending, &ctx).await; + commit_pending( + pending, + &store, + &issue_states, + &pull_statuses, + &coverage, + &search, + &records, + &resolver, + &identity, + ) + .await; + }; + + apply(parse_frame(json!({ "id": 5, "type": "identity", - "identity": { - "did": "did:plc:olaren", - "handle": "olaren.dev", - "is_active": true, - "status": "active" - } - })); - let pending = prepare_frame(frame, &ctx, now()).await; - let pending = resolve_pending(pending, &ctx).await; - commit_pending( - pending, - &store, - &issue_states, - &pull_statuses, - &coverage, - &search, - &records, - &resolver, - &identity, - ) + "identity": {"did": "did:plc:olaren", "handle": "olaren.dev"} + }))) .await; - let doc = identity - .get_by_did(&Did::new_static("did:plc:olaren").unwrap()) + .get_by_did(&olaren) .expect("identity frame should seed the resolver"); assert_eq!(doc.handle.as_ref(), "olaren.dev"); - let frame: HydrantFrame = parse_frame(json!({ + apply(parse_frame(json!({ "id": 6, "type": "identity", - "identity": { - "did": "did:plc:olaren", - "handle": "olaren.dev", - "is_active": false, - "status": "deactivated" - } - })); - let pending = prepare_frame(frame, &ctx, now()).await; - let pending = resolve_pending(pending, &ctx).await; - commit_pending( - pending, - &store, - &issue_states, - &pull_statuses, - &coverage, - &search, - &records, - &resolver, - &identity, - ) + "identity": {"did": "did:plc:olaren"} + }))) .await; + let kept = identity + .get_by_did(&olaren) + .expect("a missing handle must not drop the cached one"); + assert_eq!(kept.handle.as_ref(), "olaren.dev"); + apply(parse_frame(json!({ + "id": 7, + "type": "identity", + "identity": {"did": "did:plc:olaren", "handle": "olaren.dev"} + }))) + .await; + apply(parse_frame(json!({ + "id": 8, + "type": "account", + "account": {"did": "did:plc:olaren", "active": false, "status": "deactivated"} + }))) + .await; assert_eq!( - identity.get_by_did(&Did::new_static("did:plc:olaren").unwrap()), - Err(bobbin_resolver::IdentityResolveError::NotFound) + identity.get_by_did(&olaren), + Err(bobbin_resolver::IdentityResolveError::NotFound), + "an inactive account frame is the only liveness signal hydrant sends", ); } diff --git a/bobbin/crates/resolver/Cargo.toml b/bobbin/crates/resolver/Cargo.toml index b59645f6e..60a18f288 100644 --- a/bobbin/crates/resolver/Cargo.toml +++ b/bobbin/crates/resolver/Cargo.toml @@ -11,15 +11,19 @@ bobbin-runtime = { workspace = true } bobbin-slingshot-client = { workspace = true } jacquard-common = { workspace = true } +bytes = { workspace = true } +futures = { workspace = true } +http = { workspace = true } +url = { workspace = true } scc = { workspace = true } serde_json = { workspace = true } serde = { workspace = true } thiserror = { workspace = true } tokio = { workspace = true } +tokio-util = { workspace = true } tracing = { workspace = true } [dev-dependencies] tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } wiremock = { workspace = true } -url = { workspace = true } serde_json = { workspace = true } diff --git a/bobbin/crates/resolver/src/hydrant.rs b/bobbin/crates/resolver/src/hydrant.rs new file mode 100644 index 000000000..c13773d22 --- /dev/null +++ b/bobbin/crates/resolver/src/hydrant.rs @@ -0,0 +1,459 @@ +use std::sync::Arc; + +use bobbin_runtime::{HttpRequest, HttpResponseHead, HttpTransport, NetworkError}; +use bytes::BytesMut; +use futures::TryStreamExt; +use http::{HeaderMap, StatusCode}; +use jacquard_common::DefaultStr; +use jacquard_common::types::crypto::PublicKey; +use jacquard_common::types::did::Did; +use jacquard_common::types::string::Handle; +use serde::Deserialize; +use std::borrow::Cow; +use thiserror::Error; +use tracing::debug; +use url::Url; + +use crate::identity::MiniDoc; + +const MAX_BODY_BYTES: u64 = 64 * 1024; + +#[derive(Debug, Error)] +pub enum HydrantError { + #[error("invalid base url scheme: {0}")] + BadScheme(String), + #[error("network: {0}")] + Network(#[from] NetworkError), + #[error("upstream returned status {0}")] + Upstream(StatusCode), + #[error("response body exceeded {limit} bytes")] + BodyTooLarge { limit: u64 }, + #[error("decode response: {0}")] + Decode(#[from] serde_json::Error), +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum RepoIdentity { + Live(MiniDoc), + Dead, +} + +#[derive(Debug, Deserialize)] +struct RepoInfo { + #[serde(default)] + handle: Option>, + #[serde(default)] + pds: Option, + #[serde(default, with = "crate::identity::signing_key")] + signing_key: Option>, + #[serde(default)] + status: Option, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +enum RepoStatus { + Synced, + Desynchronized, + Throttled, + Error, + Deactivated, + Takendown, + Suspended, + Deleted, + Unknown(String), +} + +impl RepoStatus { + fn is_live(&self) -> bool { + !matches!( + self, + Self::Deactivated | Self::Takendown | Self::Suspended | Self::Deleted + ) + } +} + +impl<'de> Deserialize<'de> for RepoStatus { + fn deserialize>(deserializer: D) -> Result { + let raw = Cow::<'de, str>::deserialize(deserializer)?; + Ok(match raw.as_ref() { + "synced" => Self::Synced, + "desynchronized" => Self::Desynchronized, + "throttled" => Self::Throttled, + "deactivated" => Self::Deactivated, + "takendown" => Self::Takendown, + "suspended" => Self::Suspended, + "deleted" => Self::Deleted, + // hydrant writes the message inline as error(..) + other if other.starts_with("error(") => Self::Error, + other => Self::Unknown(other.to_owned()), + }) + } +} + +#[derive(Clone)] +pub struct HydrantClient { + http: Arc, + base: Url, +} + +impl std::fmt::Debug for HydrantClient { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("HydrantClient") + .field("base", &self.base) + .finish_non_exhaustive() + } +} + +impl HydrantClient { + // takes the websocket schemes too, callers share the stream's base url + pub fn new(base: Url, http: Arc) -> Result { + let mut base = base; + let scheme = match base.scheme() { + "http" | "ws" => "http", + "https" | "wss" => "https", + other => return Err(HydrantError::BadScheme(other.to_owned())), + }; + base.set_scheme(scheme) + .map_err(|()| HydrantError::BadScheme(scheme.to_owned()))?; + Ok(Self { http, base }) + } + + // segment-wise so the colons in a did escape instead of parsing as a scheme + fn repo_url(&self, did: &Did) -> Result { + let mut url = self.base.clone(); + url.path_segments_mut() + .map_err(|()| HydrantError::BadScheme(self.base.scheme().to_owned()))? + .pop_if_empty() + .extend(["repos", did.as_str()]); + Ok(url) + } + + /// `None` when hydrant does not track the repo or has no handle for it + pub async fn repo_identity( + &self, + did: &Did, + ) -> Result, HydrantError> { + let url = self.repo_url(did)?; + + let resp = self + .http + .execute(HttpRequest { + url, + headers: HeaderMap::new(), + }) + .await?; + match resp.status { + StatusCode::OK => {} + StatusCode::NOT_FOUND => return Ok(None), + other => return Err(HydrantError::Upstream(other)), + } + + let bytes = read_bounded(resp).await?; + let info = serde_json::from_slice::(&bytes)?; + if let Some(RepoStatus::Unknown(status)) = info.status.as_ref() { + debug!(did = did.as_str(), status, "unknown hydrant repo status"); + } + if !info.status.as_ref().is_none_or(RepoStatus::is_live) { + return Ok(Some(RepoIdentity::Dead)); + } + Ok(info.handle.map(|handle| { + RepoIdentity::Live(MiniDoc { + did: did.clone(), + handle, + pds: info.pds, + signing_key: info.signing_key, + }) + })) + } +} + +async fn read_bounded(resp: HttpResponseHead) -> Result { + if resp.content_length.is_some_and(|len| len > MAX_BODY_BYTES) { + return Err(HydrantError::BodyTooLarge { + limit: MAX_BODY_BYTES, + }); + } + resp.body + .map_err(HydrantError::Network) + .try_fold(BytesMut::new(), |mut acc, chunk| async move { + if (acc.len() as u64).saturating_add(chunk.len() as u64) > MAX_BODY_BYTES { + return Err(HydrantError::BodyTooLarge { + limit: MAX_BODY_BYTES, + }); + } + acc.extend_from_slice(&chunk); + Ok(acc) + }) + .await +} + +#[cfg(test)] +mod tests { + use super::*; + use bobbin_runtime::ReqwestHttp; + use bobbin_slingshot_client::default_http_client; + use wiremock::matchers::{method, path}; + use wiremock::{Mock, MockServer, ResponseTemplate}; + + fn did(value: &str) -> Did { + Did::new_owned(value).unwrap() + } + + fn key(value: &str) -> PublicKey<'static> { + PublicKey::decode_owned(value).unwrap() + } + + async fn client(server: &MockServer) -> HydrantClient { + HydrantClient::new( + Url::parse(&server.uri()).unwrap(), + ReqwestHttp::shared(default_http_client().unwrap()), + ) + .unwrap() + } + + fn handle_of(found: Option) -> Option> { + match found { + Some(RepoIdentity::Live(doc)) => Some(doc.handle), + Some(RepoIdentity::Dead) | None => None, + } + } + + #[tokio::test] + async fn live_repo_yields_a_whole_minidoc() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:dawn")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:dawn", + "status": "synced", + "tracked": true, + "handle": "ptr.pet", + "pds": "https://pds.example.com/", + "signing_key": "zQ3shSiLsnqpyQ4SfDTT1D8qzFEoeYT8rSDXW6o8pVY7VcRBJ", + }))) + .mount(&server) + .await; + + let subject = did("did:plc:dawn"); + let found = client(&server).await.repo_identity(&subject).await; + assert_eq!( + found.unwrap(), + Some(RepoIdentity::Live(MiniDoc { + did: subject, + handle: Handle::new_owned("ptr.pet").unwrap(), + pds: Some(Url::parse("https://pds.example.com/").unwrap()), + signing_key: Some(key("zQ3shSiLsnqpyQ4SfDTT1D8qzFEoeYT8rSDXW6o8pVY7VcRBJ")), + })) + ); + } + + #[tokio::test] + async fn untracked_repo_is_a_miss_not_an_error() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:nobody")) + .respond_with(ResponseTemplate::new(404).set_body_string("repository not found")) + .mount(&server) + .await; + + let found = client(&server) + .await + .repo_identity(&did("did:plc:nobody")) + .await; + assert_eq!(found.unwrap(), None); + } + + #[tokio::test] + async fn tracked_repo_without_a_handle_is_a_miss() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:blank")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:blank", + "status": "backfilling", + "tracked": true, + }))) + .mount(&server) + .await; + + let found = client(&server) + .await + .repo_identity(&did("did:plc:blank")) + .await; + assert_eq!(found.unwrap(), None); + } + + #[tokio::test] + async fn dead_accounts_suppress_the_handle() { + for status in ["deactivated", "takendown", "suspended", "deleted"] { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:gone")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:gone", + "status": status, + "tracked": true, + "handle": "gone.example.com", + }))) + .mount(&server) + .await; + + let found = client(&server) + .await + .repo_identity(&did("did:plc:gone")) + .await; + assert_eq!(found.unwrap(), Some(RepoIdentity::Dead), "{status}"); + } + } + + #[tokio::test] + async fn degraded_statuses_keep_the_handle() { + for status in ["desynchronized", "throttled", "error(bad commit)"] { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:behind")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:behind", + "status": status, + "tracked": true, + "handle": "behind.example.com", + }))) + .mount(&server) + .await; + + let found = client(&server) + .await + .repo_identity(&did("did:plc:behind")) + .await; + assert_eq!( + handle_of(found.unwrap()), + Some(Handle::new_owned("behind.example.com").unwrap()), + "{status}" + ); + } + } + + #[tokio::test] + async fn an_unknown_status_still_keeps_the_handle() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:future")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:future", + "status": "quarantined", + "tracked": true, + "handle": "future.example.com", + }))) + .mount(&server) + .await; + + let found = client(&server) + .await + .repo_identity(&did("did:plc:future")) + .await; + assert_eq!( + handle_of(found.unwrap()), + Some(Handle::new_owned("future.example.com").unwrap()) + ); + } + + #[tokio::test] + async fn malformed_hints_are_a_decode_error() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:junk")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:junk", + "status": "synced", + "tracked": true, + "handle": "junk.example.com", + "pds": "not a url", + }))) + .mount(&server) + .await; + + let found = client(&server) + .await + .repo_identity(&did("did:plc:junk")) + .await; + assert!(matches!(found, Err(HydrantError::Decode(_))), "{found:?}"); + } + + #[tokio::test] + async fn real_hydrant_bytes_decode() { + let body = r#"{"did":"did:web:guestbook.gaze.systems","status":"desynchronized","tracked":true,"rev":"3mrpybqh4fs2p","data":"bafyreidbnwzu3biimfmzqtuub4j3uap7xvn3kbwsa4sgwbimclgebrxc4i","handle":"guestbook.gaze.systems","pds":"https://gaze.systems/","signing_key":"zQ3shSiLsnqpyQ4SfDTT1D8qzFEoeYT8rSDXW6o8pVY7VcRBJ","last_updated_at":"2026-08-07T10:24:13Z"}"#; + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:web:guestbook.gaze.systems")) + .respond_with(ResponseTemplate::new(200).set_body_raw(body, "application/json")) + .mount(&server) + .await; + + let found = client(&server) + .await + .repo_identity(&did("did:web:guestbook.gaze.systems")) + .await; + assert_eq!( + found.unwrap(), + Some(RepoIdentity::Live(MiniDoc { + did: did("did:web:guestbook.gaze.systems"), + handle: Handle::new_owned("guestbook.gaze.systems").unwrap(), + pds: Some(Url::parse("https://gaze.systems/").unwrap()), + signing_key: Some(key("zQ3shSiLsnqpyQ4SfDTT1D8qzFEoeYT8rSDXW6o8pVY7VcRBJ")), + })) + ); + } + + #[tokio::test] + async fn did_is_escaped_into_the_path() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:web:guestbook.gaze.systems")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:web:guestbook.gaze.systems", + "status": "synced", + "tracked": true, + "handle": "guestbook.gaze.systems", + }))) + .mount(&server) + .await; + + let found = client(&server) + .await + .repo_identity(&did("did:web:guestbook.gaze.systems")) + .await; + assert!(matches!(found.unwrap(), Some(RepoIdentity::Live(_)))); + } + + #[tokio::test] + async fn websocket_base_is_accepted() { + let client = HydrantClient::new( + Url::parse("ws://127.0.0.1:3099").unwrap(), + ReqwestHttp::shared(default_http_client().unwrap()), + ) + .unwrap(); + assert_eq!(client.base.scheme(), "http"); + } + + #[tokio::test] + async fn base_forms_all_address_the_same_repo() { + let subject = did("did:plc:dawn"); + let url = |base: &str| { + HydrantClient::new( + Url::parse(base).unwrap(), + ReqwestHttp::shared(default_http_client().unwrap()), + ) + .unwrap() + .repo_url(&subject) + .unwrap() + .to_string() + }; + + assert_eq!(url("http://h:3099"), "http://h:3099/repos/did:plc:dawn"); + assert_eq!(url("http://h:3099/"), "http://h:3099/repos/did:plc:dawn"); + assert_eq!( + url("wss://h/hydrant/"), + "https://h/hydrant/repos/did:plc:dawn" + ); + } +} diff --git a/bobbin/crates/resolver/src/identity.rs b/bobbin/crates/resolver/src/identity.rs index a86e67d22..3f35ccaea 100644 --- a/bobbin/crates/resolver/src/identity.rs +++ b/bobbin/crates/resolver/src/identity.rs @@ -1,17 +1,39 @@ -use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; -use bobbin_runtime::RuntimeHasher; +use bobbin_runtime::{Clock, RuntimeHasher}; use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; use jacquard_common::DefaultStr; +use jacquard_common::types::crypto::PublicKey; use jacquard_common::types::did::Did; use jacquard_common::types::ident::AtIdentifier; use jacquard_common::types::string::Handle; use scc::HashMap as SccMap; +use scc::HashSet as SccSet; use scc::hash_map::Entry as MapEntry; use serde::{Deserialize, Serialize}; use thiserror::Error; -use tokio::sync::OnceCell; +use tokio::sync::{OnceCell, Semaphore, mpsc}; +use tokio::time::Instant; +use tokio_util::sync::CancellationToken; +use tracing::debug; +use url::Url; + +use crate::SlingshotProbe; +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); + +#[derive(Clone, Copy)] +enum WarmKind { + Fill, + Refresh, +} + +type WarmItem = (Did, WarmKind); #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] #[serde(rename_all = "camelCase")] @@ -19,7 +41,45 @@ pub struct MiniDoc { pub did: Did, pub handle: Handle, #[serde(skip_serializing_if = "Option::is_none")] - pub pds: Option, + pub pds: Option, + #[serde(default, with = "signing_key", skip_serializing_if = "Option::is_none")] + pub signing_key: Option>, +} + +pub(crate) mod signing_key { + use jacquard_common::types::crypto::{KeyCodec, PublicKey, multikey}; + use serde::de::Error as _; + use serde::{Deserialize, Deserializer, Serializer}; + + pub fn serialize( + key: &Option>, + serializer: S, + ) -> Result { + match key { + Some(key) => { + // this is here because jacquard doesnt let us use theirs :( + let code = match key.codec { + KeyCodec::Ed25519 => 0xED, + KeyCodec::Secp256k1 => 0xE7, + KeyCodec::P256 => 0x1200, + KeyCodec::Unknown(code) => code, + }; + serializer.serialize_str(&multikey(code, &key.bytes)) + } + None => serializer.serialize_none(), + } + } + + pub fn deserialize<'de, D: Deserializer<'de>>( + deserializer: D, + ) -> Result>, D::Error> { + let Some(raw) = Option::>::deserialize(deserializer)? else { + return Ok(None); + }; + PublicKey::decode_owned(raw.as_ref()) + .map(Some) + .map_err(D::Error::custom) + } } #[derive(Clone, Debug, Error, Eq, PartialEq)] @@ -41,21 +101,16 @@ impl From for IdentityResolveError { } } -// hydrant doesnt send pds so we have a separate states to have -// public resolve call upgrade from a partial-cached to a full-cached #[derive(Clone)] enum IdentityState { - // seen by bobbin from hydrant - Observed(MiniDoc), - // fetched by bobbin from slingshot - Fetched(MiniDoc), + Cached(MiniDoc), Inactive, } impl IdentityState { fn doc(&self) -> Option<&MiniDoc> { match self { - Self::Observed(doc) | Self::Fetched(doc) => Some(doc), + Self::Cached(doc) => Some(doc), Self::Inactive => None, } } @@ -67,6 +122,10 @@ pub struct IdentityResolverStatsSnapshot { pub hits: u64, pub misses: u64, pub upstream_requests: u64, + pub warm_queued: u64, + pub warm_dropped: u64, + pub warm_resolved: u64, + pub warm_failed: u64, } #[derive(Default)] @@ -74,6 +133,10 @@ struct IdentityResolverStats { hits: AtomicU64, misses: AtomicU64, upstream_requests: AtomicU64, + warm_queued: AtomicU64, + warm_dropped: AtomicU64, + warm_resolved: AtomicU64, + warm_failed: AtomicU64, } // in the future this would spill to disk probably @@ -81,25 +144,50 @@ pub struct IdentityResolver { by_did: SccMap, IdentityState, RuntimeHasher>, by_handle: SccMap, Did, RuntimeHasher>, in_flight: SccMap>>, RuntimeHasher>, - slingshot: Option, + // dids queued for background warming, keeps the queue free of duplicates + queued: SccSet, RuntimeHasher>, + // dids whose last warm failed, mapped to when we may try again + cooldown: SccMap, Instant, RuntimeHasher>, + warm_tx: mpsc::UnboundedSender, + warm_rx: Mutex>>, + hydrant: Option, + probe: Option, + clock: Option>, stats: IdentityResolverStats, } impl IdentityResolver { - pub fn with_slingshot(slingshot: SlingshotClient, hasher: RuntimeHasher) -> Self { - Self::new(Some(slingshot), hasher) + pub fn with_slingshot( + client: SlingshotClient, + clock: Arc, + hasher: RuntimeHasher, + ) -> Self { + Self::new(Some(SlingshotProbe { client, clock }), hasher) } pub fn detached(hasher: RuntimeHasher) -> Self { Self::new(None, hasher) } - fn new(slingshot: Option, hasher: RuntimeHasher) -> Self { + pub fn with_hydrant(mut self, client: HydrantClient) -> Self { + self.hydrant = Some(client); + self + } + + fn new(probe: Option, hasher: RuntimeHasher) -> Self { + let clock = probe.as_ref().map(|probe| probe.clock.clone()); + let (warm_tx, warm_rx) = mpsc::unbounded_channel(); Self { by_did: SccMap::with_hasher(hasher.clone()), by_handle: SccMap::with_hasher(hasher.clone()), - in_flight: SccMap::with_hasher(hasher), - slingshot, + in_flight: SccMap::with_hasher(hasher.clone()), + queued: SccSet::with_hasher(hasher.clone()), + cooldown: SccMap::with_hasher(hasher), + warm_tx, + warm_rx: Mutex::new(Some(warm_rx)), + hydrant: None, + probe, + clock, stats: IdentityResolverStats::default(), } } @@ -110,41 +198,156 @@ impl IdentityResolver { hits: self.stats.hits.load(Ordering::Relaxed), misses: self.stats.misses.load(Ordering::Relaxed), upstream_requests: self.stats.upstream_requests.load(Ordering::Relaxed), + warm_queued: self.stats.warm_queued.load(Ordering::Relaxed), + warm_dropped: self.stats.warm_dropped.load(Ordering::Relaxed), + warm_resolved: self.stats.warm_resolved.load(Ordering::Relaxed), + warm_failed: self.stats.warm_failed.load(Ordering::Relaxed), + } + } + + pub fn warm(&self, did: &Did) { + if self.by_did.contains_sync(did) { + return; + } + self.enqueue(did, WarmKind::Fill); + } + + pub fn refresh(&self, did: &Did) { + self.enqueue(did, WarmKind::Refresh); + } + + fn enqueue(&self, did: &Did, kind: WarmKind) { + let Some(clock) = self.clock.as_ref() else { + return; + }; + let now = clock.now_instant(); + if self + .cooldown + .read_sync(did, |_, until| *until > now) + .unwrap_or(false) + { + return; + } + if self.queued.insert_sync(did.clone()).is_err() { + return; + } + if self.warm_tx.send((did.clone(), kind)).is_err() { + self.queued.remove_sync(did); + self.stats.warm_dropped.fetch_add(1, Ordering::Relaxed); + return; + } + self.stats.warm_queued.fetch_add(1, Ordering::Relaxed); + } + + pub async fn run_warming(self: Arc, cancel: CancellationToken) { + let Some(mut rx) = self + .warm_rx + .lock() + .expect("identity warm receiver mutex poisoned") + .take() + else { + return; + }; + let permits = Arc::new(Semaphore::new(WARM_CONCURRENCY)); + loop { + let (did, kind) = tokio::select! { + biased; + _ = cancel.cancelled() => break, + next = rx.recv() => match next { + Some(next) => next, + None => break, + }, + }; + let Ok(permit) = permits.clone().acquire_owned().await else { + break; + }; + let resolver = self.clone(); + tokio::spawn(async move { + let _permit = permit; + resolver.warm_one(did, kind).await; + }); + } + } + + async fn warm_one(&self, did: Did, kind: WarmKind) { + let outcome = match kind { + WarmKind::Refresh => self + .fetch_upstream(&AtIdentifier::Did(did.clone())) + .await + .map(|doc| self.observe_doc(doc)), + WarmKind::Fill => self.fill_one(&did).await, + }; + + match outcome { + Ok(()) => { + self.cooldown.remove_sync(&did); + self.stats.warm_resolved.fetch_add(1, Ordering::Relaxed); + } + Err(error) => { + if let Some(clock) = self.clock.as_ref() { + self.cooldown + .upsert_sync(did.clone(), clock.now_instant() + WARM_RETRY_AFTER); + } + self.stats.warm_failed.fetch_add(1, Ordering::Relaxed); + debug!(did = did.as_ref(), %error, "identity warm failed"); + } + } + self.queued.remove_sync(&did); + } + + async fn fill_one(&self, did: &Did) -> Result<(), IdentityResolveError> { + let known = match self.hydrant.as_ref() { + Some(hydrant) => hydrant.repo_identity(did).await.unwrap_or_else(|error| { + debug!(did = did.as_ref(), %error, "hydrant warm failed"); + None + }), + None => None, + }; + match known { + Some(RepoIdentity::Live(doc)) => { + self.observe_doc(doc); + Ok(()) + } + Some(RepoIdentity::Dead) => { + self.deactivate(did.clone()); + Ok(()) + } + None => self + .resolve_minidoc(&AtIdentifier::Did(did.clone())) + .await + .map(drop), } } pub fn observe(&self, did: Did, handle: Handle) { + self.observe_doc(MiniDoc { + did, + handle, + pds: None, + signing_key: None, + }); + } + + fn observe_doc(&self, doc: MiniDoc) { + let (did, handle) = (doc.did.clone(), doc.handle.clone()); let mut previous_handle = None; match self.by_did.entry_sync(did.clone()) { MapEntry::Occupied(mut occupied) => { - let (pds, fetched) = match occupied.get() { - IdentityState::Observed(previous) => { + let merged = match occupied.get().doc() { + Some(previous) => { previous_handle = Some(previous.handle.clone()); - (previous.pds.clone(), false) + MiniDoc { + pds: doc.pds.or_else(|| previous.pds.clone()), + signing_key: doc.signing_key.or_else(|| previous.signing_key.clone()), + ..doc + } } - IdentityState::Fetched(previous) => { - previous_handle = Some(previous.handle.clone()); - (previous.pds.clone(), previous.handle == handle) - } - IdentityState::Inactive => (None, false), - }; - let doc = MiniDoc { - did: did.clone(), - handle: handle.clone(), - pds, + None => doc, }; - occupied.insert(if fetched { - IdentityState::Fetched(doc) - } else { - IdentityState::Observed(doc) - }); + occupied.insert(IdentityState::Cached(merged)); } MapEntry::Vacant(vacant) => { - vacant.insert_entry(IdentityState::Observed(MiniDoc { - did: did.clone(), - handle: handle.clone(), - pds: None, - })); + vacant.insert_entry(IdentityState::Cached(doc)); } } self.remove_by_handle_if_owned(&did, previous_handle.as_ref()); @@ -213,12 +416,10 @@ impl IdentityResolver { fn cached( &self, identifier: &AtIdentifier, - require_fetched: bool, ) -> Result, IdentityResolveError> { match identifier { AtIdentifier::Did(did) => match self.by_did.get_sync(did).as_deref() { - Some(IdentityState::Observed(doc)) => Ok((!require_fetched).then(|| doc.clone())), - Some(IdentityState::Fetched(doc)) => Ok(Some(doc.clone())), + Some(IdentityState::Cached(doc)) => Ok(Some(doc.clone())), Some(IdentityState::Inactive) => Err(IdentityResolveError::NotFound), None => Ok(None), }, @@ -231,17 +432,10 @@ impl IdentityResolver { return Ok(None); }; let state = self.by_did.get_sync(&did).map(|state| state.get().clone()); - match state { - Some(IdentityState::Observed(doc)) if doc.handle == *handle => { - return Ok((!require_fetched).then_some(doc)); - } - Some(IdentityState::Fetched(doc)) if doc.handle == *handle => { - return Ok(Some(doc)); - } - Some(IdentityState::Observed(_)) - | Some(IdentityState::Fetched(_)) - | Some(IdentityState::Inactive) - | None => {} + if let Some(IdentityState::Cached(doc)) = state + && doc.handle == *handle + { + return Ok(Some(doc)); } self.remove_by_handle_if_owned(&did, Some(handle)); Ok(None) @@ -252,9 +446,8 @@ impl IdentityResolver { fn get_cached( &self, identifier: &AtIdentifier, - require_fetched: bool, ) -> Result, IdentityResolveError> { - let result = self.cached(identifier, require_fetched); + let result = self.cached(identifier); match &result { Ok(Some(_)) | Err(_) => self.stats.hits.fetch_add(1, Ordering::Relaxed), Ok(None) => self.stats.misses.fetch_add(1, Ordering::Relaxed), @@ -262,26 +455,23 @@ impl IdentityResolver { result } - /// Get a Hydrant-observed DID without waiting on Slingshot. + /// reads the cache without waiting, a miss queues a background warm pub fn get_by_did(&self, did: &Did) -> Result { - self.get_cached(&AtIdentifier::Did(did.clone()), false)? - .ok_or(IdentityResolveError::NotFound) + match self.get_cached(&AtIdentifier::Did(did.clone()))? { + Some(doc) => Ok(doc), + None => { + self.warm(did); + Err(IdentityResolveError::NotFound) + } + } } - /// Resolve a minidoc, fetching Hydrant-only observations upstream first. + /// resolve a minidoc, waiting on slingshot when the cache misses pub async fn resolve_minidoc( &self, identifier: &AtIdentifier, ) -> Result { - self.resolve_with_cache(identifier, true).await - } - - async fn resolve_with_cache( - &self, - identifier: &AtIdentifier, - require_fetched: bool, - ) -> Result { - if let Some(doc) = self.get_cached(identifier, require_fetched)? { + if let Some(doc) = self.get_cached(identifier)? { return Ok(doc); } @@ -301,55 +491,50 @@ impl IdentityResolver { result } - async fn fetch_minidoc( + async fn fetch_upstream( &self, identifier: &AtIdentifier, ) -> Result { - let client = self - .slingshot - .as_ref() - .ok_or(IdentityResolveError::NotFound)?; + let probe = self.probe.as_ref().ok_or(IdentityResolveError::NotFound)?; self.stats.upstream_requests.fetch_add(1, Ordering::Relaxed); - let bytes = client + let bytes = probe + .client .resolve_mini_doc(identifier) .await .map_err(IdentityResolveError::from)?; - let doc = serde_json::from_slice::(&bytes) - .map_err(|error| IdentityResolveError::Decode(error.to_string()))?; + serde_json::from_slice::(&bytes) + .map_err(|error| IdentityResolveError::Decode(error.to_string())) + } + + async fn fetch_minidoc( + &self, + identifier: &AtIdentifier, + ) -> Result { + let doc = self.fetch_upstream(identifier).await?; self.insert_fetched_by_did(doc) } fn insert_fetched_by_did(&self, doc: MiniDoc) -> Result { - let did = doc.did.clone(); - let handle = doc.handle.clone(); - let mut previous_handle = None; - let stored = match self.by_did.entry_sync(did.clone()) { - MapEntry::Occupied(mut occupied) => match occupied.get_mut() { - IdentityState::Inactive => return Err(IdentityResolveError::NotFound), - IdentityState::Observed(observed) if observed.handle != handle => { - observed.pds = doc.pds; - return Ok(observed.clone()); - } - IdentityState::Observed(previous) | IdentityState::Fetched(previous) => { - previous_handle = Some(previous.handle.clone()); - occupied.insert(IdentityState::Fetched(doc.clone())); - doc - } + match self.by_did.entry_sync(doc.did.clone()) { + MapEntry::Occupied(occupied) => match occupied.get() { + IdentityState::Inactive => Err(IdentityResolveError::NotFound), + IdentityState::Cached(current) => Ok(current.clone()), }, MapEntry::Vacant(vacant) => { - vacant.insert_entry(IdentityState::Fetched(doc.clone())); - doc + let (did, handle) = (doc.did.clone(), doc.handle.clone()); + vacant.insert_entry(IdentityState::Cached(doc.clone())); + self.insert_by_handle(did, handle); + Ok(doc) } - }; - self.remove_by_handle_if_owned(&did, previous_handle.as_ref()); - self.insert_by_handle(did, handle); - Ok(stored) + } } } #[cfg(test)] mod tests { use super::*; + use bobbin_runtime::{ReqwestHttp, SystemClock}; + use bobbin_slingshot_client::default_http_client; use url::Url; use wiremock::matchers::{method, path, query_param}; use wiremock::{Mock, MockServer, ResponseTemplate}; @@ -390,7 +575,7 @@ mod tests { assert!( resolver - .cached(&AtIdentifier::Handle(handle("ptr.pet")), false) + .cached(&AtIdentifier::Handle(handle("ptr.pet"))) .unwrap() .is_none() ); @@ -409,7 +594,7 @@ mod tests { resolver.observe(first, handle("first.example.com")); let cached = resolver - .cached(&AtIdentifier::Handle(shared), false) + .cached(&AtIdentifier::Handle(shared)) .unwrap() .unwrap(); assert_eq!(cached.did, second); @@ -427,26 +612,57 @@ mod tests { assert!( resolver - .cached(&AtIdentifier::Handle(stale.clone()), false) + .cached(&AtIdentifier::Handle(stale.clone())) .unwrap() .is_none() ); assert!(resolver.by_handle.get_sync(&stale).is_none()); } + // the codec mapping is ours, a wrong code re-encodes to a different key + #[test] + fn signing_key_survives_a_json_round_trip() { + let encoded = "zQ3shSiLsnqpyQ4SfDTT1D8qzFEoeYT8rSDXW6o8pVY7VcRBJ"; + let doc = MiniDoc { + did: did("did:plc:dawn"), + handle: handle("ptr.pet"), + pds: Some(Url::parse("https://gaze.systems/").unwrap()), + signing_key: Some(PublicKey::decode_owned(encoded).unwrap()), + }; + + let json = serde_json::to_value(&doc).unwrap(); + assert_eq!(json["signingKey"], serde_json::json!(encoded)); + assert_eq!(json["pds"], serde_json::json!("https://gaze.systems/")); + assert_eq!(serde_json::from_value::(json).unwrap(), doc); + } + + #[test] + fn absent_optional_fields_round_trip_as_absent() { + let doc = MiniDoc { + did: did("did:plc:dawn"), + handle: handle("ptr.pet"), + pds: None, + signing_key: None, + }; + + let json = serde_json::to_value(&doc).unwrap(); + assert!(json.get("signingKey").is_none()); + assert!(json.get("pds").is_none()); + assert_eq!(serde_json::from_value::(json).unwrap(), doc); + } + #[test] - fn minidoc_lookup_keeps_a_valid_observed_handle_hint() { + fn handle_lookup_keeps_a_reverse_hint_the_forward_entry_agrees_with() { let resolver = resolver(); let identity = did("did:plc:dawn"); let handle = handle("ptr.pet"); resolver.observe(identity.clone(), handle.clone()); - assert!( - resolver - .cached(&AtIdentifier::Handle(handle.clone()), true) - .unwrap() - .is_none() - ); + let found = resolver + .cached(&AtIdentifier::Handle(handle.clone())) + .unwrap() + .unwrap(); + assert_eq!(found.did, identity); assert_eq!( resolver .by_handle @@ -457,38 +673,28 @@ mod tests { } #[tokio::test] - async fn partial_observation_fetches_and_preserves_pds() { + async fn resolve_serves_a_hydrant_observed_entry_from_cache() { let server = MockServer::start().await; Mock::given(method("GET")) .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) - .and(query_param("identifier", "did:plc:dawn")) - .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ - "did": "did:plc:dawn", - "handle": "ptr.pet", - "pds": "https://pds.example.com" - }))) - .expect(1) + .respond_with(ResponseTemplate::new(500)) + .expect(0) .mount(&server) .await; let client = SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = IdentityResolver::with_slingshot(client, hasher()); + let resolver = + IdentityResolver::with_slingshot(client, Arc::new(SystemClock::new()), hasher()); let identity = did("did:plc:dawn"); resolver.observe(identity.clone(), handle("ptr.pet")); let doc = resolver - .resolve_minidoc(&AtIdentifier::Did(identity.clone())) - .await - .unwrap(); - assert_eq!(doc.pds.as_deref(), Some("https://pds.example.com")); - - resolver.observe(identity.clone(), handle("ptr.pet")); - let cached = resolver .resolve_minidoc(&AtIdentifier::Did(identity)) .await .unwrap(); - assert_eq!(cached.pds.as_deref(), Some("https://pds.example.com")); - assert_eq!(resolver.stats().upstream_requests, 1); + assert_eq!(doc.handle, handle("ptr.pet")); + assert_eq!(doc.pds, None); + assert_eq!(resolver.stats().upstream_requests, 0); } #[test] @@ -506,15 +712,321 @@ mod tests { resolver.insert_fetched_by_did(MiniDoc { did: identity, handle: handle("ptr.pet"), - pds: Some("https://pds.example.com".to_owned()), + pds: Some(Url::parse("https://pds.example.com").unwrap()), + signing_key: None, }), Err(IdentityResolveError::NotFound) ); assert!( resolver - .cached(&AtIdentifier::Handle(handle("ptr.pet")), false) + .cached(&AtIdentifier::Handle(handle("ptr.pet"))) .unwrap() .is_none() ); } + + #[tokio::test] + async fn warming_resolves_a_did_we_only_saw_on_a_record() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) + .and(query_param("identifier", "did:plc:dawn")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:dawn", + "handle": "ptr.pet", + "pds": "https://pds.example.com" + }))) + .expect(1) + .mount(&server) + .await; + let client = + SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(); + let resolver = Arc::new(IdentityResolver::with_slingshot( + client, + Arc::new(SystemClock::new()), + hasher(), + )); + let identity = did("did:plc:dawn"); + + // a plain cache read is what feeds and enrich do, and it must not block + assert_eq!( + resolver.get_by_did(&identity), + Err(IdentityResolveError::NotFound) + ); + assert_eq!(resolver.stats().warm_queued, 1); + + let cancel = CancellationToken::new(); + let warmer = tokio::spawn(resolver.clone().run_warming(cancel.clone())); + let warmed = tokio::time::timeout(Duration::from_secs(5), async { + loop { + if let Ok(doc) = resolver.get_by_did(&identity) { + return doc; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("warmer should resolve the queued did"); + assert_eq!(warmed.handle, handle("ptr.pet")); + assert_eq!( + warmed.pds, + Some(Url::parse("https://pds.example.com").unwrap()) + ); + + cancel.cancel(); + warmer.await.unwrap(); + assert_eq!(resolver.stats().warm_resolved, 1); + } + + async fn drain_warming(resolver: &Arc, did: &Did) { + let cancel = CancellationToken::new(); + let warmer = tokio::spawn(resolver.clone().run_warming(cancel.clone())); + tokio::time::timeout(Duration::from_secs(5), async { + while resolver.stats().warm_resolved + resolver.stats().warm_failed == 0 { + tokio::task::yield_now().await; + } + }) + .await + .unwrap_or_else(|_| panic!("warmer never settled {did}")); + cancel.cancel(); + warmer.await.unwrap(); + } + + fn hydrant_backed(server: &MockServer) -> Arc { + let base = Url::parse(&server.uri()).unwrap(); + Arc::new( + IdentityResolver::with_slingshot( + SlingshotClient::with_default_http(base.clone()).unwrap(), + Arc::new(SystemClock::new()), + hasher(), + ) + .with_hydrant( + HydrantClient::new(base, ReqwestHttp::shared(default_http_client().unwrap())) + .unwrap(), + ), + ) + } + + #[tokio::test] + async fn warming_prefers_hydrant_over_slingshot() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:dawn")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:dawn", + "status": "synced", + "tracked": true, + "handle": "ptr.pet", + }))) + .expect(1) + .mount(&server) + .await; + Mock::given(method("GET")) + .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) + .respond_with(ResponseTemplate::new(500)) + .expect(0) + .mount(&server) + .await; + + let resolver = hydrant_backed(&server); + let identity = did("did:plc:dawn"); + assert_eq!( + resolver.get_by_did(&identity), + Err(IdentityResolveError::NotFound) + ); + + drain_warming(&resolver, &identity).await; + assert_eq!( + resolver.get_by_did(&identity).unwrap().handle, + handle("ptr.pet") + ); + assert_eq!(resolver.stats().upstream_requests, 0); + } + + #[tokio::test] + async fn warming_falls_back_to_slingshot_for_dids_hydrant_lacks() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:stranger")) + .respond_with(ResponseTemplate::new(404)) + .expect(1) + .mount(&server) + .await; + Mock::given(method("GET")) + .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) + .and(query_param("identifier", "did:plc:stranger")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:stranger", + "handle": "stranger.example.com", + "pds": "https://pds.example.com" + }))) + .expect(1) + .mount(&server) + .await; + + let resolver = hydrant_backed(&server); + let identity = did("did:plc:stranger"); + resolver.warm(&identity); + drain_warming(&resolver, &identity).await; + + assert_eq!( + resolver.get_by_did(&identity).unwrap().handle, + handle("stranger.example.com") + ); + assert_eq!(resolver.stats().warm_resolved, 1); + } + + #[tokio::test] + async fn warming_falls_back_when_hydrant_errors() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:dawn")) + .respond_with(ResponseTemplate::new(500)) + .mount(&server) + .await; + Mock::given(method("GET")) + .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) + .and(query_param("identifier", "did:plc:dawn")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:dawn", + "handle": "ptr.pet", + "pds": "https://pds.example.com" + }))) + .expect(1) + .mount(&server) + .await; + + let resolver = hydrant_backed(&server); + let identity = did("did:plc:dawn"); + resolver.warm(&identity); + drain_warming(&resolver, &identity).await; + + assert_eq!( + resolver.get_by_did(&identity).unwrap().handle, + handle("ptr.pet") + ); + } + + #[tokio::test] + async fn refresh_re_resolves_upstream_and_skips_hydrant() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:dawn")) + .respond_with(ResponseTemplate::new(500)) + .expect(0) + .mount(&server) + .await; + Mock::given(method("GET")) + .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) + .and(query_param("identifier", "did:plc:dawn")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:dawn", + "handle": "new.ptr.pet", + "pds": "https://pds.example.com" + }))) + .expect(1) + .mount(&server) + .await; + + let resolver = hydrant_backed(&server); + let identity = did("did:plc:dawn"); + resolver.observe(identity.clone(), handle("old.ptr.pet")); + resolver.refresh(&identity); + + assert_eq!( + resolver.get_by_did(&identity).unwrap().handle, + handle("old.ptr.pet"), + "the cached handle keeps serving while the refresh is in flight" + ); + + drain_warming(&resolver, &identity).await; + assert_eq!( + resolver.get_by_did(&identity).unwrap().handle, + handle("new.ptr.pet") + ); + assert!( + resolver + .by_handle + .get_sync(&handle("old.ptr.pet")) + .is_none() + ); + } + + #[tokio::test] + async fn a_failed_refresh_leaves_the_cached_handle_alone() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) + .respond_with(ResponseTemplate::new(500)) + .mount(&server) + .await; + + let resolver = hydrant_backed(&server); + let identity = did("did:plc:dawn"); + resolver.observe(identity.clone(), handle("ptr.pet")); + resolver.refresh(&identity); + drain_warming(&resolver, &identity).await; + + assert_eq!( + resolver.get_by_did(&identity).unwrap().handle, + handle("ptr.pet") + ); + assert_eq!(resolver.stats().warm_failed, 1); + } + + #[tokio::test] + async fn warming_a_dead_repo_caches_no_handle() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/repos/did:plc:gone")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:gone", + "status": "deleted", + "tracked": true, + "handle": "gone.example.com", + }))) + .mount(&server) + .await; + Mock::given(method("GET")) + .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) + .respond_with(ResponseTemplate::new(500)) + .expect(0) + .mount(&server) + .await; + + let resolver = hydrant_backed(&server); + let identity = did("did:plc:gone"); + resolver.warm(&identity); + drain_warming(&resolver, &identity).await; + + assert_eq!( + resolver.get_by_did(&identity), + Err(IdentityResolveError::NotFound) + ); + resolver.warm(&identity); + assert_eq!(resolver.stats().warm_queued, 1); + } + + #[tokio::test] + async fn warming_queues_dids_without_dropping() { + let resolver = Arc::new(IdentityResolver::with_slingshot( + SlingshotClient::with_default_http(Url::parse("http://127.0.0.1:1").unwrap()).unwrap(), + Arc::new(SystemClock::new()), + hasher(), + )); + for n in 0..10_000 { + resolver.warm(&did(&format!("did:plc:bulk{n}"))); + } + assert_eq!(resolver.stats().warm_queued, 10_000); + assert_eq!(resolver.stats().warm_dropped, 0); + } + + #[test] + fn warming_skips_dids_we_already_know() { + let resolver = resolver(); + let identity = did("did:plc:dawn"); + resolver.observe(identity.clone(), handle("ptr.pet")); + resolver.warm(&identity); + assert_eq!(resolver.stats().warm_queued, 0); + } } diff --git a/bobbin/crates/resolver/src/lib.rs b/bobbin/crates/resolver/src/lib.rs index ba7c4fd5a..7373481c0 100644 --- a/bobbin/crates/resolver/src/lib.rs +++ b/bobbin/crates/resolver/src/lib.rs @@ -1,7 +1,9 @@ +mod hydrant; mod identity; mod legacy_upgrade; mod normalize; +pub use hydrant::{HydrantClient, HydrantError, RepoIdentity}; pub use identity::{ IdentityResolveError, IdentityResolver, IdentityResolverStatsSnapshot, MiniDoc, }; @@ -165,9 +167,9 @@ impl ResolverStats { } } -struct SlingshotProbe { - client: SlingshotClient, - clock: Arc, +pub(crate) struct SlingshotProbe { + pub(crate) client: SlingshotClient, + pub(crate) clock: Arc, } pub struct RepoIdResolver { diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index bcd591bca..771db3dc6 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -161,6 +161,7 @@ impl AppState { ) -> Self { let identity = Arc::new(IdentityResolver::with_slingshot( slingshot.clone(), + Arc::new(bobbin_runtime::SystemClock::new()), bobbin_runtime::RuntimeHasher::default(), )); Self { -- 2.51.2