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