use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; use bobbin_runtime::RuntimeHasher; use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; use jacquard_common::DefaultStr; 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::hash_map::Entry as MapEntry; use serde::{Deserialize, Serialize}; use thiserror::Error; use tokio::sync::OnceCell; #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] #[serde(rename_all = "camelCase")] pub struct MiniDoc { pub did: Did, pub handle: Handle, #[serde(skip_serializing_if = "Option::is_none")] pub pds: Option, } #[derive(Clone, Debug, Error, Eq, PartialEq)] pub enum IdentityResolveError { #[error("identity not found")] NotFound, #[error("identity upstream: {0}")] Upstream(String), #[error("invalid identity response: {0}")] Decode(String), } impl From for IdentityResolveError { fn from(error: SlingshotError) -> Self { match error { SlingshotError::NotFound => Self::NotFound, other => Self::Upstream(other.to_string()), } } } // 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), Inactive, } impl IdentityState { fn doc(&self) -> Option<&MiniDoc> { match self { Self::Observed(doc) | Self::Fetched(doc) => Some(doc), Self::Inactive => None, } } } #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub struct IdentityResolverStatsSnapshot { pub entries: usize, pub hits: u64, pub misses: u64, pub upstream_requests: u64, } #[derive(Default)] struct IdentityResolverStats { hits: AtomicU64, misses: AtomicU64, upstream_requests: AtomicU64, } // in the future this would spill to disk probably pub struct IdentityResolver { by_did: SccMap, IdentityState, RuntimeHasher>, by_handle: SccMap, Did, RuntimeHasher>, in_flight: SccMap>>, RuntimeHasher>, slingshot: Option, stats: IdentityResolverStats, } impl IdentityResolver { pub fn with_slingshot(slingshot: SlingshotClient, hasher: RuntimeHasher) -> Self { Self::new(Some(slingshot), hasher) } pub fn detached(hasher: RuntimeHasher) -> Self { Self::new(None, hasher) } fn new(slingshot: Option, hasher: RuntimeHasher) -> Self { Self { by_did: SccMap::with_hasher(hasher.clone()), by_handle: SccMap::with_hasher(hasher.clone()), in_flight: SccMap::with_hasher(hasher), slingshot, stats: IdentityResolverStats::default(), } } pub fn stats(&self) -> IdentityResolverStatsSnapshot { IdentityResolverStatsSnapshot { entries: self.by_did.len(), hits: self.stats.hits.load(Ordering::Relaxed), misses: self.stats.misses.load(Ordering::Relaxed), upstream_requests: self.stats.upstream_requests.load(Ordering::Relaxed), } } pub fn observe(&self, did: Did, handle: Handle) { 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) => { previous_handle = Some(previous.handle.clone()); (previous.pds.clone(), false) } 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, }; occupied.insert(if fetched { IdentityState::Fetched(doc) } else { IdentityState::Observed(doc) }); } MapEntry::Vacant(vacant) => { vacant.insert_entry(IdentityState::Observed(MiniDoc { did: did.clone(), handle: handle.clone(), pds: None, })); } } self.remove_by_handle_if_owned(&did, previous_handle.as_ref()); self.insert_by_handle(did, handle); } pub fn deactivate(&self, did: Did) { let mut previous_handle = None; match self.by_did.entry_sync(did.clone()) { MapEntry::Occupied(mut occupied) => { previous_handle = occupied.get().doc().map(|doc| doc.handle.clone()); occupied.insert(IdentityState::Inactive); } MapEntry::Vacant(vacant) => { vacant.insert_entry(IdentityState::Inactive); } } self.remove_by_handle_if_owned(&did, previous_handle.as_ref()); } fn remove_by_handle_if_owned( &self, did: &Did, handle: Option<&Handle>, ) { let Some(handle) = handle else { return; }; self.by_handle.remove_if_sync(handle, |owner| owner == did); } fn remove_by_handle_for_removed_did(&self, removed: Option<(Did, IdentityState)>) { let Some((did, state)) = removed else { return; }; let Some(doc) = state.doc() else { return; }; self.remove_by_handle_if_owned(&did, Some(&doc.handle)); } fn insert_by_handle(&self, did: Did, handle: Handle) { if let Some(displaced_did) = self .by_handle .upsert_sync(handle.clone(), did.clone()) .filter(|displaced_did| displaced_did != &did) { let removed = self.by_did.remove_if_sync(&displaced_did, |state| { state.doc().is_some_and(|doc| doc.handle == handle) }); self.remove_by_handle_for_removed_did(removed); } if self.by_did_matches_handle(&did, &handle) { return; } self.remove_by_handle_if_owned(&did, Some(&handle)); } fn by_did_matches_handle(&self, did: &Did, handle: &Handle) -> bool { self.by_did .get_sync(did) .is_some_and(|state| state.get().doc().is_some_and(|doc| doc.handle == *handle)) } 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::Inactive) => Err(IdentityResolveError::NotFound), None => Ok(None), }, AtIdentifier::Handle(handle) => { let Some(did) = self .by_handle .get_sync(handle) .map(|entry| entry.get().clone()) else { 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 => {} } self.remove_by_handle_if_owned(&did, Some(handle)); Ok(None) } } } fn get_cached( &self, identifier: &AtIdentifier, require_fetched: bool, ) -> Result, IdentityResolveError> { let result = self.cached(identifier, require_fetched); match &result { Ok(Some(_)) | Err(_) => self.stats.hits.fetch_add(1, Ordering::Relaxed), Ok(None) => self.stats.misses.fetch_add(1, Ordering::Relaxed), }; result } /// Get a Hydrant-observed DID without waiting on Slingshot. pub fn get_by_did(&self, did: &Did) -> Result { self.get_cached(&AtIdentifier::Did(did.clone()), false)? .ok_or(IdentityResolveError::NotFound) } /// Resolve a minidoc, fetching Hydrant-only observations upstream first. 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)? { return Ok(doc); } let key = identifier.as_str().to_owned(); let cell = self .in_flight .entry_async(key.clone()) .await .or_insert_with(|| Arc::new(OnceCell::new())) .get() .clone(); let result = cell .get_or_init(|| async { self.fetch_minidoc(identifier).await }) .await .clone(); self.in_flight.remove_async(&key).await; result } async fn fetch_minidoc( &self, identifier: &AtIdentifier, ) -> Result { let client = self .slingshot .as_ref() .ok_or(IdentityResolveError::NotFound)?; self.stats.upstream_requests.fetch_add(1, Ordering::Relaxed); let bytes = 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()))?; 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 } }, MapEntry::Vacant(vacant) => { vacant.insert_entry(IdentityState::Fetched(doc.clone())); 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 url::Url; use wiremock::matchers::{method, path, query_param}; use wiremock::{Mock, MockServer, ResponseTemplate}; fn hasher() -> RuntimeHasher { RuntimeHasher::from_seeds(1, 2, 3, 4) } fn resolver() -> IdentityResolver { IdentityResolver::detached(hasher()) } fn did(value: &str) -> Did { Did::new_owned(value).unwrap() } fn handle(value: &str) -> Handle { Handle::new_owned(value).unwrap() } #[test] fn observed_identity_resolves_by_did_without_upstream() { let resolver = resolver(); resolver.observe(did("did:plc:dawn"), handle("ptr.pet")); let doc = resolver.get_by_did(&did("did:plc:dawn")).unwrap(); assert_eq!(doc.handle, handle("ptr.pet")); assert_eq!(doc.pds, None); assert_eq!(resolver.stats().hits, 1); } #[test] fn observed_identity_updates_existing_did_and_removes_old_handle() { let resolver = resolver(); let identity = did("did:plc:dawn"); resolver.observe(identity.clone(), handle("ptr.pet")); resolver.observe(identity.clone(), handle("new.ptr.pet")); assert!( resolver .cached(&AtIdentifier::Handle(handle("ptr.pet")), false) .unwrap() .is_none() ); let updated = resolver.get_by_did(&identity).unwrap(); assert_eq!(updated.handle, handle("new.ptr.pet")); } #[test] fn handle_reassignment_keeps_the_new_owner() { let resolver = resolver(); let first = did("did:plc:first"); let second = did("did:plc:second"); let shared = handle("shared.example.com"); resolver.observe(first.clone(), shared.clone()); resolver.observe(second.clone(), shared.clone()); resolver.observe(first, handle("first.example.com")); let cached = resolver .cached(&AtIdentifier::Handle(shared), false) .unwrap() .unwrap(); assert_eq!(cached.did, second); } #[test] fn handle_lookup_discards_an_unvalidated_reverse_hint() { let resolver = resolver(); let identity = did("did:plc:dawn"); let stale = handle("stale.example.com"); resolver.observe(identity.clone(), handle("current.example.com")); resolver .by_handle .upsert_sync(stale.clone(), identity.clone()); assert!( resolver .cached(&AtIdentifier::Handle(stale.clone()), false) .unwrap() .is_none() ); assert!(resolver.by_handle.get_sync(&stale).is_none()); } #[test] fn minidoc_lookup_keeps_a_valid_observed_handle_hint() { 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() ); assert_eq!( resolver .by_handle .get_sync(&handle) .map(|entry| entry.get().clone()), Some(identity) ); } #[tokio::test] async fn partial_observation_fetches_and_preserves_pds() { 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 = IdentityResolver::with_slingshot(client, 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); } #[test] fn inactive_identity_rejects_an_in_flight_result() { let resolver = resolver(); let identity = did("did:plc:dawn"); resolver.observe(identity.clone(), handle("ptr.pet")); resolver.deactivate(identity.clone()); assert_eq!( resolver.get_by_did(&identity), Err(IdentityResolveError::NotFound) ); assert_eq!( resolver.insert_fetched_by_did(MiniDoc { did: identity, handle: handle("ptr.pet"), pds: Some("https://pds.example.com".to_owned()), }), Err(IdentityResolveError::NotFound) ); assert!( resolver .cached(&AtIdentifier::Handle(handle("ptr.pet")), false) .unwrap() .is_none() ); } }