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
42 kB 1181 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(did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam")), None) 498 .await; 499 assert_eq!( 500 prior, 501 Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), 502 "same repoDID at a new (owner, rkey) returns the prior ident so callers can evict the stale at-uri", 503 ); 504 505 let prior = resolver 506 .observe(did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam")), None) 507 .await; 508 assert!(prior.is_none(), "re-observing the same ident is a no-op"); 509 } 510 511 #[tokio::test] 512 async fn observation_without_repo_did_does_not_track_reverse() { 513 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 514 let prior = resolver 515 .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, None) 516 .await; 517 assert!(prior.is_none()); 518 } 519 520 #[tokio::test] 521 async fn forget_clears_reverse_only_when_still_owned() { 522 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 523 resolver 524 .observe( 525 did("did:plc:nel"), 526 rkey("3liuighjy2h22"), 527 Some(did("did:plc:clam")), 528 None, 529 ) 530 .await; 531 resolver 532 .observe(did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam")), None) 533 .await; 534 535 resolver 536 .forget(&did("did:plc:nel"), &rkey("3liuighjy2h22")) 537 .await; 538 539 let prior = resolver 540 .observe( 541 did("did:plc:nel"), 542 rkey("core-renamed"), 543 Some(did("did:plc:clam")), 544 None, 545 ) 546 .await; 547 assert_eq!( 548 prior, 549 Some(RepoIdent::new(did("did:plc:nel"), rkey("core"))), 550 "stale at-uri's forget must not displace the live owner of did:plc:clam", 551 ); 552 } 553 554 #[tokio::test] 555 async fn lookup_by_name_finds_observed_rkey() { 556 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 557 resolver 558 .observe_rkey(did("did:plc:nel"), rkey("3liuighjy2h22")) 559 .await; 560 let got = resolver 561 .lookup_by_name(&did("did:plc:nel"), "3liuighjy2h22") 562 .await; 563 assert_eq!( 564 got, 565 Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), 566 ); 567 } 568 569 #[tokio::test] 570 async fn lookup_by_name_is_scoped_to_the_owner() { 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:olaren"), "3liuighjy2h22") 577 .await; 578 assert_eq!(got, None, "one owner's rkey must not answer for another's"); 579 } 580 581 #[tokio::test] 582 async fn lookup_by_name_rejects_non_rkey() { 583 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 584 let owner = did("did:plc:nel"); 585 resolver 586 .observe_rkey(owner.clone(), rkey("3liuighjy2h22")) 587 .await; 588 589 assert_eq!(resolver.lookup_by_name(&owner, "my repo").await, None); 590 } 591 592 #[tokio::test] 593 async fn lookup_by_name_finds_observed_record_name() { 594 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 595 resolver 596 .observe( 597 did("did:plc:nel"), 598 rkey("3liuighjy2h22"), 599 Some(did("did:plc:clam")), 600 Some(DefaultStr::from("ark")), 601 ) 602 .await; 603 let got = resolver.lookup_by_name(&did("did:plc:nel"), "ark").await; 604 assert_eq!( 605 got, 606 Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), 607 ); 608 } 609 610 #[tokio::test] 611 async fn lookup_by_name_record_name_is_scoped_to_the_owner() { 612 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 613 resolver 614 .observe( 615 did("did:plc:nel"), 616 rkey("3liuighjy2h22"), 617 None, 618 Some(DefaultStr::from("ark")), 619 ) 620 .await; 621 let got = resolver.lookup_by_name(&did("did:plc:olaren"), "ark").await; 622 assert_eq!(got, None, "one owner's repo name must not answer for another's"); 623 } 624 625 #[tokio::test] 626 async fn lookup_by_name_prefers_rkey_over_claimed_name() { 627 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 628 resolver 629 .observe(did("did:plc:nel"), rkey("core"), None, None) 630 .await; 631 resolver 632 .observe( 633 did("did:plc:nel"), 634 rkey("3liuighjy2h22"), 635 None, 636 Some(DefaultStr::from("core")), 637 ) 638 .await; 639 let got = resolver.lookup_by_name(&did("did:plc:nel"), "core").await; 640 assert_eq!( 641 got, 642 Some(RepoIdent::new(did("did:plc:nel"), rkey("core"))), 643 "a record naming itself after a live rkey must not shadow it", 644 ); 645 } 646 647 #[tokio::test] 648 async fn forget_clears_the_rkey() { 649 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 650 let owner = did("did:plc:nel"); 651 resolver 652 .observe_rkey(owner.clone(), rkey("3liuighjy2h22")) 653 .await; 654 resolver.forget(&owner, &rkey("3liuighjy2h22")).await; 655 assert_eq!(resolver.lookup_by_name(&owner, "3liuighjy2h22").await, None); 656 } 657 658 #[tokio::test] 659 async fn observation_with_repo_did_resolves_mapped() { 660 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 661 resolver 662 .observe( 663 did("did:plc:nel"), 664 rkey("abcabcabcabcz"), 665 Some(did("did:plc:clam")), 666 None, 667 ) 668 .await; 669 let got = resolver 670 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 671 .await; 672 assert_eq!(got, Resolution::Mapped(did("did:plc:clam"))); 673 } 674 675 #[tokio::test] 676 async fn observation_without_repo_did_resolves_no_repo_did() { 677 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 678 resolver 679 .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, None) 680 .await; 681 let got = resolver 682 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 683 .await; 684 assert_eq!( 685 got, 686 Resolution::NoRepoDid, 687 "observed but empty repoDID is a definitive answer not a lookup failure", 688 ); 689 } 690 691 #[tokio::test] 692 async fn cache_miss_without_client_is_unresolvable() { 693 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 694 let got = resolver 695 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 696 .await; 697 assert_eq!(got, Resolution::Unresolvable); 698 } 699 700 #[tokio::test] 701 async fn lookup_by_repo_did_finds_observed_ident() { 702 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 703 resolver 704 .observe( 705 did("did:plc:nel"), 706 rkey("abcabcabcabcz"), 707 Some(did("did:plc:limpet")), 708 None, 709 ) 710 .await; 711 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 712 assert_eq!( 713 got, 714 Some(RepoIdent::new(did("did:plc:nel"), rkey("abcabcabcabcz"))), 715 ); 716 } 717 718 #[tokio::test] 719 async fn lookup_by_repo_did_misses_when_unobserved() { 720 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 721 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 722 assert_eq!(got, None); 723 } 724 725 #[tokio::test] 726 async fn lookup_by_repo_did_misses_when_repo_did_was_none() { 727 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 728 resolver 729 .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, None) 730 .await; 731 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 732 assert_eq!(got, None); 733 } 734 735 #[tokio::test] 736 async fn lookup_by_repo_did_follows_move_to_new_ident() { 737 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 738 resolver 739 .observe( 740 did("did:plc:nel"), 741 rkey("abcabcabcabcz"), 742 Some(did("did:plc:limpet")), 743 None, 744 ) 745 .await; 746 resolver 747 .observe( 748 did("did:plc:olaren"), 749 rkey("xyzxyzxyzxyzx"), 750 Some(did("did:plc:limpet")), 751 None, 752 ) 753 .await; 754 let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; 755 assert_eq!( 756 got, 757 Some(RepoIdent::new(did("did:plc:olaren"), rkey("xyzxyzxyzxyzx"))), 758 ); 759 } 760 761 #[tokio::test] 762 async fn observation_overwrites_prior_value() { 763 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 764 resolver 765 .observe( 766 did("did:plc:nel"), 767 rkey("abcabcabcabcz"), 768 Some(did("did:plc:clam")), 769 None, 770 ) 771 .await; 772 resolver 773 .observe( 774 did("did:plc:nel"), 775 rkey("abcabcabcabcz"), 776 Some(did("did:plc:uni")), 777 None, 778 ) 779 .await; 780 let got = resolver 781 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 782 .await; 783 assert_eq!(got, Resolution::Mapped(did("did:plc:uni"))); 784 } 785 786 #[tokio::test] 787 async fn fill_provisional_does_not_downgrade_authoritative_mapped() { 788 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 789 let owner = did("did:plc:nel"); 790 let key = rkey("abcabcabcabcz"); 791 resolver 792 .observe(owner.clone(), key.clone(), Some(did("did:plc:clam")), None) 793 .await; 794 resolver 795 .fill_provisional( 796 RepoIdent::new(owner.clone(), key.clone()), 797 Resolution::Unresolvable, 798 ) 799 .await; 800 let got = resolver.resolve(&owner, &key).await; 801 assert_eq!( 802 got, 803 Resolution::Mapped(did("did:plc:clam")), 804 "firehose-observed mapping must outrank provisional slingshot info", 805 ); 806 } 807 808 #[tokio::test] 809 async fn fill_provisional_does_not_downgrade_authoritative_no_repo_did() { 810 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 811 let owner = did("did:plc:nel"); 812 let key = rkey("abcabcabcabcz"); 813 resolver.observe(owner.clone(), key.clone(), None, None).await; 814 resolver 815 .fill_provisional( 816 RepoIdent::new(owner.clone(), key.clone()), 817 Resolution::Mapped(did("did:plc:clam")), 818 ) 819 .await; 820 let got = resolver.resolve(&owner, &key).await; 821 assert_eq!( 822 got, 823 Resolution::NoRepoDid, 824 "an authoritative empty observation must outrank provisional slingshot info even when slingshot disagrees", 825 ); 826 } 827 828 #[tokio::test] 829 async fn slingshot_404_caches_as_unresolvable() { 830 let server = wiremock::MockServer::start().await; 831 wiremock::Mock::given(wiremock::matchers::method("GET")) 832 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 833 .respond_with(wiremock::ResponseTemplate::new(404)) 834 .expect(1) 835 .mount(&server) 836 .await; 837 838 let client = 839 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 840 let resolver = 841 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 842 843 let owner = did("did:plc:nel"); 844 let key = rkey("abcabcabcabcz"); 845 let first = resolver.resolve(&owner, &key).await; 846 let second = resolver.resolve(&owner, &key).await; 847 assert_eq!(first, Resolution::Unresolvable); 848 assert_eq!(second, Resolution::Unresolvable); 849 } 850 851 #[tokio::test] 852 async fn slingshot_malformed_envelope_caches_as_unresolvable() { 853 let server = wiremock::MockServer::start().await; 854 wiremock::Mock::given(wiremock::matchers::method("GET")) 855 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 856 .respond_with( 857 wiremock::ResponseTemplate::new(200) 858 .insert_header("content-type", "application/json") 859 .set_body_string("not json"), 860 ) 861 .expect(1) 862 .mount(&server) 863 .await; 864 865 let client = 866 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 867 let resolver = 868 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 869 870 let owner = did("did:plc:nel"); 871 let key = rkey("abcabcabcabcz"); 872 let first = resolver.resolve(&owner, &key).await; 873 let second = resolver.resolve(&owner, &key).await; 874 assert_eq!(first, Resolution::Unresolvable); 875 assert_eq!( 876 second, 877 Resolution::Unresolvable, 878 "garbage envelopes are stable across retries, so caching avoids hammering slingshot", 879 ); 880 } 881 882 #[tokio::test] 883 async fn slingshot_uri_mismatch_caches_as_unresolvable() { 884 let server = wiremock::MockServer::start().await; 885 let body = serde_json::json!({ 886 "uri": "at://did:plc:limpet/sh.tangled.repo/elsewhere", 887 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 888 "value": {"$type": "sh.tangled.repo", "knot": "oyster.cafe", "createdAt": "2026-05-01T00:00:00Z"} 889 }); 890 wiremock::Mock::given(wiremock::matchers::method("GET")) 891 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 892 .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) 893 .expect(1) 894 .mount(&server) 895 .await; 896 897 let client = 898 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 899 let resolver = 900 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 901 902 let owner = did("did:plc:nel"); 903 let key = rkey("abcabcabcabcz"); 904 let first = resolver.resolve(&owner, &key).await; 905 let second = resolver.resolve(&owner, &key).await; 906 assert_eq!(first, Resolution::Unresolvable); 907 assert_eq!(second, Resolution::Unresolvable); 908 } 909 910 #[tokio::test] 911 async fn slingshot_legacy_repo_body_resolves_no_repo_did_not_unresolvable() { 912 let server = wiremock::MockServer::start().await; 913 let body = serde_json::json!({ 914 "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", 915 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 916 "value": { 917 "$type": "sh.tangled.repo", 918 "addedAt": "2025-03-07T21:47:53Z", 919 "knot": "knot1.tangled.sh", 920 "name": "scallop", 921 "owner": "did:plc:nel", 922 }, 923 }); 924 wiremock::Mock::given(wiremock::matchers::method("GET")) 925 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 926 .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) 927 .expect(1) 928 .mount(&server) 929 .await; 930 931 let client = 932 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 933 let resolver = 934 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 935 936 let owner = did("did:plc:nel"); 937 let key = rkey("abcabcabcabcz"); 938 let got = resolver.resolve(&owner, &key).await; 939 assert_eq!( 940 got, 941 Resolution::NoRepoDid, 942 "legacy repo wires without a repo_did parse via legacy upgrade and resolve as NoRepoDid, not Unresolvable", 943 ); 944 } 945 946 #[tokio::test] 947 async fn slingshot_unparseable_repo_value_caches_as_unresolvable() { 948 let server = wiremock::MockServer::start().await; 949 let body = serde_json::json!({ 950 "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", 951 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 952 "value": {"$type": "sh.tangled.repo"} 953 }); 954 wiremock::Mock::given(wiremock::matchers::method("GET")) 955 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 956 .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) 957 .expect(1) 958 .mount(&server) 959 .await; 960 961 let client = 962 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 963 let resolver = 964 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 965 966 let owner = did("did:plc:nel"); 967 let key = rkey("abcabcabcabcz"); 968 let first = resolver.resolve(&owner, &key).await; 969 let second = resolver.resolve(&owner, &key).await; 970 assert_eq!( 971 first, 972 Resolution::Unresolvable, 973 "a repo body that fails lexicon validation must not be conflated with NoRepoDid", 974 ); 975 assert_eq!(second, Resolution::Unresolvable); 976 } 977 978 #[tokio::test] 979 async fn slingshot_transport_error_caches_with_short_ttl() { 980 let server = wiremock::MockServer::start().await; 981 wiremock::Mock::given(wiremock::matchers::method("GET")) 982 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 983 .respond_with(wiremock::ResponseTemplate::new(503)) 984 .expect(1) 985 .mount(&server) 986 .await; 987 988 let client = 989 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 990 let resolver = 991 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 992 993 let owner = did("did:plc:nel"); 994 let key = rkey("abcabcabcabcz"); 995 let first = resolver.resolve(&owner, &key).await; 996 let second = resolver.resolve(&owner, &key).await; 997 assert_eq!(first, Resolution::Unresolvable); 998 assert_eq!( 999 second, 1000 Resolution::Unresolvable, 1001 "transient TTL must suppress immediate re-hammering of a sick upstream", 1002 ); 1003 } 1004 1005 #[tokio::test] 1006 async fn slingshot_transient_recorded_separately_from_unresolvable() { 1007 let server = wiremock::MockServer::start().await; 1008 wiremock::Mock::given(wiremock::matchers::method("GET")) 1009 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 1010 .respond_with(wiremock::ResponseTemplate::new(503)) 1011 .expect(1) 1012 .mount(&server) 1013 .await; 1014 1015 let client = 1016 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 1017 let resolver = 1018 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 1019 1020 let owner = did("did:plc:nel"); 1021 let key = rkey("abcabcabcabcz"); 1022 resolver.resolve(&owner, &key).await; 1023 resolver.resolve(&owner, &key).await; 1024 1025 let snap = resolver.stats(); 1026 assert_eq!( 1027 snap.misses_transient, 1, 1028 "second resolve must hit the short-TTL cache instead of re-firing the transient miss", 1029 ); 1030 assert_eq!(snap.hits, 1, "second call hits cached transient entry"); 1031 assert_eq!( 1032 snap.misses_unresolvable, 0, 1033 "canonical unresolvable counter is reserved for cached terminal answers", 1034 ); 1035 assert!( 1036 snap.miss_latency_micros_sum > 0, 1037 "transient misses still have latency contributions", 1038 ); 1039 assert_eq!(snap.miss_count(), 1); 1040 } 1041 1042 #[tokio::test] 1043 async fn slingshot_in_flight_requests_coalesce() { 1044 let server = wiremock::MockServer::start().await; 1045 let body = serde_json::json!({ 1046 "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", 1047 "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", 1048 "value": {"$type": "sh.tangled.repo", "knot": "oyster.cafe", "createdAt": "2026-05-01T00:00:00Z", "repoDid": "did:plc:limpet"} 1049 }); 1050 wiremock::Mock::given(wiremock::matchers::method("GET")) 1051 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 1052 .respond_with( 1053 wiremock::ResponseTemplate::new(200) 1054 .set_body_json(body) 1055 .set_delay(Duration::from_millis(200)), 1056 ) 1057 .expect(1) 1058 .mount(&server) 1059 .await; 1060 1061 let client = 1062 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 1063 let resolver = Arc::new(RepoIdResolver::with_slingshot( 1064 client, 1065 test_clock(), 1066 RuntimeHasher::default(), 1067 )); 1068 1069 let owner = did("did:plc:nel"); 1070 let key = rkey("abcabcabcabcz"); 1071 let r0 = resolver.clone(); 1072 let r1 = resolver.clone(); 1073 let r2 = resolver.clone(); 1074 let o0 = owner.clone(); 1075 let o1 = owner.clone(); 1076 let o2 = owner.clone(); 1077 let k0 = key.clone(); 1078 let k1 = key.clone(); 1079 let k2 = key.clone(); 1080 let (a, b, c) = tokio::join!( 1081 tokio::spawn(async move { r0.resolve(&o0, &k0).await }), 1082 tokio::spawn(async move { r1.resolve(&o1, &k1).await }), 1083 tokio::spawn(async move { r2.resolve(&o2, &k2).await }), 1084 ); 1085 let expected = Resolution::Mapped(did("did:plc:limpet")); 1086 assert_eq!(a.unwrap(), expected); 1087 assert_eq!(b.unwrap(), expected); 1088 assert_eq!(c.unwrap(), expected); 1089 1090 let snap = resolver.stats(); 1091 assert_eq!( 1092 snap.misses_mapped, 1, 1093 "only the winning task pays the slingshot RTT", 1094 ); 1095 } 1096 1097 #[tokio::test] 1098 async fn stats_count_hits_misses_and_latency() { 1099 let server = wiremock::MockServer::start().await; 1100 wiremock::Mock::given(wiremock::matchers::method("GET")) 1101 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 1102 .respond_with(wiremock::ResponseTemplate::new(404)) 1103 .mount(&server) 1104 .await; 1105 let client = 1106 SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); 1107 let resolver = 1108 RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); 1109 1110 let owner = did("did:plc:nel"); 1111 let key = rkey("abcabcabcabcz"); 1112 resolver.resolve(&owner, &key).await; 1113 resolver.resolve(&owner, &key).await; 1114 1115 let snap = resolver.stats(); 1116 assert_eq!( 1117 snap.misses_unresolvable, 1, 1118 "first call is the slingshot miss" 1119 ); 1120 assert_eq!(snap.hits, 1, "second call hits the unresolvable cache"); 1121 assert_eq!(snap.miss_count(), 1); 1122 assert_eq!(snap.total(), 2); 1123 assert!( 1124 snap.miss_latency_micros_sum > 0, 1125 "latency recorded for slingshot miss" 1126 ); 1127 assert!(snap.miss_latency_micros_avg().unwrap() > 0); 1128 } 1129 1130 #[tokio::test] 1131 async fn stats_no_client_miss_recorded_without_latency() { 1132 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 1133 resolver 1134 .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) 1135 .await; 1136 let snap = resolver.stats(); 1137 assert_eq!(snap.misses_no_client, 1); 1138 assert_eq!(snap.miss_latency_micros_sum, 0); 1139 assert_eq!(snap.miss_latency_micros_avg(), None); 1140 } 1141 1142 #[tokio::test] 1143 async fn firehose_observe_can_demote_provisional() { 1144 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 1145 let owner = did("did:plc:nel"); 1146 let key = rkey("abcabcabcabcz"); 1147 resolver 1148 .fill_provisional( 1149 RepoIdent::new(owner.clone(), key.clone()), 1150 Resolution::Mapped(did("did:plc:clam")), 1151 ) 1152 .await; 1153 resolver.observe(owner.clone(), key.clone(), None, None).await; 1154 let got = resolver.resolve(&owner, &key).await; 1155 assert_eq!( 1156 got, 1157 Resolution::NoRepoDid, 1158 "firehose update is canonical and may legitimately remove repoDID", 1159 ); 1160 } 1161 1162 #[tokio::test] 1163 async fn forget_removes_cache_entry() { 1164 let resolver = RepoIdResolver::detached(RuntimeHasher::default()); 1165 let owner = did("did:plc:nel"); 1166 let key = rkey("abcabcabcabcz"); 1167 resolver 1168 .observe(owner.clone(), key.clone(), Some(did("did:plc:clam")), None) 1169 .await; 1170 assert_eq!( 1171 resolver.cached_resolution(&owner, &key).await, 1172 Some(Resolution::Mapped(did("did:plc:clam"))), 1173 ); 1174 resolver.forget(&owner, &key).await; 1175 assert_eq!( 1176 resolver.cached_resolution(&owner, &key).await, 1177 None, 1178 "forget must drop the entry entirely so a subsequent observe can supply fresh state", 1179 ); 1180 } 1181}