This repository has no description
1use std::sync::Arc;
2use std::sync::atomic::{AtomicU64, Ordering};
3
4use bobbin_runtime::RuntimeHasher;
5use bobbin_slingshot_client::{SlingshotClient, SlingshotError};
6use jacquard_common::DefaultStr;
7use jacquard_common::types::did::Did;
8use jacquard_common::types::ident::AtIdentifier;
9use jacquard_common::types::string::Handle;
10use scc::HashMap as SccMap;
11use scc::hash_map::Entry as MapEntry;
12use serde::{Deserialize, Serialize};
13use thiserror::Error;
14use tokio::sync::OnceCell;
15
16#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
17#[serde(rename_all = "camelCase")]
18pub struct MiniDoc {
19 pub did: Did<DefaultStr>,
20 pub handle: Handle<DefaultStr>,
21 #[serde(skip_serializing_if = "Option::is_none")]
22 pub pds: Option<String>,
23}
24
25#[derive(Clone, Debug, Error, Eq, PartialEq)]
26pub enum IdentityResolveError {
27 #[error("identity not found")]
28 NotFound,
29 #[error("identity upstream: {0}")]
30 Upstream(String),
31 #[error("invalid identity response: {0}")]
32 Decode(String),
33}
34
35impl From<SlingshotError> for IdentityResolveError {
36 fn from(error: SlingshotError) -> Self {
37 match error {
38 SlingshotError::NotFound => Self::NotFound,
39 other => Self::Upstream(other.to_string()),
40 }
41 }
42}
43
44// hydrant doesnt send pds so we have a separate states to have
45// public resolve call upgrade from a partial-cached to a full-cached
46#[derive(Clone)]
47enum IdentityState {
48 // seen by bobbin from hydrant
49 Observed(MiniDoc),
50 // fetched by bobbin from slingshot
51 Fetched(MiniDoc),
52 Inactive,
53}
54
55impl IdentityState {
56 fn doc(&self) -> Option<&MiniDoc> {
57 match self {
58 Self::Observed(doc) | Self::Fetched(doc) => Some(doc),
59 Self::Inactive => None,
60 }
61 }
62}
63
64#[derive(Clone, Copy, Debug, Eq, PartialEq)]
65pub struct IdentityResolverStatsSnapshot {
66 pub entries: usize,
67 pub hits: u64,
68 pub misses: u64,
69 pub upstream_requests: u64,
70}
71
72#[derive(Default)]
73struct IdentityResolverStats {
74 hits: AtomicU64,
75 misses: AtomicU64,
76 upstream_requests: AtomicU64,
77}
78
79// in the future this would spill to disk probably
80pub struct IdentityResolver {
81 by_did: SccMap<Did<DefaultStr>, IdentityState, RuntimeHasher>,
82 by_handle: SccMap<Handle<DefaultStr>, Did<DefaultStr>, RuntimeHasher>,
83 in_flight: SccMap<String, Arc<OnceCell<Result<MiniDoc, IdentityResolveError>>>, RuntimeHasher>,
84 slingshot: Option<SlingshotClient>,
85 stats: IdentityResolverStats,
86}
87
88impl IdentityResolver {
89 pub fn with_slingshot(slingshot: SlingshotClient, hasher: RuntimeHasher) -> Self {
90 Self::new(Some(slingshot), hasher)
91 }
92
93 pub fn detached(hasher: RuntimeHasher) -> Self {
94 Self::new(None, hasher)
95 }
96
97 fn new(slingshot: Option<SlingshotClient>, hasher: RuntimeHasher) -> Self {
98 Self {
99 by_did: SccMap::with_hasher(hasher.clone()),
100 by_handle: SccMap::with_hasher(hasher.clone()),
101 in_flight: SccMap::with_hasher(hasher),
102 slingshot,
103 stats: IdentityResolverStats::default(),
104 }
105 }
106
107 pub fn stats(&self) -> IdentityResolverStatsSnapshot {
108 IdentityResolverStatsSnapshot {
109 entries: self.by_did.len(),
110 hits: self.stats.hits.load(Ordering::Relaxed),
111 misses: self.stats.misses.load(Ordering::Relaxed),
112 upstream_requests: self.stats.upstream_requests.load(Ordering::Relaxed),
113 }
114 }
115
116 pub fn observe(&self, did: Did<DefaultStr>, handle: Handle<DefaultStr>) {
117 let mut previous_handle = None;
118 match self.by_did.entry_sync(did.clone()) {
119 MapEntry::Occupied(mut occupied) => {
120 let (pds, fetched) = match occupied.get() {
121 IdentityState::Observed(previous) => {
122 previous_handle = Some(previous.handle.clone());
123 (previous.pds.clone(), false)
124 }
125 IdentityState::Fetched(previous) => {
126 previous_handle = Some(previous.handle.clone());
127 (previous.pds.clone(), previous.handle == handle)
128 }
129 IdentityState::Inactive => (None, false),
130 };
131 let doc = MiniDoc {
132 did: did.clone(),
133 handle: handle.clone(),
134 pds,
135 };
136 occupied.insert(if fetched {
137 IdentityState::Fetched(doc)
138 } else {
139 IdentityState::Observed(doc)
140 });
141 }
142 MapEntry::Vacant(vacant) => {
143 vacant.insert_entry(IdentityState::Observed(MiniDoc {
144 did: did.clone(),
145 handle: handle.clone(),
146 pds: None,
147 }));
148 }
149 }
150 self.remove_by_handle_if_owned(&did, previous_handle.as_ref());
151 self.insert_by_handle(did, handle);
152 }
153
154 pub fn deactivate(&self, did: Did<DefaultStr>) {
155 let mut previous_handle = None;
156 match self.by_did.entry_sync(did.clone()) {
157 MapEntry::Occupied(mut occupied) => {
158 previous_handle = occupied.get().doc().map(|doc| doc.handle.clone());
159 occupied.insert(IdentityState::Inactive);
160 }
161 MapEntry::Vacant(vacant) => {
162 vacant.insert_entry(IdentityState::Inactive);
163 }
164 }
165 self.remove_by_handle_if_owned(&did, previous_handle.as_ref());
166 }
167
168 fn remove_by_handle_if_owned(
169 &self,
170 did: &Did<DefaultStr>,
171 handle: Option<&Handle<DefaultStr>>,
172 ) {
173 let Some(handle) = handle else {
174 return;
175 };
176 self.by_handle.remove_if_sync(handle, |owner| owner == did);
177 }
178
179 fn remove_by_handle_for_removed_did(&self, removed: Option<(Did<DefaultStr>, IdentityState)>) {
180 let Some((did, state)) = removed else {
181 return;
182 };
183 let Some(doc) = state.doc() else {
184 return;
185 };
186 self.remove_by_handle_if_owned(&did, Some(&doc.handle));
187 }
188
189 fn insert_by_handle(&self, did: Did<DefaultStr>, handle: Handle<DefaultStr>) {
190 if let Some(displaced_did) = self
191 .by_handle
192 .upsert_sync(handle.clone(), did.clone())
193 .filter(|displaced_did| displaced_did != &did)
194 {
195 let removed = self.by_did.remove_if_sync(&displaced_did, |state| {
196 state.doc().is_some_and(|doc| doc.handle == handle)
197 });
198 self.remove_by_handle_for_removed_did(removed);
199 }
200
201 if self.by_did_matches_handle(&did, &handle) {
202 return;
203 }
204 self.remove_by_handle_if_owned(&did, Some(&handle));
205 }
206
207 fn by_did_matches_handle(&self, did: &Did<DefaultStr>, handle: &Handle<DefaultStr>) -> bool {
208 self.by_did
209 .get_sync(did)
210 .is_some_and(|state| state.get().doc().is_some_and(|doc| doc.handle == *handle))
211 }
212
213 fn cached(
214 &self,
215 identifier: &AtIdentifier<DefaultStr>,
216 require_fetched: bool,
217 ) -> Result<Option<MiniDoc>, IdentityResolveError> {
218 match identifier {
219 AtIdentifier::Did(did) => match self.by_did.get_sync(did).as_deref() {
220 Some(IdentityState::Observed(doc)) => Ok((!require_fetched).then(|| doc.clone())),
221 Some(IdentityState::Fetched(doc)) => Ok(Some(doc.clone())),
222 Some(IdentityState::Inactive) => Err(IdentityResolveError::NotFound),
223 None => Ok(None),
224 },
225 AtIdentifier::Handle(handle) => {
226 let Some(did) = self
227 .by_handle
228 .get_sync(handle)
229 .map(|entry| entry.get().clone())
230 else {
231 return Ok(None);
232 };
233 let state = self.by_did.get_sync(&did).map(|state| state.get().clone());
234 match state {
235 Some(IdentityState::Observed(doc)) if doc.handle == *handle => {
236 return Ok((!require_fetched).then_some(doc));
237 }
238 Some(IdentityState::Fetched(doc)) if doc.handle == *handle => {
239 return Ok(Some(doc));
240 }
241 Some(IdentityState::Observed(_))
242 | Some(IdentityState::Fetched(_))
243 | Some(IdentityState::Inactive)
244 | None => {}
245 }
246 self.remove_by_handle_if_owned(&did, Some(handle));
247 Ok(None)
248 }
249 }
250 }
251
252 fn get_cached(
253 &self,
254 identifier: &AtIdentifier<DefaultStr>,
255 require_fetched: bool,
256 ) -> Result<Option<MiniDoc>, IdentityResolveError> {
257 let result = self.cached(identifier, require_fetched);
258 match &result {
259 Ok(Some(_)) | Err(_) => self.stats.hits.fetch_add(1, Ordering::Relaxed),
260 Ok(None) => self.stats.misses.fetch_add(1, Ordering::Relaxed),
261 };
262 result
263 }
264
265 /// Get a Hydrant-observed DID without waiting on Slingshot.
266 pub fn get_by_did(&self, did: &Did<DefaultStr>) -> Result<MiniDoc, IdentityResolveError> {
267 self.get_cached(&AtIdentifier::Did(did.clone()), false)?
268 .ok_or(IdentityResolveError::NotFound)
269 }
270
271 /// Resolve a minidoc, fetching Hydrant-only observations upstream first.
272 pub async fn resolve_minidoc(
273 &self,
274 identifier: &AtIdentifier<DefaultStr>,
275 ) -> Result<MiniDoc, IdentityResolveError> {
276 self.resolve_with_cache(identifier, true).await
277 }
278
279 async fn resolve_with_cache(
280 &self,
281 identifier: &AtIdentifier<DefaultStr>,
282 require_fetched: bool,
283 ) -> Result<MiniDoc, IdentityResolveError> {
284 if let Some(doc) = self.get_cached(identifier, require_fetched)? {
285 return Ok(doc);
286 }
287
288 let key = identifier.as_str().to_owned();
289 let cell = self
290 .in_flight
291 .entry_async(key.clone())
292 .await
293 .or_insert_with(|| Arc::new(OnceCell::new()))
294 .get()
295 .clone();
296 let result = cell
297 .get_or_init(|| async { self.fetch_minidoc(identifier).await })
298 .await
299 .clone();
300 self.in_flight.remove_async(&key).await;
301 result
302 }
303
304 async fn fetch_minidoc(
305 &self,
306 identifier: &AtIdentifier<DefaultStr>,
307 ) -> Result<MiniDoc, IdentityResolveError> {
308 let client = self
309 .slingshot
310 .as_ref()
311 .ok_or(IdentityResolveError::NotFound)?;
312 self.stats.upstream_requests.fetch_add(1, Ordering::Relaxed);
313 let bytes = client
314 .resolve_mini_doc(identifier)
315 .await
316 .map_err(IdentityResolveError::from)?;
317 let doc = serde_json::from_slice::<MiniDoc>(&bytes)
318 .map_err(|error| IdentityResolveError::Decode(error.to_string()))?;
319 self.insert_fetched_by_did(doc)
320 }
321
322 fn insert_fetched_by_did(&self, doc: MiniDoc) -> Result<MiniDoc, IdentityResolveError> {
323 let did = doc.did.clone();
324 let handle = doc.handle.clone();
325 let mut previous_handle = None;
326 let stored = match self.by_did.entry_sync(did.clone()) {
327 MapEntry::Occupied(mut occupied) => match occupied.get_mut() {
328 IdentityState::Inactive => return Err(IdentityResolveError::NotFound),
329 IdentityState::Observed(observed) if observed.handle != handle => {
330 observed.pds = doc.pds;
331 return Ok(observed.clone());
332 }
333 IdentityState::Observed(previous) | IdentityState::Fetched(previous) => {
334 previous_handle = Some(previous.handle.clone());
335 occupied.insert(IdentityState::Fetched(doc.clone()));
336 doc
337 }
338 },
339 MapEntry::Vacant(vacant) => {
340 vacant.insert_entry(IdentityState::Fetched(doc.clone()));
341 doc
342 }
343 };
344 self.remove_by_handle_if_owned(&did, previous_handle.as_ref());
345 self.insert_by_handle(did, handle);
346 Ok(stored)
347 }
348}
349
350#[cfg(test)]
351mod tests {
352 use super::*;
353 use url::Url;
354 use wiremock::matchers::{method, path, query_param};
355 use wiremock::{Mock, MockServer, ResponseTemplate};
356
357 fn hasher() -> RuntimeHasher {
358 RuntimeHasher::from_seeds(1, 2, 3, 4)
359 }
360
361 fn resolver() -> IdentityResolver {
362 IdentityResolver::detached(hasher())
363 }
364
365 fn did(value: &str) -> Did<DefaultStr> {
366 Did::new_owned(value).unwrap()
367 }
368
369 fn handle(value: &str) -> Handle<DefaultStr> {
370 Handle::new_owned(value).unwrap()
371 }
372
373 #[test]
374 fn observed_identity_resolves_by_did_without_upstream() {
375 let resolver = resolver();
376 resolver.observe(did("did:plc:dawn"), handle("ptr.pet"));
377
378 let doc = resolver.get_by_did(&did("did:plc:dawn")).unwrap();
379 assert_eq!(doc.handle, handle("ptr.pet"));
380 assert_eq!(doc.pds, None);
381 assert_eq!(resolver.stats().hits, 1);
382 }
383
384 #[test]
385 fn observed_identity_updates_existing_did_and_removes_old_handle() {
386 let resolver = resolver();
387 let identity = did("did:plc:dawn");
388 resolver.observe(identity.clone(), handle("ptr.pet"));
389 resolver.observe(identity.clone(), handle("new.ptr.pet"));
390
391 assert!(
392 resolver
393 .cached(&AtIdentifier::Handle(handle("ptr.pet")), false)
394 .unwrap()
395 .is_none()
396 );
397 let updated = resolver.get_by_did(&identity).unwrap();
398 assert_eq!(updated.handle, handle("new.ptr.pet"));
399 }
400
401 #[test]
402 fn handle_reassignment_keeps_the_new_owner() {
403 let resolver = resolver();
404 let first = did("did:plc:first");
405 let second = did("did:plc:second");
406 let shared = handle("shared.example.com");
407 resolver.observe(first.clone(), shared.clone());
408 resolver.observe(second.clone(), shared.clone());
409 resolver.observe(first, handle("first.example.com"));
410
411 let cached = resolver
412 .cached(&AtIdentifier::Handle(shared), false)
413 .unwrap()
414 .unwrap();
415 assert_eq!(cached.did, second);
416 }
417
418 #[test]
419 fn handle_lookup_discards_an_unvalidated_reverse_hint() {
420 let resolver = resolver();
421 let identity = did("did:plc:dawn");
422 let stale = handle("stale.example.com");
423 resolver.observe(identity.clone(), handle("current.example.com"));
424 resolver
425 .by_handle
426 .upsert_sync(stale.clone(), identity.clone());
427
428 assert!(
429 resolver
430 .cached(&AtIdentifier::Handle(stale.clone()), false)
431 .unwrap()
432 .is_none()
433 );
434 assert!(resolver.by_handle.get_sync(&stale).is_none());
435 }
436
437 #[test]
438 fn minidoc_lookup_keeps_a_valid_observed_handle_hint() {
439 let resolver = resolver();
440 let identity = did("did:plc:dawn");
441 let handle = handle("ptr.pet");
442 resolver.observe(identity.clone(), handle.clone());
443
444 assert!(
445 resolver
446 .cached(&AtIdentifier::Handle(handle.clone()), true)
447 .unwrap()
448 .is_none()
449 );
450 assert_eq!(
451 resolver
452 .by_handle
453 .get_sync(&handle)
454 .map(|entry| entry.get().clone()),
455 Some(identity)
456 );
457 }
458
459 #[tokio::test]
460 async fn partial_observation_fetches_and_preserves_pds() {
461 let server = MockServer::start().await;
462 Mock::given(method("GET"))
463 .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc"))
464 .and(query_param("identifier", "did:plc:dawn"))
465 .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
466 "did": "did:plc:dawn",
467 "handle": "ptr.pet",
468 "pds": "https://pds.example.com"
469 })))
470 .expect(1)
471 .mount(&server)
472 .await;
473 let client =
474 SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap();
475 let resolver = IdentityResolver::with_slingshot(client, hasher());
476 let identity = did("did:plc:dawn");
477 resolver.observe(identity.clone(), handle("ptr.pet"));
478
479 let doc = resolver
480 .resolve_minidoc(&AtIdentifier::Did(identity.clone()))
481 .await
482 .unwrap();
483 assert_eq!(doc.pds.as_deref(), Some("https://pds.example.com"));
484
485 resolver.observe(identity.clone(), handle("ptr.pet"));
486 let cached = resolver
487 .resolve_minidoc(&AtIdentifier::Did(identity))
488 .await
489 .unwrap();
490 assert_eq!(cached.pds.as_deref(), Some("https://pds.example.com"));
491 assert_eq!(resolver.stats().upstream_requests, 1);
492 }
493
494 #[test]
495 fn inactive_identity_rejects_an_in_flight_result() {
496 let resolver = resolver();
497 let identity = did("did:plc:dawn");
498 resolver.observe(identity.clone(), handle("ptr.pet"));
499 resolver.deactivate(identity.clone());
500
501 assert_eq!(
502 resolver.get_by_did(&identity),
503 Err(IdentityResolveError::NotFound)
504 );
505 assert_eq!(
506 resolver.insert_fetched_by_did(MiniDoc {
507 did: identity,
508 handle: handle("ptr.pet"),
509 pds: Some("https://pds.example.com".to_owned()),
510 }),
511 Err(IdentityResolveError::NotFound)
512 );
513 assert!(
514 resolver
515 .cached(&AtIdentifier::Handle(handle("ptr.pet")), false)
516 .unwrap()
517 .is_none()
518 );
519 }
520}