use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; use std::time::Duration; use bobbin_runtime::{Clock, RuntimeHasher}; use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; use bobbin_types::identity::IdentitySink; 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, 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)] pub struct MiniDoc { pub did: Did, pub handle: Handle, #[serde(skip_serializing_if = "Option::is_none")] 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)] 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()), } } } type ResolveCell = Arc>>; #[derive(Clone)] enum IdentityState { Cached(MiniDoc), Inactive { handle: Option>, }, Revalidating { handle: Option>, cell: ResolveCell, }, } impl IdentityState { fn doc(&self) -> Option<&MiniDoc> { match self { Self::Cached(doc) => Some(doc), Self::Inactive { .. } | Self::Revalidating { .. } => None, } } fn handle(&self) -> Option<&Handle> { match self { Self::Cached(doc) => Some(&doc.handle), Self::Inactive { handle } | Self::Revalidating { handle, .. } => handle.as_ref(), } } } enum Revalidation { Cached(MiniDoc), Flight(ResolveCell), Missing, Suppressed, } struct WarmQueueGuard<'a> { queued: &'a SccSet, RuntimeHasher>, did: &'a Did, } impl Drop for WarmQueueGuard<'_> { fn drop(&mut self) { self.queued.remove_sync(self.did); } } #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub struct IdentityResolverStatsSnapshot { pub entries: usize, 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)] 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 pub struct IdentityResolver { by_did: SccMap, IdentityState, RuntimeHasher>, by_handle: SccMap, Did, RuntimeHasher>, in_flight: SccMap, // 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>, sink: Option>, mutations: Mutex<()>, stats: IdentityResolverStats, } impl IdentityResolver { 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) } pub fn with_hydrant(mut self, client: HydrantClient) -> Self { self.hydrant = Some(client); self } pub fn with_sink(mut self, sink: Arc) -> Self { self.sink = Some(sink); 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.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, sink: None, mutations: Mutex::new(()), 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), 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) { let kind = match self.by_did.get_sync(did).as_deref() { Some(IdentityState::Cached(_)) => return, Some(IdentityState::Inactive { .. } | IdentityState::Revalidating { .. }) => { WarmKind::Refresh } None => WarmKind::Fill, }; self.enqueue(did, kind); } pub fn refresh(&self, did: &Did) { self.enqueue(did, WarmKind::Refresh); } fn enqueue(&self, did: &Did, kind: WarmKind) { if self.cooling_down(did) { 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); } fn cooling_down(&self, did: &Did) -> bool { let Some(clock) = self.clock.as_ref() else { return false; }; self.cooldown .read_sync(did, |_, until| *until > clock.now_instant()) .unwrap_or(false) } fn arm_cooldown(&self, did: &Did) { if let Some(clock) = self.clock.as_ref() { self.cooldown .upsert_sync(did.clone(), clock.now_instant() + WARM_RETRY_AFTER); } } 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 _queued = WarmQueueGuard { queued: &self.queued, did: &did, }; let (outcome, arm_on_error) = match kind { _ if self.needs_revalidation(&did) => { (self.revalidate_did(&did).await.map(drop), false) } WarmKind::Refresh => ( match self.fetch_upstream(&AtIdentifier::Did(did.clone())).await { Ok(doc) => self.accept_fetched(doc, true).await.map(drop), Err(error) => Err(error), }, true, ), WarmKind::Fill => (self.fill_one(&did).await, true), }; match outcome { Ok(()) => { self.cooldown.remove_sync(&did); self.stats.warm_resolved.fetch_add(1, Ordering::Relaxed); } Err(error) => { if arm_on_error { self.arm_cooldown(&did); } self.stats.warm_failed.fetch_add(1, Ordering::Relaxed); debug!(did = did.as_ref(), %error, "identity warm failed"); } } } 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.accept_fetched(doc, true).await.map(drop), 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) { let observed = MiniDoc { did: did.clone(), handle, pds: None, signing_key: None, }; let refresh = { let _mutation = self .mutations .lock() .expect("identity mutation mutex poisoned"); let state = self.by_did.get_sync(&did).map(|state| state.get().clone()); match state { Some(IdentityState::Inactive { .. }) => true, Some(IdentityState::Revalidating { .. }) => true, Some(IdentityState::Cached(doc)) => { let incomplete = doc.pds.is_none(); self.observe_doc_locked(observed); incomplete } None => { self.observe_doc_locked(observed); false } } }; if refresh { self.refresh(&did); } } fn observe_doc_locked(&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) => { previous_handle = occupied.get().handle().cloned(); let merged = match occupied.get().doc() { Some(previous) => MiniDoc { pds: doc.pds.or_else(|| previous.pds.clone()), signing_key: doc.signing_key.or_else(|| previous.signing_key.clone()), ..doc }, None => doc, }; occupied.insert(IdentityState::Cached(merged)); } MapEntry::Vacant(vacant) => { vacant.insert_entry(IdentityState::Cached(doc)); } } self.remove_by_handle_if_owned(&did, previous_handle.as_ref()); self.insert_by_handle(did.clone(), handle.clone()); if let Some(sink) = self.sink.as_ref() { sink.identity_changed(&did, previous_handle.as_ref(), &handle); } } pub fn deactivate(&self, did: Did) { let _mutation = self .mutations .lock() .expect("identity mutation mutex poisoned"); let previous_handle = self .by_did .get_sync(&did) .and_then(|state| state.get().handle().cloned()); self.by_did.upsert_sync( did.clone(), IdentityState::Inactive { handle: previous_handle.clone(), }, ); if let Some(sink) = self.sink.as_ref() { sink.identity_removed(&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(handle) = state.handle() else { return; }; self.remove_by_handle_if_owned(&did, Some(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, ) -> Result, IdentityResolveError> { match identifier { AtIdentifier::Did(did) => Ok(self .by_did .get_sync(did) .and_then(|state| state.get().doc().cloned())), 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.as_ref() { Some(IdentityState::Cached(doc)) if doc.handle == *handle => { Ok(Some(doc.clone())) } Some(state) if state.handle() == Some(handle) => Ok(None), _ => { self.remove_by_handle_if_owned(&did, Some(handle)); Ok(None) } } } } } fn get_cached( &self, identifier: &AtIdentifier, ) -> Result, IdentityResolveError> { 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), }; result } fn needs_revalidation(&self, did: &Did) -> bool { self.by_did.get_sync(did).is_some_and(|state| { matches!( state.get(), IdentityState::Inactive { .. } | IdentityState::Revalidating { .. } ) }) } fn inactive_did_for_handle(&self, handle: &Handle) -> Option> { let did = self.by_handle.get_sync(handle)?.get().clone(); self.by_did .get_sync(&did) .is_some_and(|state| { matches!( state.get(), IdentityState::Inactive { .. } | IdentityState::Revalidating { .. } ) && state.get().handle() == Some(handle) }) .then_some(did) } /// reads the cache without waiting, a miss queues a background warm pub fn get_by_did(&self, did: &Did) -> Result { match self.get_cached(&AtIdentifier::Did(did.clone()))? { Some(doc) => Ok(doc), None => { self.warm(did); Err(IdentityResolveError::NotFound) } } } /// resolve a minidoc, waiting on slingshot when the cache misses pub async fn resolve_minidoc( &self, identifier: &AtIdentifier, ) -> Result { if let Some(doc) = self.get_cached(identifier)? { return Ok(doc); } match identifier { AtIdentifier::Did(did) if self.needs_revalidation(did) => { self.revalidate_did(did).await } AtIdentifier::Handle(handle) => { if let Some(did) = self.inactive_did_for_handle(handle) { let doc = self.revalidate_did(&did).await?; return (doc.handle == *handle) .then_some(doc) .ok_or(IdentityResolveError::NotFound); } self.resolve_uncached(identifier).await } AtIdentifier::Did(_) => self.resolve_uncached(identifier).await, } } async fn resolve_uncached( &self, identifier: &AtIdentifier, ) -> Result { 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_if_async(&key, |current| Arc::ptr_eq(current, &cell)) .await; result } async fn fetch_upstream( &self, identifier: &AtIdentifier, ) -> Result { let probe = self.probe.as_ref().ok_or(IdentityResolveError::NotFound)?; self.stats.upstream_requests.fetch_add(1, Ordering::Relaxed); let bytes = probe .client .resolve_mini_doc(identifier) .await .map_err(IdentityResolveError::from)?; 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.accept_fetched(doc, false).await } async fn accept_fetched( &self, doc: MiniDoc, overwrite_cached: bool, ) -> Result { { let _mutation = self .mutations .lock() .expect("identity mutation mutex poisoned"); let state = self .by_did .get_sync(&doc.did) .map(|state| state.get().clone()); match state { Some(IdentityState::Inactive { .. } | IdentityState::Revalidating { .. }) => {} Some(IdentityState::Cached(current)) if !overwrite_cached => return Ok(current), Some(IdentityState::Cached(_)) | None => { self.observe_doc_locked(doc.clone()); return Ok(doc); } } } self.revalidate_did(&doc.did).await } fn start_revalidation(&self, did: &Did) -> Revalidation { let _mutation = self .mutations .lock() .expect("identity mutation mutex poisoned"); match self.by_did.entry_sync(did.clone()) { MapEntry::Occupied(mut occupied) => match occupied.get() { IdentityState::Cached(doc) => Revalidation::Cached(doc.clone()), IdentityState::Inactive { .. } if self.cooling_down(did) => { Revalidation::Suppressed } IdentityState::Inactive { handle } => { let handle = handle.clone(); let cell = Arc::new(OnceCell::new()); occupied.insert(IdentityState::Revalidating { handle, cell: cell.clone(), }); Revalidation::Flight(cell) } IdentityState::Revalidating { cell, .. } => Revalidation::Flight(cell.clone()), }, MapEntry::Vacant(_) => Revalidation::Missing, } } async fn revalidate_did(&self, did: &Did) -> Result { match self.start_revalidation(did) { Revalidation::Cached(doc) => Ok(doc), Revalidation::Flight(cell) => { let owner = cell.clone(); // share whether the fetch was accepted, not just what it returned cell.get_or_init(|| async { let result = self.fetch_upstream(&AtIdentifier::Did(did.clone())).await; self.finish_revalidation(did, &owner, result) }) .await .clone() } Revalidation::Missing | Revalidation::Suppressed => Err(IdentityResolveError::NotFound), } } fn finish_revalidation( &self, did: &Did, owner: &ResolveCell, result: Result, ) -> Result { let _mutation = self .mutations .lock() .expect("identity mutation mutex poisoned"); let state = self.by_did.get_sync(did).map(|state| state.get().clone()); let handle = match state { Some(IdentityState::Revalidating { handle, cell }) if Arc::ptr_eq(&cell, owner) => { handle } Some(IdentityState::Cached(doc)) => return Ok(doc), Some(IdentityState::Inactive { .. } | IdentityState::Revalidating { .. }) | None => { return Err(IdentityResolveError::NotFound); } }; match result { Ok(doc) if doc.did != *did => { self.by_did .upsert_sync(did.clone(), IdentityState::Inactive { handle }); self.arm_cooldown(did); Err(IdentityResolveError::Decode(format!( "resolved {} while revalidating {}", doc.did.as_str(), did.as_str() ))) } Ok(doc) => { self.cooldown.remove_sync(did); self.observe_doc_locked(doc.clone()); Ok(doc) } Err(error) => { self.by_did .upsert_sync(did.clone(), IdentityState::Inactive { handle }); self.arm_cooldown(did); Err(error) } } } } #[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}; 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"))) .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)) .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())) .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["signing_key"], 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("signing_key").is_none()); assert!(json.get("pds").is_none()); assert_eq!(serde_json::from_value::(json).unwrap(), doc); } #[test] 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()); let found = resolver .cached(&AtIdentifier::Handle(handle.clone())) .unwrap() .unwrap(); assert_eq!(found.did, identity); assert_eq!( resolver .by_handle .get_sync(&handle) .map(|entry| entry.get().clone()), Some(identity) ); } #[tokio::test] async fn resolve_serves_a_hydrant_observed_entry_from_cache() { let server = MockServer::start().await; Mock::given(method("GET")) .and(path("/xrpc/blue.microcosm.identity.resolveMiniDoc")) .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, 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)) .await .unwrap(); assert_eq!(doc.handle, handle("ptr.pet")); assert_eq!(doc.pds, None); assert_eq!(resolver.stats().upstream_requests, 0); } #[test] fn deactivation_cancels_an_in_flight_revalidation() { let resolver = resolver(); let identity = did("did:plc:dawn"); resolver.observe(identity.clone(), handle("ptr.pet")); resolver.deactivate(identity.clone()); let Revalidation::Flight(owner) = resolver.start_revalidation(&identity) else { panic!("inactive identity did not start revalidation"); }; resolver.deactivate(identity.clone()); assert_eq!( resolver.finish_revalidation( &identity, &owner, Ok(MiniDoc { did: identity.clone(), handle: handle("ptr.pet"), pds: Some(Url::parse("https://pds.example.com").unwrap()), signing_key: None, }) ), Err(IdentityResolveError::NotFound) ); assert!( resolver .cached(&AtIdentifier::Handle(handle("ptr.pet"))) .unwrap() .is_none() ); } #[test] fn an_identity_frame_does_not_cancel_a_live_revalidation() { let resolver = resolver(); let identity = did("did:plc:dawn"); resolver.observe(identity.clone(), handle("ptr.pet")); resolver.deactivate(identity.clone()); let Revalidation::Flight(owner) = resolver.start_revalidation(&identity) else { panic!("inactive identity did not start revalidation"); }; resolver.observe(identity.clone(), handle("ptr.pet")); let doc = MiniDoc { did: identity.clone(), handle: handle("ptr.pet"), pds: Some(Url::parse("https://pds.example.com").unwrap()), signing_key: None, }; assert_eq!( resolver.finish_revalidation(&identity, &owner, Ok(doc.clone())), Ok(doc) ); assert_eq!( resolver.get_by_did(&identity).unwrap().handle, handle("ptr.pet") ); } #[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/blue.microcosm.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 settled_before = resolver.stats().warm_resolved + resolver.stats().warm_failed; 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 == settled_before { 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/blue.microcosm.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/blue.microcosm.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/blue.microcosm.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/blue.microcosm.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/blue.microcosm.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 a_deactivated_did_reverifies_on_read() { let server = MockServer::start().await; Mock::given(method("GET")) .and(path("/xrpc/blue.microcosm.identity.resolveMiniDoc")) .and(query_param("identifier", "did:plc:dawn")) .respond_with( ResponseTemplate::new(200) .set_delay(Duration::from_millis(100)) .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.deactivate(identity.clone()); assert_eq!( resolver.get_by_did(&identity), Err(IdentityResolveError::NotFound), "the read queues a revalidation instead of waiting on it" ); assert_eq!(resolver.stats().warm_queued, 1); let cancel = CancellationToken::new(); let warmer = tokio::spawn(resolver.clone().run_warming(cancel.clone())); tokio::time::timeout(Duration::from_secs(5), async { while !matches!( resolver.by_did.get_sync(&identity).as_deref(), Some(IdentityState::Revalidating { .. }) ) { tokio::task::yield_now().await; } }) .await .expect("background revalidation never started"); let doc = resolver .resolve_minidoc(&AtIdentifier::Did(identity.clone())) .await .unwrap(); tokio::time::timeout(Duration::from_secs(5), async { while resolver.stats().warm_resolved == 0 { tokio::task::yield_now().await; } }) .await .expect("background revalidation never settled"); cancel.cancel(); warmer.await.unwrap(); assert_eq!(doc.handle, handle("ptr.pet")); assert_eq!(resolver.stats().upstream_requests, 1); assert_eq!( doc.pds, Some(Url::parse("https://pds.example.com").unwrap()) ); } #[tokio::test] async fn a_deactivated_did_that_stays_dead_is_not_refetched_per_read() { let server = MockServer::start().await; Mock::given(method("GET")) .and(path("/xrpc/blue.microcosm.identity.resolveMiniDoc")) .respond_with(ResponseTemplate::new(404)) .expect(1) .mount(&server) .await; let resolver = hydrant_backed(&server); let identity = did("did:plc:gone"); resolver.deactivate(identity.clone()); for _ in 0..3 { assert_eq!( resolver.get_by_did(&identity), Err(IdentityResolveError::NotFound) ); } assert_eq!(resolver.stats().warm_queued, 1); drain_warming(&resolver, &identity).await; for _ in 0..3 { assert_eq!( resolver.get_by_did(&identity), Err(IdentityResolveError::NotFound) ); } let stats = resolver.stats(); assert_eq!(stats.warm_queued, 1); assert_eq!(stats.warm_failed, 1); } #[tokio::test] async fn failed_foreground_revalidation_obeys_the_cooldown() { let server = MockServer::start().await; Mock::given(method("GET")) .and(path("/xrpc/blue.microcosm.identity.resolveMiniDoc")) .and(query_param("identifier", "did:plc:gone")) .respond_with(ResponseTemplate::new(404)) .expect(1) .mount(&server) .await; let resolver = hydrant_backed(&server); let identity = did("did:plc:gone"); resolver.deactivate(identity.clone()); let identifier = AtIdentifier::Did(identity); for _ in 0..3 { assert_eq!( resolver.resolve_minidoc(&identifier).await, Err(IdentityResolveError::NotFound) ); } assert_eq!(resolver.stats().upstream_requests, 1); } #[tokio::test] async fn did_and_handle_lookups_share_an_inactive_revalidation() { let server = MockServer::start().await; Mock::given(method("GET")) .and(path("/xrpc/blue.microcosm.identity.resolveMiniDoc")) .and(query_param("identifier", "dawn.example.com")) .respond_with(ResponseTemplate::new(404)) .expect(0) .mount(&server) .await; Mock::given(method("GET")) .and(path("/xrpc/blue.microcosm.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": "dawn.example.com", "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("dawn.example.com")); resolver.deactivate(identity.clone()); let by_handle = AtIdentifier::Handle(handle("dawn.example.com")); let by_did = AtIdentifier::Did(identity); let (handle_result, did_result) = tokio::join!( resolver.resolve_minidoc(&by_handle), resolver.resolve_minidoc(&by_did) ); assert_eq!(handle_result.unwrap().handle, handle("dawn.example.com")); assert_eq!(did_result.unwrap().handle, handle("dawn.example.com")); assert_eq!(resolver.stats().upstream_requests, 1); assert_eq!(resolver.stats().warm_queued, 0); } #[tokio::test] async fn an_inactive_handle_hint_cannot_claim_a_renamed_identity() { let server = MockServer::start().await; Mock::given(method("GET")) .and(path("/xrpc/blue.microcosm.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.example.com", "pds": "https://pds.example.com" }))) .expect(1) .mount(&server) .await; let resolver = hydrant_backed(&server); let identity = did("did:plc:dawn"); let old_handle = handle("old.example.com"); resolver.observe(identity.clone(), old_handle.clone()); resolver.deactivate(identity); assert_eq!( resolver .resolve_minidoc(&AtIdentifier::Handle(old_handle.clone())) .await, Err(IdentityResolveError::NotFound) ); assert!(resolver.by_handle.get_sync(&old_handle).is_none()); assert_eq!( resolver .cached(&AtIdentifier::Handle(handle("new.example.com"))) .unwrap() .unwrap() .handle, handle("new.example.com") ); } #[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/blue.microcosm.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); } #[derive(Clone, Debug, Eq, PartialEq)] enum SinkEvent { Changed { did: Did, previous_handle: Option>, handle: Handle, }, Removed { did: Did, previous_handle: Option>, }, } #[derive(Default)] struct RecordingSink { events: Mutex>, } impl IdentitySink for RecordingSink { fn identity_changed( &self, did: &Did, previous_handle: Option<&Handle>, handle: &Handle, ) { self.events.lock().unwrap().push(SinkEvent::Changed { did: did.clone(), previous_handle: previous_handle.cloned(), handle: handle.clone(), }); } fn identity_removed( &self, did: &Did, previous_handle: Option<&Handle>, ) { self.events.lock().unwrap().push(SinkEvent::Removed { did: did.clone(), previous_handle: previous_handle.cloned(), }); } } struct IdentityCheckingSink { resolver: Arc>>, events: Mutex>, } impl IdentitySink for IdentityCheckingSink { fn identity_changed( &self, did: &Did, previous_handle: Option<&Handle>, handle: &Handle, ) { let resolver = self.resolver.get().expect("resolver initialized"); let doc = resolver.get_by_did(did).expect("did must resolve"); assert_eq!(&doc.handle, handle); let by_handle = resolver .cached(&AtIdentifier::Handle(handle.clone())) .expect("cache lookup must succeed") .expect("handle must resolve"); assert_eq!(&by_handle.did, did); if let Some(prev) = previous_handle && prev != handle { assert!( resolver .cached(&AtIdentifier::Handle(prev.clone())) .expect("cache lookup must succeed") .is_none() ); } self.events.lock().unwrap().push(SinkEvent::Changed { did: did.clone(), previous_handle: previous_handle.cloned(), handle: handle.clone(), }); } fn identity_removed( &self, did: &Did, previous_handle: Option<&Handle>, ) { let resolver = self.resolver.get().expect("resolver initialized"); assert_eq!( resolver.get_by_did(did), Err(IdentityResolveError::NotFound) ); if let Some(prev) = previous_handle { assert!( resolver .cached(&AtIdentifier::Handle(prev.clone())) .expect("cache lookup must succeed") .is_none() ); } self.events.lock().unwrap().push(SinkEvent::Removed { did: did.clone(), previous_handle: previous_handle.cloned(), }); } } #[test] fn callback_runs_after_identity_maps_are_updated() { let cell = Arc::new(OnceCell::new()); let sink = Arc::new(IdentityCheckingSink { resolver: cell.clone(), events: Mutex::new(Vec::new()), }); let resolver = Arc::new(resolver().with_sink(sink.clone())); cell.set(resolver.clone()).ok().expect("cell set"); let identity = did("did:plc:dawn"); let initial_handle = handle("old.ptr.pet"); let updated_handle = handle("new.ptr.pet"); // 1. Initial observation resolver.observe(identity.clone(), initial_handle.clone()); // 2. Rename resolver.observe(identity.clone(), updated_handle.clone()); // 3. Deactivate resolver.deactivate(identity.clone()); let checked = sink.events.lock().unwrap().clone(); assert_eq!(checked.len(), 3); assert_eq!( checked[0], SinkEvent::Changed { did: identity.clone(), previous_handle: None, handle: initial_handle.clone(), } ); assert_eq!( checked[1], SinkEvent::Changed { did: identity.clone(), previous_handle: Some(initial_handle), handle: updated_handle.clone(), } ); assert_eq!( checked[2], SinkEvent::Removed { did: identity, previous_handle: Some(updated_handle), } ); } #[tokio::test] async fn hydrant_results_are_sent_to_identity_sink() { let server = MockServer::start().await; Mock::given(method("GET")) .and(path("/repos/did:plc:live")) .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ "did": "did:plc:live", "status": "active", "tracked": true, "handle": "live.example.com", }))) .mount(&server) .await; Mock::given(method("GET")) .and(path("/repos/did:plc:dead")) .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ "did": "did:plc:dead", "status": "deleted", "tracked": true, "handle": "dead.example.com", }))) .mount(&server) .await; let sink = Arc::new(RecordingSink::default()); let base = Url::parse(&server.uri()).unwrap(); let resolver = 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(), ) .with_sink(sink.clone()), ); let live_did = did("did:plc:live"); resolver.fill_one(&live_did).await.unwrap(); let dead_did = did("did:plc:dead"); resolver.fill_one(&dead_did).await.unwrap(); let events = sink.events.lock().unwrap().clone(); assert_eq!( events, vec![ SinkEvent::Changed { did: live_did, previous_handle: None, handle: handle("live.example.com"), }, SinkEvent::Removed { did: dead_did, previous_handle: None, }, ] ); } }