diff --git a/crates/ingest/Cargo.toml b/crates/ingest/Cargo.toml index 26563f9..1ef2e8c 100644 --- a/crates/ingest/Cargo.toml +++ b/crates/ingest/Cargo.toml @@ -9,16 +9,23 @@ rust-version.workspace = true bobbin-types = { workspace = true } bobbin-edge-index = { workspace = true } bobbin-record-lru = { workspace = true } +bobbin-slingshot-client = { workspace = true } jacquard-common = { workspace = true } +bytes = { workspace = true } futures = { workspace = true } +scc = { workspace = true } serde = { workspace = true } serde_json = { workspace = true } thiserror = { workspace = true } tokio = { workspace = true } tokio-tungstenite = { workspace = true } +tokio-util = { workspace = true } tracing = { workspace = true } url = { workspace = true } [dev-dependencies] -tracing-subscriber = { version = "0.3", features = ["env-filter"] } +tracing-subscriber = { workspace = true } +tokio-util = { workspace = true } +tokio = { workspace = true, features = ["test-util"] } +wiremock = { workspace = true } diff --git a/crates/ingest/src/resolver.rs b/crates/ingest/src/resolver.rs new file mode 100644 index 0000000..9eeaf22 --- /dev/null +++ b/crates/ingest/src/resolver.rs @@ -0,0 +1,450 @@ +use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; +use bobbin_types::edges::{ExtractError, Record}; +use bobbin_types::ids::nsid_static; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::nsid::Nsid; +use jacquard_common::types::recordkey::Rkey; +use scc::HashMap as SccMap; +use tracing::warn; + +const REPO_COLLECTION: &str = "sh.tangled.repo"; + +#[derive(Clone, Debug, Eq, Hash, PartialEq)] +struct RepoRef { + owner: Did, + rkey: Rkey, +} + +impl RepoRef { + fn new(owner: Did, rkey: Rkey) -> Self { + Self { owner, rkey } + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum Resolution { + Mapped(Did), + NoRepoDid, + Unresolvable, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +enum AuthoritativeResolution { + Mapped(Did), + NoRepoDid, +} + +impl AuthoritativeResolution { + fn from_repo_did(repo_did: Option>) -> Self { + match repo_did { + Some(did) => Self::Mapped(did), + None => Self::NoRepoDid, + } + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +enum CacheEntry { + Authoritative(AuthoritativeResolution), + Provisional(Resolution), +} + +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, + } + } +} + +pub struct RepoIdResolver { + cache: SccMap, + client: Option, +} + +impl RepoIdResolver { + pub fn with_slingshot(client: SlingshotClient) -> Self { + Self { + cache: SccMap::new(), + client: Some(client), + } + } + + pub fn detached() -> Self { + Self { + cache: SccMap::new(), + client: None, + } + } + + pub async fn observe( + &self, + owner: Did, + rkey: Rkey, + repo_did: Option>, + ) { + let entry = CacheEntry::Authoritative(AuthoritativeResolution::from_repo_did(repo_did)); + self.cache + .entry_async(RepoRef::new(owner, rkey)) + .await + .and_modify(|existing| *existing = entry.clone()) + .or_insert(entry); + } + + async fn fill_provisional(&self, key: RepoRef, 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); + } + + pub async fn resolve(&self, owner: &Did, rkey: &Rkey) -> Resolution { + let key = RepoRef::new(owner.clone(), rkey.clone()); + if let Some(entry) = self.cache.get_async(&key).await { + return entry.get().clone().into_resolution(); + } + let Some(client) = self.client.as_ref() else { + return Resolution::Unresolvable; + }; + let nsid: Nsid = nsid_static(REPO_COLLECTION); + let provisional = match client.get_record(owner, &nsid, rkey).await { + Ok(body) => match repo_did_from_body(&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(), + "slingshot transient failure during repoDID lookup, will retry", + ); + return Resolution::Unresolvable; + } + }; + self.fill_provisional(key, 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(body: &[u8]) -> Result>, ExtractError> { + match Record::from_json_bytes(REPO_COLLECTION, body)? { + Record::Repo(repo) => Ok(repo.repo_did), + _ => Ok(None), + } +} + +#[cfg(test)] +mod tests { + use super::*; + use jacquard_common::types::did::Did; + use jacquard_common::types::recordkey::Rkey; + + fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() + } + + fn rkey(s: &str) -> Rkey { + Rkey::new_owned(s).unwrap() + } + + #[tokio::test] + async fn observation_with_repo_did_resolves_mapped() { + let resolver = RepoIdResolver::detached(); + resolver + .observe( + did("did:plc:nel"), + rkey("abcabcabcabcz"), + Some(did("did:plc:abalone")), + ) + .await; + let got = resolver + .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) + .await; + assert_eq!(got, Resolution::Mapped(did("did:plc:abalone"))); + } + + #[tokio::test] + async fn observation_without_repo_did_resolves_no_repo_did() { + let resolver = RepoIdResolver::detached(); + resolver + .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None) + .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(); + let got = resolver + .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) + .await; + assert_eq!(got, Resolution::Unresolvable); + } + + #[tokio::test] + async fn observation_overwrites_prior_value() { + let resolver = RepoIdResolver::detached(); + resolver + .observe( + did("did:plc:nel"), + rkey("abcabcabcabcz"), + Some(did("did:plc:abalone")), + ) + .await; + resolver + .observe( + did("did:plc:nel"), + rkey("abcabcabcabcz"), + Some(did("did:plc:uni")), + ) + .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(); + let owner = did("did:plc:nel"); + let key = rkey("abcabcabcabcz"); + resolver + .observe(owner.clone(), key.clone(), Some(did("did:plc:abalone"))) + .await; + resolver + .fill_provisional( + RepoRef::new(owner.clone(), key.clone()), + Resolution::Unresolvable, + ) + .await; + let got = resolver.resolve(&owner, &key).await; + assert_eq!( + got, + Resolution::Mapped(did("did:plc:abalone")), + "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(); + let owner = did("did:plc:nel"); + let key = rkey("abcabcabcabcz"); + resolver.observe(owner.clone(), key.clone(), None).await; + resolver + .fill_provisional( + RepoRef::new(owner.clone(), key.clone()), + Resolution::Mapped(did("did:plc:abalone")), + ) + .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::new(url::Url::parse(&server.uri()).unwrap()).unwrap(); + let resolver = RepoIdResolver::with_slingshot(client); + + 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::new(url::Url::parse(&server.uri()).unwrap()).unwrap(); + let resolver = RepoIdResolver::with_slingshot(client); + + 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::new(url::Url::parse(&server.uri()).unwrap()).unwrap(); + let resolver = RepoIdResolver::with_slingshot(client); + + 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_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::new(url::Url::parse(&server.uri()).unwrap()).unwrap(); + let resolver = RepoIdResolver::with_slingshot(client); + + 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_does_not_cache() { + 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(2) + .mount(&server) + .await; + + let client = SlingshotClient::new(url::Url::parse(&server.uri()).unwrap()).unwrap(); + let resolver = RepoIdResolver::with_slingshot(client); + + 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 firehose_observe_can_demote_provisional() { + let resolver = RepoIdResolver::detached(); + let owner = did("did:plc:nel"); + let key = rkey("abcabcabcabcz"); + resolver + .fill_provisional( + RepoRef::new(owner.clone(), key.clone()), + Resolution::Mapped(did("did:plc:abalone")), + ) + .await; + resolver.observe(owner.clone(), key.clone(), None).await; + let got = resolver.resolve(&owner, &key).await; + assert_eq!( + got, + Resolution::NoRepoDid, + "firehose update is canonical and may legitimately remove repoDID", + ); + } +}