Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156mod legacy_upgrade;mod normalize;
pub use legacy_upgrade::{ DecodedRecord, decode_canon_or_upgrade, decode_canon_or_upgrade_bytes, normalize_record_fields, scrub_record_bytes, synthesize_created_at, upgrade, upgrade_wire_bytes,};pub use normalize::NormalizeRepoRefs;
use std::sync::Arc;use std::sync::atomic::{AtomicU64, Ordering};use std::time::Duration;
use tokio::time::Instant;
use bobbin_runtime::{Clock, RuntimeHasher};use bobbin_slingshot_client::{SlingshotClient, SlingshotError};use bobbin_types::edges::{ExtractError, Record};use bobbin_types::ids::{RepoIdent, nsid_static};use jacquard_common::DefaultStr;use jacquard_common::types::datetime::Datetime;use jacquard_common::types::did::Did;use jacquard_common::types::nsid::Nsid;use jacquard_common::types::recordkey::Rkey;use scc::HashMap as SccMap;use tokio::sync::OnceCell;use tracing::warn;
const REPO_COLLECTION: &str = "sh.tangled.repo";const TRANSIENT_TTL: Duration = Duration::from_secs(60);
#[derive(Clone, Debug, Eq, PartialEq)]pub enum Resolution { Mapped(Did<DefaultStr>), NoRepoDid, Unresolvable,}
#[derive(Clone, Debug, Eq, PartialEq)]enum AuthoritativeResolution { Mapped(Did<DefaultStr>), NoRepoDid,}
impl AuthoritativeResolution { fn from_repo_did(repo_did: Option<Did<DefaultStr>>) -> Self { match repo_did { Some(did) => Self::Mapped(did), None => Self::NoRepoDid, } }}
#[derive(Clone, Debug, Eq, PartialEq)]struct RepoDidHolder { ident: RepoIdent, created_at: Datetime,}
impl RepoDidHolder { // a backfill replays a repo in rkey order, so the stamp decides and not arrival fn yields_to(&self, incoming: &Self) -> bool { self.ident == incoming.ident || incoming.created_at >= self.created_at }}
#[derive(Clone, Debug, Eq, PartialEq)]pub enum RepoClaim { Current { displaced: Option<RepoIdent> }, Superseded { by: RepoIdent },}
#[derive(Clone, Debug, Eq, PartialEq)]enum CacheEntry { Authoritative(AuthoritativeResolution), Provisional(Resolution), Transient { expires_at: Instant },}
impl CacheEntry { fn into_resolution(self) -> Resolution { match self { Self::Authoritative(AuthoritativeResolution::Mapped(did)) => Resolution::Mapped(did), Self::Authoritative(AuthoritativeResolution::NoRepoDid) => Resolution::NoRepoDid, Self::Provisional(r) => r, Self::Transient { .. } => Resolution::Unresolvable, } }
fn is_expired_transient(&self, now: Instant) -> bool { matches!(self, Self::Transient { expires_at } if *expires_at <= now) }}
#[derive(Default)]pub struct ResolverStats { hits: AtomicU64, misses_mapped: AtomicU64, misses_no_repo_did: AtomicU64, misses_unresolvable: AtomicU64, misses_transient: AtomicU64, misses_no_client: AtomicU64, miss_latency_micros_sum: AtomicU64, miss_latency_micros_max: AtomicU64,}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]pub struct ResolverStatsSnapshot { pub hits: u64, pub misses_mapped: u64, pub misses_no_repo_did: u64, pub misses_unresolvable: u64, pub misses_transient: u64, pub misses_no_client: u64, pub miss_latency_micros_sum: u64, pub miss_latency_micros_max: u64,}
impl ResolverStatsSnapshot { pub fn miss_count(&self) -> u64 { self.misses_mapped + self.misses_no_repo_did + self.misses_unresolvable + self.misses_transient + self.misses_no_client }
pub fn total(&self) -> u64 { self.hits + self.miss_count() }
pub fn miss_latency_micros_avg(&self) -> Option<u64> { let misses = self.miss_count() - self.misses_no_client; (misses > 0).then(|| self.miss_latency_micros_sum / misses) }}
#[derive(Clone, Copy)]enum MissKind { Mapped, NoRepoDid, Unresolvable, Transient, NoClient,}
impl ResolverStats { fn record_hit(&self) { self.hits.fetch_add(1, Ordering::Relaxed); }
fn record_miss(&self, kind: MissKind, latency: Option<Duration>) { let counter = match kind { MissKind::Mapped => &self.misses_mapped, MissKind::NoRepoDid => &self.misses_no_repo_did, MissKind::Unresolvable => &self.misses_unresolvable, MissKind::Transient => &self.misses_transient, MissKind::NoClient => &self.misses_no_client, }; counter.fetch_add(1, Ordering::Relaxed); if let Some(latency) = latency { let micros = u64::try_from(latency.as_micros()).unwrap_or(u64::MAX); self.miss_latency_micros_sum .fetch_add(micros, Ordering::Relaxed); self.miss_latency_micros_max .fetch_max(micros, Ordering::Relaxed); } }
pub fn snapshot(&self) -> ResolverStatsSnapshot { ResolverStatsSnapshot { hits: self.hits.load(Ordering::Relaxed), misses_mapped: self.misses_mapped.load(Ordering::Relaxed), misses_no_repo_did: self.misses_no_repo_did.load(Ordering::Relaxed), misses_unresolvable: self.misses_unresolvable.load(Ordering::Relaxed), misses_transient: self.misses_transient.load(Ordering::Relaxed), misses_no_client: self.misses_no_client.load(Ordering::Relaxed), miss_latency_micros_sum: self.miss_latency_micros_sum.load(Ordering::Relaxed), miss_latency_micros_max: self.miss_latency_micros_max.load(Ordering::Relaxed), } }}
struct SlingshotProbe { client: SlingshotClient, clock: Arc<dyn Clock>,}
pub struct RepoIdResolver { cache: SccMap<RepoIdent, CacheEntry, RuntimeHasher>, by_repo_did: SccMap<Did<DefaultStr>, RepoDidHolder, RuntimeHasher>, in_flight: SccMap<RepoIdent, Arc<OnceCell<Resolution>>, RuntimeHasher>, probe: Option<SlingshotProbe>, stats: ResolverStats,}
impl RepoIdResolver { pub fn with_slingshot( client: SlingshotClient, clock: Arc<dyn Clock>, hasher: RuntimeHasher, ) -> Self { Self { cache: SccMap::with_hasher(hasher.clone()), by_repo_did: SccMap::with_hasher(hasher.clone()), in_flight: SccMap::with_hasher(hasher), probe: Some(SlingshotProbe { client, clock }), stats: ResolverStats::default(), } }
pub fn detached(hasher: RuntimeHasher) -> Self { Self { cache: SccMap::with_hasher(hasher.clone()), by_repo_did: SccMap::with_hasher(hasher.clone()), in_flight: SccMap::with_hasher(hasher), probe: None, stats: ResolverStats::default(), } }
pub fn stats(&self) -> ResolverStatsSnapshot { self.stats.snapshot() }
pub async fn cached_resolution( &self, owner: &Did<DefaultStr>, rkey: &Rkey<DefaultStr>, ) -> Option<Resolution> { let key = RepoIdent::new(owner.clone(), rkey.clone()); let entry = self.cache.get_async(&key).await?; let now = self.probe.as_ref().map(|p| p.clock.now_instant()); if let Some(now) = now && entry.get().is_expired_transient(now) { return None; } Some(entry.get().clone().into_resolution()) }
pub async fn lookup_by_repo_did(&self, repo_did: &Did<DefaultStr>) -> Option<RepoIdent> { self.by_repo_did .get_async(repo_did) .await .map(|e| e.get().ident.clone()) }
pub async fn observe( &self, owner: Did<DefaultStr>, rkey: Rkey<DefaultStr>, repo_did: Option<Did<DefaultStr>>, created_at: &Datetime, ) -> RepoClaim { let ident = RepoIdent::new(owner, rkey); let entry = CacheEntry::Authoritative(AuthoritativeResolution::from_repo_did(repo_did.clone())); self.cache .entry_async(ident.clone()) .await .and_modify(|existing| *existing = entry.clone()) .or_insert(entry);
let Some(repo_did) = repo_did else { return RepoClaim::Current { displaced: None }; }; let incoming = RepoDidHolder { ident, created_at: created_at.clone(), }; let mut outcome = RepoClaim::Current { displaced: None }; self.by_repo_did .entry_async(repo_did) .await .and_modify(|held| { let yields = held.yields_to(&incoming); outcome = yields .then(|| RepoClaim::Current { displaced: (held.ident != incoming.ident).then(|| held.ident.clone()), }) .unwrap_or_else(|| RepoClaim::Superseded { by: held.ident.clone(), }); if yields { *held = incoming.clone(); } }) .or_insert(incoming); outcome }
pub async fn forget(&self, owner: &Did<DefaultStr>, rkey: &Rkey<DefaultStr>) { let ident = RepoIdent::new(owner.clone(), rkey.clone()); let prior_resolution = self .cache .remove_async(&ident) .await .map(|(_, entry)| entry.into_resolution()); if let Some(Resolution::Mapped(repo_did)) = prior_resolution { self.by_repo_did .remove_if_async(&repo_did, |held| held.ident == ident) .await; } }
async fn fill_provisional(&self, key: RepoIdent, resolution: Resolution) { let entry = CacheEntry::Provisional(resolution); self.cache .entry_async(key) .await .and_modify(|existing| { if matches!(existing, CacheEntry::Authoritative(_)) { return; } *existing = entry.clone(); }) .or_insert(entry); }
async fn fill_transient(&self, key: RepoIdent, expires_at: Instant) { let entry = CacheEntry::Transient { expires_at }; self.cache .entry_async(key) .await .and_modify(|existing| { if matches!(existing, CacheEntry::Authoritative(_)) { return; } *existing = entry.clone(); }) .or_insert(entry); }
pub async fn resolve(&self, owner: &Did<DefaultStr>, rkey: &Rkey<DefaultStr>) -> Resolution { let key = RepoIdent::new(owner.clone(), rkey.clone());
let Some(probe) = self.probe.as_ref() else { if let Some(entry) = self.cache.get_async(&key).await { self.stats.record_hit(); return entry.get().clone().into_resolution(); } self.stats.record_miss(MissKind::NoClient, None); return Resolution::Unresolvable; };
let now = probe.clock.now_instant(); if let Some(entry) = self.cache.get_async(&key).await && !entry.get().is_expired_transient(now) { self.stats.record_hit(); return entry.get().clone().into_resolution(); }
let cell: Arc<OnceCell<Resolution>> = self .in_flight .entry_async(key.clone()) .await .or_insert_with(|| Arc::new(OnceCell::new())) .get() .clone();
let result = cell .get_or_init(|| async { self.fetch_repo_did(owner, rkey, &key).await }) .await .clone();
self.in_flight.remove_async(&key).await;
result }
async fn fetch_repo_did( &self, owner: &Did<DefaultStr>, rkey: &Rkey<DefaultStr>, key: &RepoIdent, ) -> Resolution { let probe = self .probe .as_ref() .expect("fetch_repo_did is only called when a probe is present"); let started = probe.clock.now_instant(); let nsid: Nsid<DefaultStr> = nsid_static(REPO_COLLECTION); let provisional = match probe.client.get_record(owner, &nsid, rkey).await { Ok(body) => match repo_did_from_body(&nsid, &body.value) { Ok(Some(did)) => Resolution::Mapped(did), Ok(None) => Resolution::NoRepoDid, Err(e) => { warn!( error = ?e, owner = owner.as_ref(), rkey = rkey.as_ref(), "slingshot returned unparseable repo body, caching as unresolvable", ); Resolution::Unresolvable } }, Err(SlingshotError::NotFound) => { warn!( owner = owner.as_ref(), rkey = rkey.as_ref(), "no repo record on slingshot, caching as unresolvable", ); Resolution::Unresolvable } Err(ref e) if is_garbage_response(e) => { warn!( error = ?e, owner = owner.as_ref(), rkey = rkey.as_ref(), "slingshot returned malformed response, caching as unresolvable", ); Resolution::Unresolvable } Err(e) => { warn!( error = ?e, owner = owner.as_ref(), rkey = rkey.as_ref(), "caching transient slingshot failure for repoDID lookup under short TTL", ); let elapsed = probe.clock.now_instant().duration_since(started); self.stats.record_miss(MissKind::Transient, Some(elapsed)); let expires_at = probe.clock.now_instant() + TRANSIENT_TTL; self.fill_transient(key.clone(), expires_at).await; return Resolution::Unresolvable; } }; let elapsed = probe.clock.now_instant().duration_since(started); let kind = match &provisional { Resolution::Mapped(_) => MissKind::Mapped, Resolution::NoRepoDid => MissKind::NoRepoDid, Resolution::Unresolvable => MissKind::Unresolvable, }; self.stats.record_miss(kind, Some(elapsed)); self.fill_provisional(key.clone(), provisional.clone()) .await; provisional }}
fn is_garbage_response(err: &SlingshotError) -> bool { matches!( err, SlingshotError::Decode(_) | SlingshotError::MissingField(_) | SlingshotError::InvalidAtUri(_) | SlingshotError::InvalidCid(_) | SlingshotError::UriMismatch { .. }, )}
fn repo_did_from_body( nsid: &Nsid<DefaultStr>, body: &[u8],) -> Result<Option<Did<DefaultStr>>, ExtractError> { match DecodedRecord::try_decode(nsid, body)? { DecodedRecord::Canon(Record::Repo(repo)) => Ok(repo.repo_did), DecodedRecord::Canon(_) | DecodedRecord::Legacy(_) => Ok(None), }}
#[cfg(test)]mod tests { use super::*; use bobbin_runtime::SystemClock; use jacquard_common::types::did::Did; use jacquard_common::types::recordkey::Rkey;
fn did(s: &str) -> Did<DefaultStr> { Did::new_owned(s).unwrap() }
fn rkey(s: &str) -> Rkey<DefaultStr> { Rkey::new_owned(s).unwrap() }
fn test_clock() -> Arc<dyn Clock> { Arc::new(SystemClock::new()) }
fn made(day: u32) -> Datetime { Datetime::raw_str(format!("2026-05-{day:02}T00:00:00Z")) }
#[tokio::test] async fn observation_returns_prior_ident_when_repo_did_moves() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let claim = resolver .observe( did("did:plc:nel"), rkey("3liuighjy2h22"), Some(did("did:plc:clam")), &made(1), ) .await; assert_eq!( claim, RepoClaim::Current { displaced: None }, "first observation has no prior" );
let claim = resolver .observe( did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam")), &made(2), ) .await; assert_eq!( claim, RepoClaim::Current { displaced: Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), }, "same repoDID at a new (owner, rkey) returns the prior ident so callers can evict the stale at-uri", );
let claim = resolver .observe( did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam")), &made(2), ) .await; assert_eq!( claim, RepoClaim::Current { displaced: None }, "re-observing the same ident is a no-op" ); }
#[tokio::test] async fn a_record_older_than_the_holder_cannot_take_the_repo_did() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let owner = did("did:plc:oppi"); let repo_did = did("did:plc:clam"); resolver .observe( owner.clone(), rkey("stinkpot"), Some(repo_did.clone()), &made(24), ) .await; let claim = resolver .observe( owner.clone(), rkey("tortu"), Some(repo_did.clone()), &made(20), ) .await;
let stinkpot = RepoIdent::new(owner, rkey("stinkpot")); assert_eq!( claim, RepoClaim::Superseded { by: stinkpot.clone() }, "the older record is the alias a rename left behind", ); assert_eq!(resolver.lookup_by_repo_did(&repo_did).await, Some(stinkpot)); }
#[tokio::test] async fn observation_without_repo_did_does_not_track_reverse() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let claim = resolver .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, &made(1)) .await; assert_eq!(claim, RepoClaim::Current { displaced: None }); }
#[tokio::test] async fn forget_clears_reverse_only_when_still_owned() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); resolver .observe( did("did:plc:nel"), rkey("3liuighjy2h22"), Some(did("did:plc:clam")), &made(1), ) .await; resolver .observe( did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam")), &made(1), ) .await;
resolver .forget(&did("did:plc:nel"), &rkey("3liuighjy2h22")) .await;
let claim = resolver .observe( did("did:plc:nel"), rkey("core-renamed"), Some(did("did:plc:clam")), &made(1), ) .await; assert_eq!( claim, RepoClaim::Current { displaced: Some(RepoIdent::new(did("did:plc:nel"), rkey("core"))), }, "stale at-uri's forget must not displace the live owner of did:plc:clam", ); }
#[tokio::test] async fn observation_with_repo_did_resolves_mapped() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); resolver .observe( did("did:plc:nel"), rkey("abcabcabcabcz"), Some(did("did:plc:clam")), &made(1), ) .await; let got = resolver .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) .await; assert_eq!(got, Resolution::Mapped(did("did:plc:clam"))); }
#[tokio::test] async fn observation_without_repo_did_resolves_no_repo_did() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); resolver .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, &made(1)) .await; let got = resolver .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) .await; assert_eq!( got, Resolution::NoRepoDid, "observed but empty repoDID is a definitive answer not a lookup failure", ); }
#[tokio::test] async fn cache_miss_without_client_is_unresolvable() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let got = resolver .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) .await; assert_eq!(got, Resolution::Unresolvable); }
#[tokio::test] async fn lookup_by_repo_did_finds_observed_ident() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); resolver .observe( did("did:plc:nel"), rkey("abcabcabcabcz"), Some(did("did:plc:limpet")), &made(1), ) .await; let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; assert_eq!( got, Some(RepoIdent::new(did("did:plc:nel"), rkey("abcabcabcabcz"))), ); }
#[tokio::test] async fn lookup_by_repo_did_misses_when_unobserved() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; assert_eq!(got, None); }
#[tokio::test] async fn lookup_by_repo_did_misses_when_repo_did_was_none() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); resolver .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, &made(1)) .await; let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; assert_eq!(got, None); }
#[tokio::test] async fn lookup_by_repo_did_follows_move_to_new_ident() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); resolver .observe( did("did:plc:nel"), rkey("abcabcabcabcz"), Some(did("did:plc:limpet")), &made(1), ) .await; resolver .observe( did("did:plc:olaren"), rkey("xyzxyzxyzxyzx"), Some(did("did:plc:limpet")), &made(1), ) .await; let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; assert_eq!( got, Some(RepoIdent::new(did("did:plc:olaren"), rkey("xyzxyzxyzxyzx"))), ); }
#[tokio::test] async fn observation_overwrites_prior_value() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); resolver .observe( did("did:plc:nel"), rkey("abcabcabcabcz"), Some(did("did:plc:clam")), &made(1), ) .await; resolver .observe( did("did:plc:nel"), rkey("abcabcabcabcz"), Some(did("did:plc:uni")), &made(1), ) .await; let got = resolver .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) .await; assert_eq!(got, Resolution::Mapped(did("did:plc:uni"))); }
#[tokio::test] async fn fill_provisional_does_not_downgrade_authoritative_mapped() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver .observe( owner.clone(), key.clone(), Some(did("did:plc:clam")), &made(1), ) .await; resolver .fill_provisional( RepoIdent::new(owner.clone(), key.clone()), Resolution::Unresolvable, ) .await; let got = resolver.resolve(&owner, &key).await; assert_eq!( got, Resolution::Mapped(did("did:plc:clam")), "firehose-observed mapping must outrank provisional slingshot info", ); }
#[tokio::test] async fn fill_provisional_does_not_downgrade_authoritative_no_repo_did() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver .observe(owner.clone(), key.clone(), None, &made(1)) .await; resolver .fill_provisional( RepoIdent::new(owner.clone(), key.clone()), Resolution::Mapped(did("did:plc:clam")), ) .await; let got = resolver.resolve(&owner, &key).await; assert_eq!( got, Resolution::NoRepoDid, "an authoritative empty observation must outrank provisional slingshot info even when slingshot disagrees", ); }
#[tokio::test] async fn slingshot_404_caches_as_unresolvable() { let server = wiremock::MockServer::start().await; wiremock::Mock::given(wiremock::matchers::method("GET")) .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) .respond_with(wiremock::ResponseTemplate::new(404)) .expect(1) .mount(&server) .await;
let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default());
let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); let first = resolver.resolve(&owner, &key).await; let second = resolver.resolve(&owner, &key).await; assert_eq!(first, Resolution::Unresolvable); assert_eq!(second, Resolution::Unresolvable); }
#[tokio::test] async fn slingshot_malformed_envelope_caches_as_unresolvable() { let server = wiremock::MockServer::start().await; wiremock::Mock::given(wiremock::matchers::method("GET")) .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) .respond_with( wiremock::ResponseTemplate::new(200) .insert_header("content-type", "application/json") .set_body_string("not json"), ) .expect(1) .mount(&server) .await;
let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default());
let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); let first = resolver.resolve(&owner, &key).await; let second = resolver.resolve(&owner, &key).await; assert_eq!(first, Resolution::Unresolvable); assert_eq!( second, Resolution::Unresolvable, "garbage envelopes are stable across retries, so caching avoids hammering slingshot", ); }
#[tokio::test] async fn slingshot_uri_mismatch_caches_as_unresolvable() { let server = wiremock::MockServer::start().await; let body = serde_json::json!({ "uri": "at://did:plc:limpet/sh.tangled.repo/elsewhere", "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", "value": {"$type": "sh.tangled.repo", "knot": "oyster.cafe", "createdAt": "2026-05-01T00:00:00Z"} }); wiremock::Mock::given(wiremock::matchers::method("GET")) .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) .expect(1) .mount(&server) .await;
let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default());
let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); let first = resolver.resolve(&owner, &key).await; let second = resolver.resolve(&owner, &key).await; assert_eq!(first, Resolution::Unresolvable); assert_eq!(second, Resolution::Unresolvable); }
#[tokio::test] async fn slingshot_legacy_repo_body_resolves_no_repo_did_not_unresolvable() { let server = wiremock::MockServer::start().await; let body = serde_json::json!({ "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", "value": { "$type": "sh.tangled.repo", "addedAt": "2025-03-07T21:47:53Z", "knot": "knot1.tangled.sh", "name": "scallop", "owner": "did:plc:nel", }, }); wiremock::Mock::given(wiremock::matchers::method("GET")) .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) .expect(1) .mount(&server) .await;
let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default());
let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); let got = resolver.resolve(&owner, &key).await; assert_eq!( got, Resolution::NoRepoDid, "legacy repo wires without a repo_did parse via legacy upgrade and resolve as NoRepoDid, not Unresolvable", ); }
#[tokio::test] async fn slingshot_unparseable_repo_value_caches_as_unresolvable() { let server = wiremock::MockServer::start().await; let body = serde_json::json!({ "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", "value": {"$type": "sh.tangled.repo"} }); wiremock::Mock::given(wiremock::matchers::method("GET")) .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) .respond_with(wiremock::ResponseTemplate::new(200).set_body_json(body)) .expect(1) .mount(&server) .await;
let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default());
let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); let first = resolver.resolve(&owner, &key).await; let second = resolver.resolve(&owner, &key).await; assert_eq!( first, Resolution::Unresolvable, "a repo body that fails lexicon validation must not be conflated with NoRepoDid", ); assert_eq!(second, Resolution::Unresolvable); }
#[tokio::test] async fn slingshot_transport_error_caches_with_short_ttl() { let server = wiremock::MockServer::start().await; wiremock::Mock::given(wiremock::matchers::method("GET")) .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) .respond_with(wiremock::ResponseTemplate::new(503)) .expect(1) .mount(&server) .await;
let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default());
let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); let first = resolver.resolve(&owner, &key).await; let second = resolver.resolve(&owner, &key).await; assert_eq!(first, Resolution::Unresolvable); assert_eq!( second, Resolution::Unresolvable, "transient TTL must suppress immediate re-hammering of a sick upstream", ); }
#[tokio::test] async fn slingshot_transient_recorded_separately_from_unresolvable() { let server = wiremock::MockServer::start().await; wiremock::Mock::given(wiremock::matchers::method("GET")) .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) .respond_with(wiremock::ResponseTemplate::new(503)) .expect(1) .mount(&server) .await;
let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default());
let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver.resolve(&owner, &key).await; resolver.resolve(&owner, &key).await;
let snap = resolver.stats(); assert_eq!( snap.misses_transient, 1, "second resolve must hit the short-TTL cache instead of re-firing the transient miss", ); assert_eq!(snap.hits, 1, "second call hits cached transient entry"); assert_eq!( snap.misses_unresolvable, 0, "canonical unresolvable counter is reserved for cached terminal answers", ); assert!( snap.miss_latency_micros_sum > 0, "transient misses still have latency contributions", ); assert_eq!(snap.miss_count(), 1); }
#[tokio::test] async fn slingshot_in_flight_requests_coalesce() { let server = wiremock::MockServer::start().await; let body = serde_json::json!({ "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", "value": {"$type": "sh.tangled.repo", "knot": "oyster.cafe", "createdAt": "2026-05-01T00:00:00Z", "repoDid": "did:plc:limpet"} }); wiremock::Mock::given(wiremock::matchers::method("GET")) .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) .respond_with( wiremock::ResponseTemplate::new(200) .set_body_json(body) .set_delay(Duration::from_millis(200)), ) .expect(1) .mount(&server) .await;
let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); let resolver = Arc::new(RepoIdResolver::with_slingshot( client, test_clock(), RuntimeHasher::default(), ));
let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); let r0 = resolver.clone(); let r1 = resolver.clone(); let r2 = resolver.clone(); let o0 = owner.clone(); let o1 = owner.clone(); let o2 = owner.clone(); let k0 = key.clone(); let k1 = key.clone(); let k2 = key.clone(); let (a, b, c) = tokio::join!( tokio::spawn(async move { r0.resolve(&o0, &k0).await }), tokio::spawn(async move { r1.resolve(&o1, &k1).await }), tokio::spawn(async move { r2.resolve(&o2, &k2).await }), ); let expected = Resolution::Mapped(did("did:plc:limpet")); assert_eq!(a.unwrap(), expected); assert_eq!(b.unwrap(), expected); assert_eq!(c.unwrap(), expected);
let snap = resolver.stats(); assert_eq!( snap.misses_mapped, 1, "only the winning task pays the slingshot RTT", ); }
#[tokio::test] async fn stats_count_hits_misses_and_latency() { let server = wiremock::MockServer::start().await; wiremock::Mock::given(wiremock::matchers::method("GET")) .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) .respond_with(wiremock::ResponseTemplate::new(404)) .mount(&server) .await; let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default());
let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver.resolve(&owner, &key).await; resolver.resolve(&owner, &key).await;
let snap = resolver.stats(); assert_eq!( snap.misses_unresolvable, 1, "first call is the slingshot miss" ); assert_eq!(snap.hits, 1, "second call hits the unresolvable cache"); assert_eq!(snap.miss_count(), 1); assert_eq!(snap.total(), 2); assert!( snap.miss_latency_micros_sum > 0, "latency recorded for slingshot miss" ); assert!(snap.miss_latency_micros_avg().unwrap() > 0); }
#[tokio::test] async fn stats_no_client_miss_recorded_without_latency() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); resolver .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) .await; let snap = resolver.stats(); assert_eq!(snap.misses_no_client, 1); assert_eq!(snap.miss_latency_micros_sum, 0); assert_eq!(snap.miss_latency_micros_avg(), None); }
#[tokio::test] async fn firehose_observe_can_demote_provisional() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver .fill_provisional( RepoIdent::new(owner.clone(), key.clone()), Resolution::Mapped(did("did:plc:clam")), ) .await; resolver .observe(owner.clone(), key.clone(), None, &made(1)) .await; let got = resolver.resolve(&owner, &key).await; assert_eq!( got, Resolution::NoRepoDid, "firehose update is canonical and may legitimately remove repoDID", ); }
#[tokio::test] async fn forget_removes_cache_entry() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver .observe( owner.clone(), key.clone(), Some(did("did:plc:clam")), &made(1), ) .await; assert_eq!( resolver.cached_resolution(&owner, &key).await, Some(Resolution::Mapped(did("did:plc:clam"))), ); resolver.forget(&owner, &key).await; assert_eq!( resolver.cached_resolution(&owner, &key).await, None, "forget must drop the entry entirely so a subsequent observe can supply fresh state", ); }}