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
39 kB 1099 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 in_flight: SccMap<RepoIdent, Arc<OnceCell<Resolution>>, RuntimeHasher>, 174 probe: Option<SlingshotProbe>, 175 stats: ResolverStats, 176} 177 178impl RepoIdResolver { 179 pub fn with_slingshot( 180 client: SlingshotClient, 181 clock: Arc<dyn Clock>, 182 hasher: RuntimeHasher, 183 ) -> Self { 184 Self { 185 cache: SccMap::with_hasher(hasher.clone()), 186 by_repo_did: SccMap::with_hasher(hasher.clone()), 187 by_rkey: SccSet::with_hasher(hasher.clone()), 188 in_flight: SccMap::with_hasher(hasher), 189 probe: Some(SlingshotProbe { client, clock }), 190 stats: ResolverStats::default(), 191 } 192 } 193 194 pub fn detached(hasher: RuntimeHasher) -> Self { 195 Self { 196 cache: SccMap::with_hasher(hasher.clone()), 197 by_repo_did: SccMap::with_hasher(hasher.clone()), 198 by_rkey: SccSet::with_hasher(hasher.clone()), 199 in_flight: SccMap::with_hasher(hasher), 200 probe: None, 201 stats: ResolverStats::default(), 202 } 203 } 204 205 pub fn stats(&self) -> ResolverStatsSnapshot { 206 self.stats.snapshot() 207 } 208 209 pub async fn cached_resolution( 210 &self, 211 owner: &Did<DefaultStr>, 212 rkey: &Rkey<DefaultStr>, 213 ) -> Option<Resolution> { 214 let key = RepoIdent::new(owner.clone(), rkey.clone()); 215 let entry = self.cache.get_async(&key).await?; 216 let now = self.probe.as_ref().map(|p| p.clock.now_instant()); 217 if let Some(now) = now 218 && entry.get().is_expired_transient(now) 219 { 220 return None; 221 } 222 Some(entry.get().clone().into_resolution()) 223 } 224 225 pub async fn lookup_by_repo_did(&self, repo_did: &Did<DefaultStr>) -> Option<RepoIdent> { 226 self.by_repo_did 227 .get_async(repo_did) 228 .await 229 .map(|e| e.get().clone()) 230 } 231 232 pub async fn lookup_by_name(&self, owner: &Did<DefaultStr>, name: &str) -> Option<RepoIdent> { 233 let rkey = Rkey::new_owned(name).ok()?; 234 let ident = RepoIdent::new(owner.clone(), rkey); 235 self.by_rkey.contains_async(&ident).await.then_some(ident) 236 } 237 238 pub async fn observe_rkey(&self, owner: Did<DefaultStr>, rkey: Rkey<DefaultStr>) { 239 let _ = self.by_rkey.insert_async(RepoIdent::new(owner, rkey)).await; 240 } 241 242 pub async fn observe( 243 &self, 244 owner: Did<DefaultStr>, 245 rkey: Rkey<DefaultStr>, 246 repo_did: Option<Did<DefaultStr>>, 247 ) -> Option<RepoIdent> { 248 let ident = RepoIdent::new(owner, rkey); 249 self.observe_rkey(ident.owner.clone(), ident.rkey.clone()) 250 .await; 251 let entry = 252 CacheEntry::Authoritative(AuthoritativeResolution::from_repo_did(repo_did.clone())); 253 self.cache 254 .entry_async(ident.clone()) 255 .await 256 .and_modify(|existing| *existing = entry.clone()) 257 .or_insert(entry); 258 259 let repo_did = repo_did?; 260 let mut prior: Option<RepoIdent> = None; 261 self.by_repo_did 262 .entry_async(repo_did) 263 .await 264 .and_modify(|existing| { 265 if *existing != ident { 266 prior = Some(existing.clone()); 267 *existing = ident.clone(); 268 } 269 }) 270 .or_insert(ident); 271 prior 272 } 273 274 pub async fn forget(&self, owner: &Did<DefaultStr>, rkey: &Rkey<DefaultStr>) { 275 let ident = RepoIdent::new(owner.clone(), rkey.clone()); 276 let prior_resolution = self 277 .cache 278 .remove_async(&ident) 279 .await 280 .map(|(_, entry)| entry.into_resolution()); 281 if let Some(Resolution::Mapped(repo_did)) = prior_resolution { 282 self.by_repo_did 283 .remove_if_async(&repo_did, |existing| *existing == ident) 284 .await; 285 } 286 self.by_rkey.remove_async(&ident).await; 287 } 288 289 async fn fill_provisional(&self, key: RepoIdent, resolution: Resolution) { 290 let entry = CacheEntry::Provisional(resolution); 291 self.cache 292 .entry_async(key) 293 .await 294 .and_modify(|existing| { 295 if matches!(existing, CacheEntry::Authoritative(_)) { 296 return; 297 } 298 *existing = entry.clone(); 299 }) 300 .or_insert(entry); 301 } 302 303 async fn fill_transient(&self, key: RepoIdent, expires_at: Instant) { 304 let entry = CacheEntry::Transient { expires_at }; 305 self.cache 306 .entry_async(key) 307 .await 308 .and_modify(|existing| { 309 if matches!(existing, CacheEntry::Authoritative(_)) { 310 return; 311 } 312 *existing = entry.clone(); 313 }) 314 .or_insert(entry); 315 } 316 317 pub async fn resolve(&self, owner: &Did<DefaultStr>, rkey: &Rkey<DefaultStr>) -> Resolution { 318 let key = RepoIdent::new(owner.clone(), rkey.clone()); 319 320 let Some(probe) = self.probe.as_ref() else { 321 if let Some(entry) = self.cache.get_async(&key).await { 322 self.stats.record_hit(); 323 return entry.get().clone().into_resolution(); 324 } 325 self.stats.record_miss(MissKind::NoClient, None); 326 return Resolution::Unresolvable; 327 }; 328 329 let now = probe.clock.now_instant(); 330 if let Some(entry) = self.cache.get_async(&key).await 331 && !entry.get().is_expired_transient(now) 332 { 333 self.stats.record_hit(); 334 return entry.get().clone().into_resolution(); 335 } 336 337 let cell: Arc<OnceCell<Resolution>> = self 338 .in_flight 339 .entry_async(key.clone()) 340 .await 341 .or_insert_with(|| Arc::new(OnceCell::new())) 342 .get() 343 .clone(); 344 345 let result = cell 346 .get_or_init(|| async { self.fetch_repo_did(owner, rkey, &key).await }) 347 .await 348 .clone(); 349 350 self.in_flight.remove_async(&key).await; 351 352 result 353 } 354 355 async fn fetch_repo_did( 356 &self, 357 owner: &Did<DefaultStr>, 358 rkey: &Rkey<DefaultStr>, 359 key: &RepoIdent, 360 ) -> Resolution { 361 let probe = self 362 .probe 363 .as_ref() 364 .expect("fetch_repo_did is only called when a probe is present"); 365 let started = probe.clock.now_instant(); 366 let nsid: Nsid<DefaultStr> = nsid_static(REPO_COLLECTION); 367 let provisional = match probe.client.get_record(owner, &nsid, rkey).await { 368 Ok(body) => match repo_did_from_body(&nsid, &body.value) { 369 Ok(Some(did)) => Resolution::Mapped(did), 370 Ok(None) => Resolution::NoRepoDid, 371 Err(e) => { 372 warn!( 373 error = ?e, 374 owner = owner.as_ref(), 375 rkey = rkey.as_ref(), 376 "slingshot returned unparseable repo body, caching as unresolvable", 377 ); 378 Resolution::Unresolvable 379 } 380 }, 381 Err(SlingshotError::NotFound) => { 382 warn!( 383 owner = owner.as_ref(), 384 rkey = rkey.as_ref(), 385 "no repo record on slingshot, caching as unresolvable", 386 ); 387 Resolution::Unresolvable 388 } 389 Err(ref e) if is_garbage_response(e) => { 390 warn!( 391 error = ?e, 392 owner = owner.as_ref(), 393 rkey = rkey.as_ref(), 394 "slingshot returned malformed response, caching as unresolvable", 395 ); 396 Resolution::Unresolvable 397 } 398 Err(e) => { 399 warn!( 400 error = ?e, 401 owner = owner.as_ref(), 402 rkey = rkey.as_ref(), 403 "caching transient slingshot failure for repoDID lookup under short TTL", 404 ); 405 let elapsed = probe.clock.now_instant().duration_since(started); 406 self.stats.record_miss(MissKind::Transient, Some(elapsed)); 407 let expires_at = probe.clock.now_instant() + TRANSIENT_TTL; 408 self.fill_transient(key.clone(), expires_at).await; 409 return Resolution::Unresolvable; 410 } 411 }; 412 let elapsed = probe.clock.now_instant().duration_since(started); 413 let kind = match &provisional { 414 Resolution::Mapped(_) => MissKind::Mapped, 415 Resolution::NoRepoDid => MissKind::NoRepoDid, 416 Resolution::Unresolvable => MissKind::Unresolvable, 417 }; 418 self.stats.record_miss(kind, Some(elapsed)); 419 self.fill_provisional(key.clone(), provisional.clone()) 420 .await; 421 provisional 422 } 423} 424 425fn is_garbage_response(err: &SlingshotError) -> bool { 426 matches!( 427 err, 428 SlingshotError::Decode(_) 429 | SlingshotError::MissingField(_) 430 | SlingshotError::InvalidAtUri(_) 431 | SlingshotError::InvalidCid(_) 432 | SlingshotError::UriMismatch { .. }, 433 ) 434} 435 436fn repo_did_from_body( 437 nsid: &Nsid<DefaultStr>, 438 body: &[u8], 439) -> Result<Option<Did<DefaultStr>>, ExtractError> { 440 match DecodedRecord::try_decode(nsid, body)? { 441 DecodedRecord::Canon(Record::Repo(repo)) => Ok(repo.repo_did), 442 DecodedRecord::Canon(_) | DecodedRecord::Legacy(_) => Ok(None), 443 } 444} 445 446#[cfg(test)] 447mod tests { 448 use super::*; 449 use bobbin_runtime::SystemClock; 450 use jacquard_common::types::did::Did; 451 use jacquard_common::types::recordkey::Rkey; 452 453 fn did(s: &str) -> Did<DefaultStr> { 454 Did::new_owned(s).unwrap() 455 } 456 457 fn rkey(s: &str) -> Rkey<DefaultStr> { 458 Rkey::new_owned(s).unwrap() 459 } 460 461 fn test_clock() -> Arc<dyn Clock> { 462 Arc::new(SystemClock::new()) 463 } 464 465 #[tokio::test] 466 async fn observation_returns_prior_ident_when_repo_did_moves() { 467 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 468 let prior = resolver 469 .observe( 470 did("did:plc:nel"), 471 rkey("3liuighjy2h22"), 472 Some(did("did:plc:clam")), 473 ) 474 .await; 475 assert!(prior.is_none(), "first observation has no prior"); 476 477 let prior = resolver 478 .observe(did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam"))) 479 .await; 480 assert_eq!( 481 prior, 482 Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), 483 "same repoDID at a new (owner, rkey) returns the prior ident so callers can evict the stale at-uri", 484 ); 485 486 let prior = resolver 487 .observe(did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam"))) 488 .await; 489 assert!(prior.is_none(), "re-observing the same ident is a no-op"); 490 } 491 492 #[tokio::test] 493 async fn observation_without_repo_did_does_not_track_reverse() { 494 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 495 let prior = resolver 496 .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None) 497 .await; 498 assert!(prior.is_none()); 499 } 500 501 #[tokio::test] 502 async fn forget_clears_reverse_only_when_still_owned() { 503 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 504 resolver 505 .observe( 506 did("did:plc:nel"), 507 rkey("3liuighjy2h22"), 508 Some(did("did:plc:clam")), 509 ) 510 .await; 511 resolver 512 .observe(did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam"))) 513 .await; 514 515 resolver 516 .forget(&did("did:plc:nel"), &rkey("3liuighjy2h22")) 517 .await; 518 519 let prior = resolver 520 .observe( 521 did("did:plc:nel"), 522 rkey("core-renamed"), 523 Some(did("did:plc:clam")), 524 ) 525 .await; 526 assert_eq!( 527 prior, 528 Some(RepoIdent::new(did("did:plc:nel"), rkey("core"))), 529 "stale at-uri's forget must not displace the live owner of did:plc:clam", 530 ); 531 } 532 533 #[tokio::test] 534 async fn lookup_by_name_finds_observed_rkey() { 535 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 536 resolver 537 .observe_rkey(did("did:plc:nel"), rkey("3liuighjy2h22")) 538 .await; 539 let got = resolver 540 .lookup_by_name(&did("did:plc:nel"), "3liuighjy2h22") 541 .await; 542 assert_eq!( 543 got, 544 Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), 545 ); 546 } 547 548 #[tokio::test] 549 async fn lookup_by_name_is_scoped_to_the_owner() { 550 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 551 resolver 552 .observe_rkey(did("did:plc:nel"), rkey("3liuighjy2h22")) 553 .await; 554 let got = resolver 555 .lookup_by_name(&did("did:plc:olaren"), "3liuighjy2h22") 556 .await; 557 assert_eq!(got, None, "one owner's rkey must not answer for another's"); 558 } 559 560 #[tokio::test] 561 async fn lookup_by_name_rejects_non_rkey() { 562 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 563 let owner = did("did:plc:nel"); 564 resolver 565 .observe_rkey(owner.clone(), rkey("3liuighjy2h22")) 566 .await; 567 568 assert_eq!(resolver.lookup_by_name(&owner, "my repo").await, None); 569 } 570 571 #[tokio::test] 572 async fn forget_clears_the_rkey() { 573 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 574 let owner = did("did:plc:nel"); 575 resolver 576 .observe_rkey(owner.clone(), rkey("3liuighjy2h22")) 577 .await; 578 resolver.forget(&owner, &rkey("3liuighjy2h22")).await; 579 assert_eq!(resolver.lookup_by_name(&owner, "3liuighjy2h22").await, None); 580 } 581 582 #[tokio::test] 583 async fn observation_with_repo_did_resolves_mapped() { 584 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 585 resolver 586 .observe( 587 did("did:plc:nel"), 588 rkey("abcabcabcabcz"), 589 Some(did("did:plc:clam")), 590 ) 591 .await; 592 let got = resolver 593 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 594 .await; 595 assert_eq!(got, Resolution::Mapped(did("did:plc:clam"))); 596 } 597 598 #[tokio::test] 599 async fn observation_without_repo_did_resolves_no_repo_did() { 600 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 601 resolver 602 .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None) 603 .await; 604 let got = resolver 605 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 606 .await; 607 assert_eq!( 608 got, 609 Resolution::NoRepoDid, 610 "observed but empty repoDID is a definitive answer not a lookup failure", 611 ); 612 } 613 614 #[tokio::test] 615 async fn cache_miss_without_client_is_unresolvable() { 616 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 617 let got = resolver 618 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 619 .await; 620 assert_eq!(got, Resolution::Unresolvable); 621 } 622 623 #[tokio::test] 624 async fn lookup_by_repo_did_finds_observed_ident() { 625 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 626 resolver 627 .observe( 628 did("did:plc:nel"), 629 rkey("abcabcabcabcz"), 630 Some(did("did:plc:limpet")), 631 ) 632 .await; 633 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 634 assert_eq!( 635 got, 636 Some(RepoIdent::new(did("did:plc:nel"), rkey("abcabcabcabcz"))), 637 ); 638 } 639 640 #[tokio::test] 641 async fn lookup_by_repo_did_misses_when_unobserved() { 642 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 643 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 644 assert_eq!(got, None); 645 } 646 647 #[tokio::test] 648 async fn lookup_by_repo_did_misses_when_repo_did_was_none() { 649 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 650 resolver 651 .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None) 652 .await; 653 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 654 assert_eq!(got, None); 655 } 656 657 #[tokio::test] 658 async fn lookup_by_repo_did_follows_move_to_new_ident() { 659 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 660 resolver 661 .observe( 662 did("did:plc:nel"), 663 rkey("abcabcabcabcz"), 664 Some(did("did:plc:limpet")), 665 ) 666 .await; 667 resolver 668 .observe( 669 did("did:plc:olaren"), 670 rkey("xyzxyzxyzxyzx"), 671 Some(did("did:plc:limpet")), 672 ) 673 .await; 674 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 675 assert_eq!( 676 got, 677 Some(RepoIdent::new(did("did:plc:olaren"), rkey("xyzxyzxyzxyzx"))), 678 ); 679 } 680 681 #[tokio::test] 682 async fn observation_overwrites_prior_value() { 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 ) 690 .await; 691 resolver 692 .observe( 693 did("did:plc:nel"), 694 rkey("abcabcabcabcz"), 695 Some(did("did:plc:uni")), 696 ) 697 .await; 698 let got = resolver 699 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 700 .await; 701 assert_eq!(got, Resolution::Mapped(did("did:plc:uni"))); 702 } 703 704 #[tokio::test] 705 async fn fill_provisional_does_not_downgrade_authoritative_mapped() { 706 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 707 let owner = did("did:plc:nel"); 708 let key = rkey("abcabcabcabcz"); 709 resolver 710 .observe(owner.clone(), key.clone(), Some(did("did:plc:clam"))) 711 .await; 712 resolver 713 .fill_provisional( 714 RepoIdent::new(owner.clone(), key.clone()), 715 Resolution::Unresolvable, 716 ) 717 .await; 718 let got = resolver.resolve(&owner, &key).await; 719 assert_eq!( 720 got, 721 Resolution::Mapped(did("did:plc:clam")), 722 "firehose-observed mapping must outrank provisional slingshot info", 723 ); 724 } 725 726 #[tokio::test] 727 async fn fill_provisional_does_not_downgrade_authoritative_no_repo_did() { 728 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 729 let owner = did("did:plc:nel"); 730 let key = rkey("abcabcabcabcz"); 731 resolver.observe(owner.clone(), key.clone(), None).await; 732 resolver 733 .fill_provisional( 734 RepoIdent::new(owner.clone(), key.clone()), 735 Resolution::Mapped(did("did:plc:clam")), 736 ) 737 .await; 738 let got = resolver.resolve(&owner, &key).await; 739 assert_eq!( 740 got, 741 Resolution::NoRepoDid, 742 "an authoritative empty observation must outrank provisional slingshot info even when slingshot disagrees", 743 ); 744 } 745 746 #[tokio::test] 747 async fn slingshot_404_caches_as_unresolvable() { 748 let server = wiremock::MockServer::start().await; 749 wiremock::Mock::given(wiremock::matchers::method("GET")) 750 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 751 .respond_with(wiremock::ResponseTemplate::new(404)) 752 .expect(1) 753 .mount(&server) 754 .await; 755 756 let client = 757 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 758 let resolver = 759 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 760 761 let owner = did("did:plc:nel"); 762 let key = rkey("abcabcabcabcz"); 763 let first = resolver.resolve(&owner, &key).await; 764 let second = resolver.resolve(&owner, &key).await; 765 assert_eq!(first, Resolution::Unresolvable); 766 assert_eq!(second, Resolution::Unresolvable); 767 } 768 769 #[tokio::test] 770 async fn slingshot_malformed_envelope_caches_as_unresolvable() { 771 let server = wiremock::MockServer::start().await; 772 wiremock::Mock::given(wiremock::matchers::method("GET")) 773 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 774 .respond_with( 775 wiremock::ResponseTemplate::new(200) 776 .insert_header("content-type", "application/json") 777 .set_body_string("not json"), 778 ) 779 .expect(1) 780 .mount(&server) 781 .await; 782 783 let client = 784 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 785 let resolver = 786 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 787 788 let owner = did("did:plc:nel"); 789 let key = rkey("abcabcabcabcz"); 790 let first = resolver.resolve(&owner, &key).await; 791 let second = resolver.resolve(&owner, &key).await; 792 assert_eq!(first, Resolution::Unresolvable); 793 assert_eq!( 794 second, 795 Resolution::Unresolvable, 796 "garbage envelopes are stable across retries, so caching avoids hammering slingshot", 797 ); 798 } 799 800 #[tokio::test] 801 async fn slingshot_uri_mismatch_caches_as_unresolvable() { 802 let server = wiremock::MockServer::start().await; 803 let body = serde_json::json!({ 804 "uri": "at://did:plc:limpet/sh.tangled.repo/elsewhere", 805 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 806 "value": {"$type": "sh.tangled.repo", "knot": "oyster.cafe", "createdAt": "2026-05-01T00:00:00Z"} 807 }); 808 wiremock::Mock::given(wiremock::matchers::method("GET")) 809 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 810 .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) 811 .expect(1) 812 .mount(&server) 813 .await; 814 815 let client = 816 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 817 let resolver = 818 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 819 820 let owner = did("did:plc:nel"); 821 let key = rkey("abcabcabcabcz"); 822 let first = resolver.resolve(&owner, &key).await; 823 let second = resolver.resolve(&owner, &key).await; 824 assert_eq!(first, Resolution::Unresolvable); 825 assert_eq!(second, Resolution::Unresolvable); 826 } 827 828 #[tokio::test] 829 async fn slingshot_legacy_repo_body_resolves_no_repo_did_not_unresolvable() { 830 let server = wiremock::MockServer::start().await; 831 let body = serde_json::json!({ 832 "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", 833 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 834 "value": { 835 "$type": "sh.tangled.repo", 836 "addedAt": "2025-03-07T21:47:53Z", 837 "knot": "knot1.tangled.sh", 838 "name": "scallop", 839 "owner": "did:plc:nel", 840 }, 841 }); 842 wiremock::Mock::given(wiremock::matchers::method("GET")) 843 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 844 .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) 845 .expect(1) 846 .mount(&server) 847 .await; 848 849 let client = 850 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 851 let resolver = 852 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 853 854 let owner = did("did:plc:nel"); 855 let key = rkey("abcabcabcabcz"); 856 let got = resolver.resolve(&owner, &key).await; 857 assert_eq!( 858 got, 859 Resolution::NoRepoDid, 860 "legacy repo wires without a repo_did parse via legacy upgrade and resolve as NoRepoDid, not Unresolvable", 861 ); 862 } 863 864 #[tokio::test] 865 async fn slingshot_unparseable_repo_value_caches_as_unresolvable() { 866 let server = wiremock::MockServer::start().await; 867 let body = serde_json::json!({ 868 "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", 869 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 870 "value": {"$type": "sh.tangled.repo"} 871 }); 872 wiremock::Mock::given(wiremock::matchers::method("GET")) 873 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 874 .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) 875 .expect(1) 876 .mount(&server) 877 .await; 878 879 let client = 880 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 881 let resolver = 882 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 883 884 let owner = did("did:plc:nel"); 885 let key = rkey("abcabcabcabcz"); 886 let first = resolver.resolve(&owner, &key).await; 887 let second = resolver.resolve(&owner, &key).await; 888 assert_eq!( 889 first, 890 Resolution::Unresolvable, 891 "a repo body that fails lexicon validation must not be conflated with NoRepoDid", 892 ); 893 assert_eq!(second, Resolution::Unresolvable); 894 } 895 896 #[tokio::test] 897 async fn slingshot_transport_error_caches_with_short_ttl() { 898 let server = wiremock::MockServer::start().await; 899 wiremock::Mock::given(wiremock::matchers::method("GET")) 900 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 901 .respond_with(wiremock::ResponseTemplate::new(503)) 902 .expect(1) 903 .mount(&server) 904 .await; 905 906 let client = 907 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 908 let resolver = 909 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 910 911 let owner = did("did:plc:nel"); 912 let key = rkey("abcabcabcabcz"); 913 let first = resolver.resolve(&owner, &key).await; 914 let second = resolver.resolve(&owner, &key).await; 915 assert_eq!(first, Resolution::Unresolvable); 916 assert_eq!( 917 second, 918 Resolution::Unresolvable, 919 "transient TTL must suppress immediate re-hammering of a sick upstream", 920 ); 921 } 922 923 #[tokio::test] 924 async fn slingshot_transient_recorded_separately_from_unresolvable() { 925 let server = wiremock::MockServer::start().await; 926 wiremock::Mock::given(wiremock::matchers::method("GET")) 927 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 928 .respond_with(wiremock::ResponseTemplate::new(503)) 929 .expect(1) 930 .mount(&server) 931 .await; 932 933 let client = 934 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 935 let resolver = 936 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 937 938 let owner = did("did:plc:nel"); 939 let key = rkey("abcabcabcabcz"); 940 resolver.resolve(&owner, &key).await; 941 resolver.resolve(&owner, &key).await; 942 943 let snap = resolver.stats(); 944 assert_eq!( 945 snap.misses_transient, 1, 946 "second resolve must hit the short-TTL cache instead of re-firing the transient miss", 947 ); 948 assert_eq!(snap.hits, 1, "second call hits cached transient entry"); 949 assert_eq!( 950 snap.misses_unresolvable, 0, 951 "canonical unresolvable counter is reserved for cached terminal answers", 952 ); 953 assert!( 954 snap.miss_latency_micros_sum > 0, 955 "transient misses still have latency contributions", 956 ); 957 assert_eq!(snap.miss_count(), 1); 958 } 959 960 #[tokio::test] 961 async fn slingshot_in_flight_requests_coalesce() { 962 let server = wiremock::MockServer::start().await; 963 let body = serde_json::json!({ 964 "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", 965 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 966 "value": {"$type": "sh.tangled.repo", "knot": "oyster.cafe", "createdAt": "2026-05-01T00:00:00Z", "repoDid": "did:plc:limpet"} 967 }); 968 wiremock::Mock::given(wiremock::matchers::method("GET")) 969 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 970 .respond_with( 971 wiremock::ResponseTemplate::new(200) 972 .set_body_json(body) 973 .set_delay(Duration::from_millis(200)), 974 ) 975 .expect(1) 976 .mount(&server) 977 .await; 978 979 let client = 980 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 981 let resolver = Arc::new(RepoIdResolver::with_slingshot( 982 client, 983 test_clock(), 984 RuntimeHasher::default(), 985 )); 986 987 let owner = did("did:plc:nel"); 988 let key = rkey("abcabcabcabcz"); 989 let r0 = resolver.clone(); 990 let r1 = resolver.clone(); 991 let r2 = resolver.clone(); 992 let o0 = owner.clone(); 993 let o1 = owner.clone(); 994 let o2 = owner.clone(); 995 let k0 = key.clone(); 996 let k1 = key.clone(); 997 let k2 = key.clone(); 998 let (a, b, c) = tokio::join!( 999 tokio::spawn(async move { r0.resolve(&o0, &k0).await }), 1000 tokio::spawn(async move { r1.resolve(&o1, &k1).await }), 1001 tokio::spawn(async move { r2.resolve(&o2, &k2).await }), 1002 ); 1003 let expected = Resolution::Mapped(did("did:plc:limpet")); 1004 assert_eq!(a.unwrap(), expected); 1005 assert_eq!(b.unwrap(), expected); 1006 assert_eq!(c.unwrap(), expected); 1007 1008 let snap = resolver.stats(); 1009 assert_eq!( 1010 snap.misses_mapped, 1, 1011 "only the winning task pays the slingshot RTT", 1012 ); 1013 } 1014 1015 #[tokio::test] 1016 async fn stats_count_hits_misses_and_latency() { 1017 let server = wiremock::MockServer::start().await; 1018 wiremock::Mock::given(wiremock::matchers::method("GET")) 1019 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 1020 .respond_with(wiremock::ResponseTemplate::new(404)) 1021 .mount(&server) 1022 .await; 1023 let client = 1024 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 1025 let resolver = 1026 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 1027 1028 let owner = did("did:plc:nel"); 1029 let key = rkey("abcabcabcabcz"); 1030 resolver.resolve(&owner, &key).await; 1031 resolver.resolve(&owner, &key).await; 1032 1033 let snap = resolver.stats(); 1034 assert_eq!( 1035 snap.misses_unresolvable, 1, 1036 "first call is the slingshot miss" 1037 ); 1038 assert_eq!(snap.hits, 1, "second call hits the unresolvable cache"); 1039 assert_eq!(snap.miss_count(), 1); 1040 assert_eq!(snap.total(), 2); 1041 assert!( 1042 snap.miss_latency_micros_sum > 0, 1043 "latency recorded for slingshot miss" 1044 ); 1045 assert!(snap.miss_latency_micros_avg().unwrap() > 0); 1046 } 1047 1048 #[tokio::test] 1049 async fn stats_no_client_miss_recorded_without_latency() { 1050 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 1051 resolver 1052 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 1053 .await; 1054 let snap = resolver.stats(); 1055 assert_eq!(snap.misses_no_client, 1); 1056 assert_eq!(snap.miss_latency_micros_sum, 0); 1057 assert_eq!(snap.miss_latency_micros_avg(), None); 1058 } 1059 1060 #[tokio::test] 1061 async fn firehose_observe_can_demote_provisional() { 1062 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 1063 let owner = did("did:plc:nel"); 1064 let key = rkey("abcabcabcabcz"); 1065 resolver 1066 .fill_provisional( 1067 RepoIdent::new(owner.clone(), key.clone()), 1068 Resolution::Mapped(did("did:plc:clam")), 1069 ) 1070 .await; 1071 resolver.observe(owner.clone(), key.clone(), None).await; 1072 let got = resolver.resolve(&owner, &key).await; 1073 assert_eq!( 1074 got, 1075 Resolution::NoRepoDid, 1076 "firehose update is canonical and may legitimately remove repoDID", 1077 ); 1078 } 1079 1080 #[tokio::test] 1081 async fn forget_removes_cache_entry() { 1082 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 1083 let owner = did("did:plc:nel"); 1084 let key = rkey("abcabcabcabcz"); 1085 resolver 1086 .observe(owner.clone(), key.clone(), Some(did("did:plc:clam"))) 1087 .await; 1088 assert_eq!( 1089 resolver.cached_resolution(&owner, &key).await, 1090 Some(Resolution::Mapped(did("did:plc:clam"))), 1091 ); 1092 resolver.forget(&owner, &key).await; 1093 assert_eq!( 1094 resolver.cached_resolution(&owner, &key).await, 1095 None, 1096 "forget must drop the entry entirely so a subsequent observe can supply fresh state", 1097 ); 1098 } 1099}