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 / lib.rs
43 kB 1207 lines
1mod identity; 2mod legacy_upgrade; 3mod normalize; 4 5pub use identity::{ 6 IdentityResolveError, IdentityResolver, IdentityResolverStatsSnapshot, MiniDoc, 7}; 8pub use legacy_upgrade::{ 9 DecodedRecord, decode_canon_or_upgrade, decode_canon_or_upgrade_bytes, normalize_record_fields, 10 scrub_record_bytes, synthesize_created_at, upgrade, upgrade_wire_bytes, 11}; 12pub use normalize::NormalizeRepoRefs; 13 14use std::sync::Arc; 15use std::sync::atomic::{AtomicU64, Ordering}; 16use std::time::Duration; 17 18use tokio::time::Instant; 19 20use bobbin_runtime::{Clock, RuntimeHasher}; 21use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; 22use bobbin_types::edges::{ExtractError, Record}; 23use bobbin_types::ids::{RepoIdent, nsid_static}; 24use jacquard_common::DefaultStr; 25use jacquard_common::types::did::Did; 26use jacquard_common::types::nsid::Nsid; 27use jacquard_common::types::recordkey::Rkey; 28use scc::{HashMap as SccMap, HashSet as SccSet}; 29use tokio::sync::OnceCell; 30use tracing::warn; 31 32const REPO_COLLECTION: &str = "sh.tangled.repo"; 33const TRANSIENT_TTL: Duration = Duration::from_secs(60); 34 35#[derive(Clone, Debug, Eq, PartialEq)] 36pub enum Resolution { 37 Mapped(Did<DefaultStr>), 38 NoRepoDid, 39 Unresolvable, 40} 41 42#[derive(Clone, Debug, Eq, PartialEq)] 43enum AuthoritativeResolution { 44 Mapped(Did<DefaultStr>), 45 NoRepoDid, 46} 47 48impl AuthoritativeResolution { 49 fn from_repo_did(repo_did: Option<Did<DefaultStr>>) -> Self { 50 match repo_did { 51 Some(did) => Self::Mapped(did), 52 None => Self::NoRepoDid, 53 } 54 } 55} 56 57#[derive(Clone, Debug, Eq, PartialEq)] 58enum CacheEntry { 59 Authoritative(AuthoritativeResolution), 60 Provisional(Resolution), 61 Transient { expires_at: Instant }, 62} 63 64impl CacheEntry { 65 fn into_resolution(self) -> Resolution { 66 match self { 67 Self::Authoritative(AuthoritativeResolution::Mapped(did)) => Resolution::Mapped(did), 68 Self::Authoritative(AuthoritativeResolution::NoRepoDid) => Resolution::NoRepoDid, 69 Self::Provisional(r) => r, 70 Self::Transient { .. } => Resolution::Unresolvable, 71 } 72 } 73 74 fn is_expired_transient(&self, now: Instant) -> bool { 75 matches!(self, Self::Transient { expires_at } if *expires_at <= now) 76 } 77} 78 79#[derive(Default)] 80pub struct ResolverStats { 81 hits: AtomicU64, 82 misses_mapped: AtomicU64, 83 misses_no_repo_did: AtomicU64, 84 misses_unresolvable: AtomicU64, 85 misses_transient: AtomicU64, 86 misses_no_client: AtomicU64, 87 miss_latency_micros_sum: AtomicU64, 88 miss_latency_micros_max: AtomicU64, 89} 90 91#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] 92pub struct ResolverStatsSnapshot { 93 pub hits: u64, 94 pub misses_mapped: u64, 95 pub misses_no_repo_did: u64, 96 pub misses_unresolvable: u64, 97 pub misses_transient: u64, 98 pub misses_no_client: u64, 99 pub miss_latency_micros_sum: u64, 100 pub miss_latency_micros_max: u64, 101} 102 103impl ResolverStatsSnapshot { 104 pub fn miss_count(&self) -> u64 { 105 self.misses_mapped 106 + self.misses_no_repo_did 107 + self.misses_unresolvable 108 + self.misses_transient 109 + self.misses_no_client 110 } 111 112 pub fn total(&self) -> u64 { 113 self.hits + self.miss_count() 114 } 115 116 pub fn miss_latency_micros_avg(&self) -> Option<u64> { 117 let misses = self.miss_count() - self.misses_no_client; 118 (misses > 0).then(|| self.miss_latency_micros_sum / misses) 119 } 120} 121 122#[derive(Clone, Copy)] 123enum MissKind { 124 Mapped, 125 NoRepoDid, 126 Unresolvable, 127 Transient, 128 NoClient, 129} 130 131impl ResolverStats { 132 fn record_hit(&self) { 133 self.hits.fetch_add(1, Ordering::Relaxed); 134 } 135 136 fn record_miss(&self, kind: MissKind, latency: Option<Duration>) { 137 let counter = match kind { 138 MissKind::Mapped => &self.misses_mapped, 139 MissKind::NoRepoDid => &self.misses_no_repo_did, 140 MissKind::Unresolvable => &self.misses_unresolvable, 141 MissKind::Transient => &self.misses_transient, 142 MissKind::NoClient => &self.misses_no_client, 143 }; 144 counter.fetch_add(1, Ordering::Relaxed); 145 if let Some(latency) = latency { 146 let micros = u64::try_from(latency.as_micros()).unwrap_or(u64::MAX); 147 self.miss_latency_micros_sum 148 .fetch_add(micros, Ordering::Relaxed); 149 self.miss_latency_micros_max 150 .fetch_max(micros, Ordering::Relaxed); 151 } 152 } 153 154 pub fn snapshot(&self) -> ResolverStatsSnapshot { 155 ResolverStatsSnapshot { 156 hits: self.hits.load(Ordering::Relaxed), 157 misses_mapped: self.misses_mapped.load(Ordering::Relaxed), 158 misses_no_repo_did: self.misses_no_repo_did.load(Ordering::Relaxed), 159 misses_unresolvable: self.misses_unresolvable.load(Ordering::Relaxed), 160 misses_transient: self.misses_transient.load(Ordering::Relaxed), 161 misses_no_client: self.misses_no_client.load(Ordering::Relaxed), 162 miss_latency_micros_sum: self.miss_latency_micros_sum.load(Ordering::Relaxed), 163 miss_latency_micros_max: self.miss_latency_micros_max.load(Ordering::Relaxed), 164 } 165 } 166} 167 168struct SlingshotProbe { 169 client: SlingshotClient, 170 clock: Arc<dyn Clock>, 171} 172 173pub struct RepoIdResolver { 174 cache: SccMap<RepoIdent, CacheEntry, RuntimeHasher>, 175 by_repo_did: SccMap<Did<DefaultStr>, RepoIdent, RuntimeHasher>, 176 by_rkey: SccSet<RepoIdent, RuntimeHasher>, 177 by_name: SccMap<(Did<DefaultStr>, DefaultStr), RepoIdent, RuntimeHasher>, 178 in_flight: SccMap<RepoIdent, Arc<OnceCell<Resolution>>, RuntimeHasher>, 179 probe: Option<SlingshotProbe>, 180 stats: ResolverStats, 181} 182 183impl RepoIdResolver { 184 pub fn with_slingshot( 185 client: SlingshotClient, 186 clock: Arc<dyn Clock>, 187 hasher: RuntimeHasher, 188 ) -> Self { 189 Self { 190 cache: SccMap::with_hasher(hasher.clone()), 191 by_repo_did: SccMap::with_hasher(hasher.clone()), 192 by_rkey: SccSet::with_hasher(hasher.clone()), 193 by_name: SccMap::with_hasher(hasher.clone()), 194 in_flight: SccMap::with_hasher(hasher), 195 probe: Some(SlingshotProbe { client, clock }), 196 stats: ResolverStats::default(), 197 } 198 } 199 200 pub fn detached(hasher: RuntimeHasher) -> Self { 201 Self { 202 cache: SccMap::with_hasher(hasher.clone()), 203 by_repo_did: SccMap::with_hasher(hasher.clone()), 204 by_rkey: SccSet::with_hasher(hasher.clone()), 205 by_name: SccMap::with_hasher(hasher.clone()), 206 in_flight: SccMap::with_hasher(hasher), 207 probe: None, 208 stats: ResolverStats::default(), 209 } 210 } 211 212 pub fn stats(&self) -> ResolverStatsSnapshot { 213 self.stats.snapshot() 214 } 215 216 pub async fn cached_resolution( 217 &self, 218 owner: &Did<DefaultStr>, 219 rkey: &Rkey<DefaultStr>, 220 ) -> Option<Resolution> { 221 let key = RepoIdent::new(owner.clone(), rkey.clone()); 222 let entry = self.cache.get_async(&key).await?; 223 let now = self.probe.as_ref().map(|p| p.clock.now_instant()); 224 if let Some(now) = now 225 && entry.get().is_expired_transient(now) 226 { 227 return None; 228 } 229 Some(entry.get().clone().into_resolution()) 230 } 231 232 pub async fn lookup_by_repo_did(&self, repo_did: &Did<DefaultStr>) -> Option<RepoIdent> { 233 self.by_repo_did 234 .get_async(repo_did) 235 .await 236 .map(|e| e.get().clone()) 237 } 238 239 pub async fn lookup_by_name(&self, owner: &Did<DefaultStr>, name: &str) -> Option<RepoIdent> { 240 // rkey addressing wins so a record can't hijack another repo's rkey via its name 241 if let Ok(rkey) = Rkey::new_owned(name) { 242 let ident = RepoIdent::new(owner.clone(), rkey); 243 if self.by_rkey.contains_async(&ident).await { 244 return Some(ident); 245 } 246 } 247 self.by_name 248 .get_async(&(owner.clone(), DefaultStr::from(name))) 249 .await 250 .map(|e| e.get().clone()) 251 } 252 253 pub async fn observe_rkey(&self, owner: Did<DefaultStr>, rkey: Rkey<DefaultStr>) { 254 let _ = self.by_rkey.insert_async(RepoIdent::new(owner, rkey)).await; 255 } 256 257 pub async fn observe( 258 &self, 259 owner: Did<DefaultStr>, 260 rkey: Rkey<DefaultStr>, 261 repo_did: Option<Did<DefaultStr>>, 262 name: Option<DefaultStr>, 263 ) -> Option<RepoIdent> { 264 let ident = RepoIdent::new(owner, rkey); 265 self.observe_rkey(ident.owner.clone(), ident.rkey.clone()) 266 .await; 267 if let Some(name) = name { 268 let _ = self 269 .by_name 270 .insert_async((ident.owner.clone(), name), ident.clone()) 271 .await; 272 } 273 let entry = 274 CacheEntry::Authoritative(AuthoritativeResolution::from_repo_did(repo_did.clone())); 275 self.cache 276 .entry_async(ident.clone()) 277 .await 278 .and_modify(|existing| *existing = entry.clone()) 279 .or_insert(entry); 280 281 let repo_did = repo_did?; 282 let mut prior: Option<RepoIdent> = None; 283 self.by_repo_did 284 .entry_async(repo_did) 285 .await 286 .and_modify(|existing| { 287 if *existing != ident { 288 prior = Some(existing.clone()); 289 *existing = ident.clone(); 290 } 291 }) 292 .or_insert(ident); 293 prior 294 } 295 296 pub async fn forget(&self, owner: &Did<DefaultStr>, rkey: &Rkey<DefaultStr>) { 297 let ident = RepoIdent::new(owner.clone(), rkey.clone()); 298 let prior_resolution = self 299 .cache 300 .remove_async(&ident) 301 .await 302 .map(|(_, entry)| entry.into_resolution()); 303 if let Some(Resolution::Mapped(repo_did)) = prior_resolution { 304 self.by_repo_did 305 .remove_if_async(&repo_did, |existing| *existing == ident) 306 .await; 307 } 308 self.by_rkey.remove_async(&ident).await; 309 } 310 311 async fn fill_provisional(&self, key: RepoIdent, resolution: Resolution) { 312 let entry = CacheEntry::Provisional(resolution); 313 self.cache 314 .entry_async(key) 315 .await 316 .and_modify(|existing| { 317 if matches!(existing, CacheEntry::Authoritative(_)) { 318 return; 319 } 320 *existing = entry.clone(); 321 }) 322 .or_insert(entry); 323 } 324 325 async fn fill_transient(&self, key: RepoIdent, expires_at: Instant) { 326 let entry = CacheEntry::Transient { expires_at }; 327 self.cache 328 .entry_async(key) 329 .await 330 .and_modify(|existing| { 331 if matches!(existing, CacheEntry::Authoritative(_)) { 332 return; 333 } 334 *existing = entry.clone(); 335 }) 336 .or_insert(entry); 337 } 338 339 pub async fn resolve(&self, owner: &Did<DefaultStr>, rkey: &Rkey<DefaultStr>) -> Resolution { 340 let key = RepoIdent::new(owner.clone(), rkey.clone()); 341 342 let Some(probe) = self.probe.as_ref() else { 343 if let Some(entry) = self.cache.get_async(&key).await { 344 self.stats.record_hit(); 345 return entry.get().clone().into_resolution(); 346 } 347 self.stats.record_miss(MissKind::NoClient, None); 348 return Resolution::Unresolvable; 349 }; 350 351 let now = probe.clock.now_instant(); 352 if let Some(entry) = self.cache.get_async(&key).await 353 && !entry.get().is_expired_transient(now) 354 { 355 self.stats.record_hit(); 356 return entry.get().clone().into_resolution(); 357 } 358 359 let cell: Arc<OnceCell<Resolution>> = self 360 .in_flight 361 .entry_async(key.clone()) 362 .await 363 .or_insert_with(|| Arc::new(OnceCell::new())) 364 .get() 365 .clone(); 366 367 let result = cell 368 .get_or_init(|| async { self.fetch_repo_did(owner, rkey, &key).await }) 369 .await 370 .clone(); 371 372 self.in_flight.remove_async(&key).await; 373 374 result 375 } 376 377 async fn fetch_repo_did( 378 &self, 379 owner: &Did<DefaultStr>, 380 rkey: &Rkey<DefaultStr>, 381 key: &RepoIdent, 382 ) -> Resolution { 383 let probe = self 384 .probe 385 .as_ref() 386 .expect("fetch_repo_did is only called when a probe is present"); 387 let started = probe.clock.now_instant(); 388 let nsid: Nsid<DefaultStr> = nsid_static(REPO_COLLECTION); 389 let provisional = match probe.client.get_record(owner, &nsid, rkey).await { 390 Ok(body) => match repo_did_from_body(&nsid, &body.value) { 391 Ok(Some(did)) => Resolution::Mapped(did), 392 Ok(None) => Resolution::NoRepoDid, 393 Err(e) => { 394 warn!( 395 error = ?e, 396 owner = owner.as_ref(), 397 rkey = rkey.as_ref(), 398 "slingshot returned unparseable repo body, caching as unresolvable", 399 ); 400 Resolution::Unresolvable 401 } 402 }, 403 Err(SlingshotError::NotFound) => { 404 warn!( 405 owner = owner.as_ref(), 406 rkey = rkey.as_ref(), 407 "no repo record on slingshot, caching as unresolvable", 408 ); 409 Resolution::Unresolvable 410 } 411 Err(ref e) if is_garbage_response(e) => { 412 warn!( 413 error = ?e, 414 owner = owner.as_ref(), 415 rkey = rkey.as_ref(), 416 "slingshot returned malformed response, caching as unresolvable", 417 ); 418 Resolution::Unresolvable 419 } 420 Err(e) => { 421 warn!( 422 error = ?e, 423 owner = owner.as_ref(), 424 rkey = rkey.as_ref(), 425 "caching transient slingshot failure for repoDID lookup under short TTL", 426 ); 427 let elapsed = probe.clock.now_instant().duration_since(started); 428 self.stats.record_miss(MissKind::Transient, Some(elapsed)); 429 let expires_at = probe.clock.now_instant() + TRANSIENT_TTL; 430 self.fill_transient(key.clone(), expires_at).await; 431 return Resolution::Unresolvable; 432 } 433 }; 434 let elapsed = probe.clock.now_instant().duration_since(started); 435 let kind = match &provisional { 436 Resolution::Mapped(_) => MissKind::Mapped, 437 Resolution::NoRepoDid => MissKind::NoRepoDid, 438 Resolution::Unresolvable => MissKind::Unresolvable, 439 }; 440 self.stats.record_miss(kind, Some(elapsed)); 441 self.fill_provisional(key.clone(), provisional.clone()) 442 .await; 443 provisional 444 } 445} 446 447fn is_garbage_response(err: &SlingshotError) -> bool { 448 matches!( 449 err, 450 SlingshotError::Decode(_) 451 | SlingshotError::MissingField(_) 452 | SlingshotError::InvalidAtUri(_) 453 | SlingshotError::InvalidCid(_) 454 | SlingshotError::UriMismatch { .. }, 455 ) 456} 457 458fn repo_did_from_body( 459 nsid: &Nsid<DefaultStr>, 460 body: &[u8], 461) -> Result<Option<Did<DefaultStr>>, ExtractError> { 462 match DecodedRecord::try_decode(nsid, body)? { 463 DecodedRecord::Canon(Record::Repo(repo)) => Ok(repo.repo_did), 464 DecodedRecord::Canon(_) | DecodedRecord::Legacy(_) => Ok(None), 465 } 466} 467 468#[cfg(test)] 469mod tests { 470 use super::*; 471 use bobbin_runtime::SystemClock; 472 use jacquard_common::types::did::Did; 473 use jacquard_common::types::recordkey::Rkey; 474 475 fn did(s: &str) -> Did<DefaultStr> { 476 Did::new_owned(s).unwrap() 477 } 478 479 fn rkey(s: &str) -> Rkey<DefaultStr> { 480 Rkey::new_owned(s).unwrap() 481 } 482 483 fn test_clock() -> Arc<dyn Clock> { 484 Arc::new(SystemClock::new()) 485 } 486 487 #[tokio::test] 488 async fn observation_returns_prior_ident_when_repo_did_moves() { 489 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 490 let prior = resolver 491 .observe( 492 did("did:plc:nel"), 493 rkey("3liuighjy2h22"), 494 Some(did("did:plc:clam")), 495 None, 496 ) 497 .await; 498 assert!(prior.is_none(), "first observation has no prior"); 499 500 let prior = resolver 501 .observe( 502 did("did:plc:nel"), 503 rkey("core"), 504 Some(did("did:plc:clam")), 505 None, 506 ) 507 .await; 508 assert_eq!( 509 prior, 510 Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), 511 "same repoDID at a new (owner, rkey) returns the prior ident so callers can evict the stale at-uri", 512 ); 513 514 let prior = resolver 515 .observe( 516 did("did:plc:nel"), 517 rkey("core"), 518 Some(did("did:plc:clam")), 519 None, 520 ) 521 .await; 522 assert!(prior.is_none(), "re-observing the same ident is a no-op"); 523 } 524 525 #[tokio::test] 526 async fn observation_without_repo_did_does_not_track_reverse() { 527 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 528 let prior = resolver 529 .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, None) 530 .await; 531 assert!(prior.is_none()); 532 } 533 534 #[tokio::test] 535 async fn forget_clears_reverse_only_when_still_owned() { 536 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 537 resolver 538 .observe( 539 did("did:plc:nel"), 540 rkey("3liuighjy2h22"), 541 Some(did("did:plc:clam")), 542 None, 543 ) 544 .await; 545 resolver 546 .observe( 547 did("did:plc:nel"), 548 rkey("core"), 549 Some(did("did:plc:clam")), 550 None, 551 ) 552 .await; 553 554 resolver 555 .forget(&did("did:plc:nel"), &rkey("3liuighjy2h22")) 556 .await; 557 558 let prior = resolver 559 .observe( 560 did("did:plc:nel"), 561 rkey("core-renamed"), 562 Some(did("did:plc:clam")), 563 None, 564 ) 565 .await; 566 assert_eq!( 567 prior, 568 Some(RepoIdent::new(did("did:plc:nel"), rkey("core"))), 569 "stale at-uri's forget must not displace the live owner of did:plc:clam", 570 ); 571 } 572 573 #[tokio::test] 574 async fn lookup_by_name_finds_observed_rkey() { 575 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 576 resolver 577 .observe_rkey(did("did:plc:nel"), rkey("3liuighjy2h22")) 578 .await; 579 let got = resolver 580 .lookup_by_name(&did("did:plc:nel"), "3liuighjy2h22") 581 .await; 582 assert_eq!( 583 got, 584 Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), 585 ); 586 } 587 588 #[tokio::test] 589 async fn lookup_by_name_is_scoped_to_the_owner() { 590 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 591 resolver 592 .observe_rkey(did("did:plc:nel"), rkey("3liuighjy2h22")) 593 .await; 594 let got = resolver 595 .lookup_by_name(&did("did:plc:olaren"), "3liuighjy2h22") 596 .await; 597 assert_eq!(got, None, "one owner's rkey must not answer for another's"); 598 } 599 600 #[tokio::test] 601 async fn lookup_by_name_rejects_non_rkey() { 602 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 603 let owner = did("did:plc:nel"); 604 resolver 605 .observe_rkey(owner.clone(), rkey("3liuighjy2h22")) 606 .await; 607 608 assert_eq!(resolver.lookup_by_name(&owner, "my repo").await, None); 609 } 610 611 #[tokio::test] 612 async fn lookup_by_name_finds_observed_record_name() { 613 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 614 resolver 615 .observe( 616 did("did:plc:nel"), 617 rkey("3liuighjy2h22"), 618 Some(did("did:plc:clam")), 619 Some(DefaultStr::from("ark")), 620 ) 621 .await; 622 let got = resolver.lookup_by_name(&did("did:plc:nel"), "ark").await; 623 assert_eq!( 624 got, 625 Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), 626 ); 627 } 628 629 #[tokio::test] 630 async fn lookup_by_name_record_name_is_scoped_to_the_owner() { 631 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 632 resolver 633 .observe( 634 did("did:plc:nel"), 635 rkey("3liuighjy2h22"), 636 None, 637 Some(DefaultStr::from("ark")), 638 ) 639 .await; 640 let got = resolver.lookup_by_name(&did("did:plc:olaren"), "ark").await; 641 assert_eq!( 642 got, None, 643 "one owner's repo name must not answer for another's" 644 ); 645 } 646 647 #[tokio::test] 648 async fn lookup_by_name_prefers_rkey_over_claimed_name() { 649 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 650 resolver 651 .observe(did("did:plc:nel"), rkey("core"), None, None) 652 .await; 653 resolver 654 .observe( 655 did("did:plc:nel"), 656 rkey("3liuighjy2h22"), 657 None, 658 Some(DefaultStr::from("core")), 659 ) 660 .await; 661 let got = resolver.lookup_by_name(&did("did:plc:nel"), "core").await; 662 assert_eq!( 663 got, 664 Some(RepoIdent::new(did("did:plc:nel"), rkey("core"))), 665 "a record naming itself after a live rkey must not shadow it", 666 ); 667 } 668 669 #[tokio::test] 670 async fn forget_clears_the_rkey() { 671 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 672 let owner = did("did:plc:nel"); 673 resolver 674 .observe_rkey(owner.clone(), rkey("3liuighjy2h22")) 675 .await; 676 resolver.forget(&owner, &rkey("3liuighjy2h22")).await; 677 assert_eq!(resolver.lookup_by_name(&owner, "3liuighjy2h22").await, None); 678 } 679 680 #[tokio::test] 681 async fn observation_with_repo_did_resolves_mapped() { 682 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 683 resolver 684 .observe( 685 did("did:plc:nel"), 686 rkey("abcabcabcabcz"), 687 Some(did("did:plc:clam")), 688 None, 689 ) 690 .await; 691 let got = resolver 692 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 693 .await; 694 assert_eq!(got, Resolution::Mapped(did("did:plc:clam"))); 695 } 696 697 #[tokio::test] 698 async fn observation_without_repo_did_resolves_no_repo_did() { 699 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 700 resolver 701 .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, None) 702 .await; 703 let got = resolver 704 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 705 .await; 706 assert_eq!( 707 got, 708 Resolution::NoRepoDid, 709 "observed but empty repoDID is a definitive answer not a lookup failure", 710 ); 711 } 712 713 #[tokio::test] 714 async fn cache_miss_without_client_is_unresolvable() { 715 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 716 let got = resolver 717 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 718 .await; 719 assert_eq!(got, Resolution::Unresolvable); 720 } 721 722 #[tokio::test] 723 async fn lookup_by_repo_did_finds_observed_ident() { 724 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 725 resolver 726 .observe( 727 did("did:plc:nel"), 728 rkey("abcabcabcabcz"), 729 Some(did("did:plc:limpet")), 730 None, 731 ) 732 .await; 733 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 734 assert_eq!( 735 got, 736 Some(RepoIdent::new(did("did:plc:nel"), rkey("abcabcabcabcz"))), 737 ); 738 } 739 740 #[tokio::test] 741 async fn lookup_by_repo_did_misses_when_unobserved() { 742 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 743 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 744 assert_eq!(got, None); 745 } 746 747 #[tokio::test] 748 async fn lookup_by_repo_did_misses_when_repo_did_was_none() { 749 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 750 resolver 751 .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, None) 752 .await; 753 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 754 assert_eq!(got, None); 755 } 756 757 #[tokio::test] 758 async fn lookup_by_repo_did_follows_move_to_new_ident() { 759 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 760 resolver 761 .observe( 762 did("did:plc:nel"), 763 rkey("abcabcabcabcz"), 764 Some(did("did:plc:limpet")), 765 None, 766 ) 767 .await; 768 resolver 769 .observe( 770 did("did:plc:olaren"), 771 rkey("xyzxyzxyzxyzx"), 772 Some(did("did:plc:limpet")), 773 None, 774 ) 775 .await; 776 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 777 assert_eq!( 778 got, 779 Some(RepoIdent::new(did("did:plc:olaren"), rkey("xyzxyzxyzxyzx"))), 780 ); 781 } 782 783 #[tokio::test] 784 async fn observation_overwrites_prior_value() { 785 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 786 resolver 787 .observe( 788 did("did:plc:nel"), 789 rkey("abcabcabcabcz"), 790 Some(did("did:plc:clam")), 791 None, 792 ) 793 .await; 794 resolver 795 .observe( 796 did("did:plc:nel"), 797 rkey("abcabcabcabcz"), 798 Some(did("did:plc:uni")), 799 None, 800 ) 801 .await; 802 let got = resolver 803 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 804 .await; 805 assert_eq!(got, Resolution::Mapped(did("did:plc:uni"))); 806 } 807 808 #[tokio::test] 809 async fn fill_provisional_does_not_downgrade_authoritative_mapped() { 810 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 811 let owner = did("did:plc:nel"); 812 let key = rkey("abcabcabcabcz"); 813 resolver 814 .observe(owner.clone(), key.clone(), Some(did("did:plc:clam")), None) 815 .await; 816 resolver 817 .fill_provisional( 818 RepoIdent::new(owner.clone(), key.clone()), 819 Resolution::Unresolvable, 820 ) 821 .await; 822 let got = resolver.resolve(&owner, &key).await; 823 assert_eq!( 824 got, 825 Resolution::Mapped(did("did:plc:clam")), 826 "firehose-observed mapping must outrank provisional slingshot info", 827 ); 828 } 829 830 #[tokio::test] 831 async fn fill_provisional_does_not_downgrade_authoritative_no_repo_did() { 832 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 833 let owner = did("did:plc:nel"); 834 let key = rkey("abcabcabcabcz"); 835 resolver 836 .observe(owner.clone(), key.clone(), None, None) 837 .await; 838 resolver 839 .fill_provisional( 840 RepoIdent::new(owner.clone(), key.clone()), 841 Resolution::Mapped(did("did:plc:clam")), 842 ) 843 .await; 844 let got = resolver.resolve(&owner, &key).await; 845 assert_eq!( 846 got, 847 Resolution::NoRepoDid, 848 "an authoritative empty observation must outrank provisional slingshot info even when slingshot disagrees", 849 ); 850 } 851 852 #[tokio::test] 853 async fn slingshot_404_caches_as_unresolvable() { 854 let server = wiremock::MockServer::start().await; 855 wiremock::Mock::given(wiremock::matchers::method("GET")) 856 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 857 .respond_with(wiremock::ResponseTemplate::new(404)) 858 .expect(1) 859 .mount(&server) 860 .await; 861 862 let client = 863 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 864 let resolver = 865 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 866 867 let owner = did("did:plc:nel"); 868 let key = rkey("abcabcabcabcz"); 869 let first = resolver.resolve(&owner, &key).await; 870 let second = resolver.resolve(&owner, &key).await; 871 assert_eq!(first, Resolution::Unresolvable); 872 assert_eq!(second, Resolution::Unresolvable); 873 } 874 875 #[tokio::test] 876 async fn slingshot_malformed_envelope_caches_as_unresolvable() { 877 let server = wiremock::MockServer::start().await; 878 wiremock::Mock::given(wiremock::matchers::method("GET")) 879 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 880 .respond_with( 881 wiremock::ResponseTemplate::new(200) 882 .insert_header("content-type", "application/json") 883 .set_body_string("not json"), 884 ) 885 .expect(1) 886 .mount(&server) 887 .await; 888 889 let client = 890 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 891 let resolver = 892 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 893 894 let owner = did("did:plc:nel"); 895 let key = rkey("abcabcabcabcz"); 896 let first = resolver.resolve(&owner, &key).await; 897 let second = resolver.resolve(&owner, &key).await; 898 assert_eq!(first, Resolution::Unresolvable); 899 assert_eq!( 900 second, 901 Resolution::Unresolvable, 902 "garbage envelopes are stable across retries, so caching avoids hammering slingshot", 903 ); 904 } 905 906 #[tokio::test] 907 async fn slingshot_uri_mismatch_caches_as_unresolvable() { 908 let server = wiremock::MockServer::start().await; 909 let body = serde_json::json!({ 910 "uri": "at://did:plc:limpet/sh.tangled.repo/elsewhere", 911 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 912 "value": {"$type": "sh.tangled.repo", "knot": "oyster.cafe", "createdAt": "2026-05-01T00:00:00Z"} 913 }); 914 wiremock::Mock::given(wiremock::matchers::method("GET")) 915 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 916 .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) 917 .expect(1) 918 .mount(&server) 919 .await; 920 921 let client = 922 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 923 let resolver = 924 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 925 926 let owner = did("did:plc:nel"); 927 let key = rkey("abcabcabcabcz"); 928 let first = resolver.resolve(&owner, &key).await; 929 let second = resolver.resolve(&owner, &key).await; 930 assert_eq!(first, Resolution::Unresolvable); 931 assert_eq!(second, Resolution::Unresolvable); 932 } 933 934 #[tokio::test] 935 async fn slingshot_legacy_repo_body_resolves_no_repo_did_not_unresolvable() { 936 let server = wiremock::MockServer::start().await; 937 let body = serde_json::json!({ 938 "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", 939 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 940 "value": { 941 "$type": "sh.tangled.repo", 942 "addedAt": "2025-03-07T21:47:53Z", 943 "knot": "knot1.tangled.sh", 944 "name": "scallop", 945 "owner": "did:plc:nel", 946 }, 947 }); 948 wiremock::Mock::given(wiremock::matchers::method("GET")) 949 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 950 .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) 951 .expect(1) 952 .mount(&server) 953 .await; 954 955 let client = 956 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 957 let resolver = 958 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 959 960 let owner = did("did:plc:nel"); 961 let key = rkey("abcabcabcabcz"); 962 let got = resolver.resolve(&owner, &key).await; 963 assert_eq!( 964 got, 965 Resolution::NoRepoDid, 966 "legacy repo wires without a repo_did parse via legacy upgrade and resolve as NoRepoDid, not Unresolvable", 967 ); 968 } 969 970 #[tokio::test] 971 async fn slingshot_unparseable_repo_value_caches_as_unresolvable() { 972 let server = wiremock::MockServer::start().await; 973 let body = serde_json::json!({ 974 "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", 975 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 976 "value": {"$type": "sh.tangled.repo"} 977 }); 978 wiremock::Mock::given(wiremock::matchers::method("GET")) 979 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 980 .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) 981 .expect(1) 982 .mount(&server) 983 .await; 984 985 let client = 986 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 987 let resolver = 988 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 989 990 let owner = did("did:plc:nel"); 991 let key = rkey("abcabcabcabcz"); 992 let first = resolver.resolve(&owner, &key).await; 993 let second = resolver.resolve(&owner, &key).await; 994 assert_eq!( 995 first, 996 Resolution::Unresolvable, 997 "a repo body that fails lexicon validation must not be conflated with NoRepoDid", 998 ); 999 assert_eq!(second, Resolution::Unresolvable); 1000 } 1001 1002 #[tokio::test] 1003 async fn slingshot_transport_error_caches_with_short_ttl() { 1004 let server = wiremock::MockServer::start().await; 1005 wiremock::Mock::given(wiremock::matchers::method("GET")) 1006 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 1007 .respond_with(wiremock::ResponseTemplate::new(503)) 1008 .expect(1) 1009 .mount(&server) 1010 .await; 1011 1012 let client = 1013 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 1014 let resolver = 1015 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 1016 1017 let owner = did("did:plc:nel"); 1018 let key = rkey("abcabcabcabcz"); 1019 let first = resolver.resolve(&owner, &key).await; 1020 let second = resolver.resolve(&owner, &key).await; 1021 assert_eq!(first, Resolution::Unresolvable); 1022 assert_eq!( 1023 second, 1024 Resolution::Unresolvable, 1025 "transient TTL must suppress immediate re-hammering of a sick upstream", 1026 ); 1027 } 1028 1029 #[tokio::test] 1030 async fn slingshot_transient_recorded_separately_from_unresolvable() { 1031 let server = wiremock::MockServer::start().await; 1032 wiremock::Mock::given(wiremock::matchers::method("GET")) 1033 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 1034 .respond_with(wiremock::ResponseTemplate::new(503)) 1035 .expect(1) 1036 .mount(&server) 1037 .await; 1038 1039 let client = 1040 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 1041 let resolver = 1042 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 1043 1044 let owner = did("did:plc:nel"); 1045 let key = rkey("abcabcabcabcz"); 1046 resolver.resolve(&owner, &key).await; 1047 resolver.resolve(&owner, &key).await; 1048 1049 let snap = resolver.stats(); 1050 assert_eq!( 1051 snap.misses_transient, 1, 1052 "second resolve must hit the short-TTL cache instead of re-firing the transient miss", 1053 ); 1054 assert_eq!(snap.hits, 1, "second call hits cached transient entry"); 1055 assert_eq!( 1056 snap.misses_unresolvable, 0, 1057 "canonical unresolvable counter is reserved for cached terminal answers", 1058 ); 1059 assert!( 1060 snap.miss_latency_micros_sum > 0, 1061 "transient misses still have latency contributions", 1062 ); 1063 assert_eq!(snap.miss_count(), 1); 1064 } 1065 1066 #[tokio::test] 1067 async fn slingshot_in_flight_requests_coalesce() { 1068 let server = wiremock::MockServer::start().await; 1069 let body = serde_json::json!({ 1070 "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", 1071 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 1072 "value": {"$type": "sh.tangled.repo", "knot": "oyster.cafe", "createdAt": "2026-05-01T00:00:00Z", "repoDid": "did:plc:limpet"} 1073 }); 1074 wiremock::Mock::given(wiremock::matchers::method("GET")) 1075 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 1076 .respond_with( 1077 wiremock::ResponseTemplate::new(200) 1078 .set_body_json(body) 1079 .set_delay(Duration::from_millis(200)), 1080 ) 1081 .expect(1) 1082 .mount(&server) 1083 .await; 1084 1085 let client = 1086 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 1087 let resolver = Arc::new(RepoIdResolver::with_slingshot( 1088 client, 1089 test_clock(), 1090 RuntimeHasher::default(), 1091 )); 1092 1093 let owner = did("did:plc:nel"); 1094 let key = rkey("abcabcabcabcz"); 1095 let r0 = resolver.clone(); 1096 let r1 = resolver.clone(); 1097 let r2 = resolver.clone(); 1098 let o0 = owner.clone(); 1099 let o1 = owner.clone(); 1100 let o2 = owner.clone(); 1101 let k0 = key.clone(); 1102 let k1 = key.clone(); 1103 let k2 = key.clone(); 1104 let (a, b, c) = tokio::join!( 1105 tokio::spawn(async move { r0.resolve(&o0, &k0).await }), 1106 tokio::spawn(async move { r1.resolve(&o1, &k1).await }), 1107 tokio::spawn(async move { r2.resolve(&o2, &k2).await }), 1108 ); 1109 let expected = Resolution::Mapped(did("did:plc:limpet")); 1110 assert_eq!(a.unwrap(), expected); 1111 assert_eq!(b.unwrap(), expected); 1112 assert_eq!(c.unwrap(), expected); 1113 1114 let snap = resolver.stats(); 1115 assert_eq!( 1116 snap.misses_mapped, 1, 1117 "only the winning task pays the slingshot RTT", 1118 ); 1119 } 1120 1121 #[tokio::test] 1122 async fn stats_count_hits_misses_and_latency() { 1123 let server = wiremock::MockServer::start().await; 1124 wiremock::Mock::given(wiremock::matchers::method("GET")) 1125 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 1126 .respond_with(wiremock::ResponseTemplate::new(404)) 1127 .mount(&server) 1128 .await; 1129 let client = 1130 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 1131 let resolver = 1132 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 1133 1134 let owner = did("did:plc:nel"); 1135 let key = rkey("abcabcabcabcz"); 1136 resolver.resolve(&owner, &key).await; 1137 resolver.resolve(&owner, &key).await; 1138 1139 let snap = resolver.stats(); 1140 assert_eq!( 1141 snap.misses_unresolvable, 1, 1142 "first call is the slingshot miss" 1143 ); 1144 assert_eq!(snap.hits, 1, "second call hits the unresolvable cache"); 1145 assert_eq!(snap.miss_count(), 1); 1146 assert_eq!(snap.total(), 2); 1147 assert!( 1148 snap.miss_latency_micros_sum > 0, 1149 "latency recorded for slingshot miss" 1150 ); 1151 assert!(snap.miss_latency_micros_avg().unwrap() > 0); 1152 } 1153 1154 #[tokio::test] 1155 async fn stats_no_client_miss_recorded_without_latency() { 1156 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 1157 resolver 1158 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 1159 .await; 1160 let snap = resolver.stats(); 1161 assert_eq!(snap.misses_no_client, 1); 1162 assert_eq!(snap.miss_latency_micros_sum, 0); 1163 assert_eq!(snap.miss_latency_micros_avg(), None); 1164 } 1165 1166 #[tokio::test] 1167 async fn firehose_observe_can_demote_provisional() { 1168 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 1169 let owner = did("did:plc:nel"); 1170 let key = rkey("abcabcabcabcz"); 1171 resolver 1172 .fill_provisional( 1173 RepoIdent::new(owner.clone(), key.clone()), 1174 Resolution::Mapped(did("did:plc:clam")), 1175 ) 1176 .await; 1177 resolver 1178 .observe(owner.clone(), key.clone(), None, None) 1179 .await; 1180 let got = resolver.resolve(&owner, &key).await; 1181 assert_eq!( 1182 got, 1183 Resolution::NoRepoDid, 1184 "firehose update is canonical and may legitimately remove repoDID", 1185 ); 1186 } 1187 1188 #[tokio::test] 1189 async fn forget_removes_cache_entry() { 1190 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 1191 let owner = did("did:plc:nel"); 1192 let key = rkey("abcabcabcabcz"); 1193 resolver 1194 .observe(owner.clone(), key.clone(), Some(did("did:plc:clam")), None) 1195 .await; 1196 assert_eq!( 1197 resolver.cached_resolution(&owner, &key).await, 1198 Some(Resolution::Mapped(did("did:plc:clam"))), 1199 ); 1200 resolver.forget(&owner, &key).await; 1201 assert_eq!( 1202 resolver.cached_resolution(&owner, &key).await, 1203 None, 1204 "forget must drop the entry entirely so a subsequent observe can supply fresh state", 1205 ); 1206 } 1207}