This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / bobbin / crates / resolver / src / identity.rs
18 kB 520 lines
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}