Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716use 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 didconst WARM_RETRY_AFTER: Duration = Duration::from_secs(300);
#[derive(Clone, Copy)]enum WarmKind { Fill, Refresh,}
type WarmItem = (Did<DefaultStr>, WarmKind);
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]pub struct MiniDoc { pub did: Did<DefaultStr>, pub handle: Handle<DefaultStr>, #[serde(skip_serializing_if = "Option::is_none")] pub pds: Option<Url>, #[serde(default, with = "signing_key", skip_serializing_if = "Option::is_none")] pub signing_key: Option<PublicKey<'static>>,}
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<S: Serializer>( key: &Option<PublicKey<'static>>, serializer: S, ) -> Result<S::Ok, S::Error> { 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<Option<PublicKey<'static>>, D::Error> { let Some(raw) = Option::<std::borrow::Cow<'de, str>>::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<SlingshotError> for IdentityResolveError { fn from(error: SlingshotError) -> Self { match error { SlingshotError::NotFound => Self::NotFound, other => Self::Upstream(other.to_string()), } }}
type ResolveCell = Arc<OnceCell<Result<MiniDoc, IdentityResolveError>>>;
#[derive(Clone)]enum IdentityState { Cached(MiniDoc), Inactive { handle: Option<Handle<DefaultStr>>, }, Revalidating { handle: Option<Handle<DefaultStr>>, 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<DefaultStr>> { 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<Did<DefaultStr>, RuntimeHasher>, did: &'a Did<DefaultStr>,}
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 probablypub struct IdentityResolver { by_did: SccMap<Did<DefaultStr>, IdentityState, RuntimeHasher>, by_handle: SccMap<Handle<DefaultStr>, Did<DefaultStr>, RuntimeHasher>, in_flight: SccMap<String, ResolveCell, RuntimeHasher>, // dids queued for background warming, keeps the queue free of duplicates queued: SccSet<Did<DefaultStr>, RuntimeHasher>, // dids whose last warm failed, mapped to when we may try again cooldown: SccMap<Did<DefaultStr>, Instant, RuntimeHasher>, warm_tx: mpsc::UnboundedSender<WarmItem>, warm_rx: Mutex<Option<mpsc::UnboundedReceiver<WarmItem>>>, hydrant: Option<HydrantClient>, probe: Option<SlingshotProbe>, clock: Option<Arc<dyn Clock>>, sink: Option<Arc<dyn IdentitySink>>, mutations: Mutex<()>, stats: IdentityResolverStats,}
impl IdentityResolver { pub fn with_slingshot( client: SlingshotClient, clock: Arc<dyn Clock>, 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<dyn IdentitySink>) -> Self { self.sink = Some(sink); self }
fn new(probe: Option<SlingshotProbe>, 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<DefaultStr>) { 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<DefaultStr>) { self.enqueue(did, WarmKind::Refresh); }
fn enqueue(&self, did: &Did<DefaultStr>, 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<DefaultStr>) -> 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<DefaultStr>) { 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<Self>, 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<DefaultStr>, 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<DefaultStr>) -> 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<DefaultStr>, handle: Handle<DefaultStr>) { 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<DefaultStr>) { 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<DefaultStr>, handle: Option<&Handle<DefaultStr>>, ) { 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<DefaultStr>, 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<DefaultStr>, handle: Handle<DefaultStr>) { 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<DefaultStr>, handle: &Handle<DefaultStr>) -> 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<DefaultStr>, ) -> Result<Option<MiniDoc>, 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<DefaultStr>, ) -> Result<Option<MiniDoc>, 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<DefaultStr>) -> 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<DefaultStr>) -> Option<Did<DefaultStr>> { 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<DefaultStr>) -> Result<MiniDoc, IdentityResolveError> { 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<DefaultStr>, ) -> Result<MiniDoc, IdentityResolveError> { 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<DefaultStr>, ) -> Result<MiniDoc, IdentityResolveError> { 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<DefaultStr>, ) -> Result<MiniDoc, IdentityResolveError> { 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::<MiniDoc>(&bytes) .map_err(|error| IdentityResolveError::Decode(error.to_string())) }
async fn fetch_minidoc( &self, identifier: &AtIdentifier<DefaultStr>, ) -> Result<MiniDoc, IdentityResolveError> { let doc = self.fetch_upstream(identifier).await?; self.accept_fetched(doc, false).await }
async fn accept_fetched( &self, doc: MiniDoc, overwrite_cached: bool, ) -> Result<MiniDoc, IdentityResolveError> { { 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<DefaultStr>) -> 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<DefaultStr>) -> Result<MiniDoc, IdentityResolveError> { 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<DefaultStr>, owner: &ResolveCell, result: Result<MiniDoc, IdentityResolveError>, ) -> Result<MiniDoc, IdentityResolveError> { 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<DefaultStr> { Did::new_owned(value).unwrap() }
fn handle(value: &str) -> Handle<DefaultStr> { 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::<MiniDoc>(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::<MiniDoc>(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<IdentityResolver>, did: &Did<DefaultStr>) { 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<IdentityResolver> { 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<DefaultStr>, previous_handle: Option<Handle<DefaultStr>>, handle: Handle<DefaultStr>, }, Removed { did: Did<DefaultStr>, previous_handle: Option<Handle<DefaultStr>>, }, }
#[derive(Default)] struct RecordingSink { events: Mutex<Vec<SinkEvent>>, }
impl IdentitySink for RecordingSink { fn identity_changed( &self, did: &Did<DefaultStr>, previous_handle: Option<&Handle<DefaultStr>>, handle: &Handle<DefaultStr>, ) { self.events.lock().unwrap().push(SinkEvent::Changed { did: did.clone(), previous_handle: previous_handle.cloned(), handle: handle.clone(), }); }
fn identity_removed( &self, did: &Did<DefaultStr>, previous_handle: Option<&Handle<DefaultStr>>, ) { self.events.lock().unwrap().push(SinkEvent::Removed { did: did.clone(), previous_handle: previous_handle.cloned(), }); } }
struct IdentityCheckingSink { resolver: Arc<OnceCell<Arc<IdentityResolver>>>, events: Mutex<Vec<SinkEvent>>, }
impl IdentitySink for IdentityCheckingSink { fn identity_changed( &self, did: &Did<DefaultStr>, previous_handle: Option<&Handle<DefaultStr>>, handle: &Handle<DefaultStr>, ) { 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<DefaultStr>, previous_handle: Option<&Handle<DefaultStr>>, ) { 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, }, ] ); }}