This repository has no description
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}