Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460use 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<Handle<DefaultStr>>, #[serde(default)] pds: Option<Url>, #[serde(default, with = "crate::identity::signing_key")] signing_key: Option<PublicKey<'static>>, #[serde(default)] status: Option<RepoStatus>,}
#[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<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> { 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<dyn HttpTransport>, 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<dyn HttpTransport>) -> Result<Self, HydrantError> { 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<DefaultStr>) -> Result<Url, HydrantError> { 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<DefaultStr>, ) -> Result<Option<RepoIdentity>, 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::<RepoInfo>(&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<BytesMut, HydrantError> { 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<DefaultStr> { 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<RepoIdentity>) -> Option<Handle<DefaultStr>> { 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" ); }}