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" ); } }