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