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