From 5c698066f505b466cb1893bc137b53bc5749aa9e Mon Sep 17 00:00:00 2001 From: Lewis Date: Mon, 11 May 2026 09:51:06 +0300 Subject: [PATCH] feat(resolver): housekeeping bs Lewis: May this revision serve well! --- Cargo.lock | 35 +- Cargo.toml | 4 +- crates/resolver/Cargo.toml | 23 ++ crates/resolver/src/legacy_upgrade.rs | 466 ++++++++++++++++++++++++++ crates/resolver/src/normalize.rs | 155 +++++++++ 5 files changed, 665 insertions(+), 18 deletions(-) create mode 100644 crates/resolver/Cargo.toml create mode 100644 crates/resolver/src/legacy_upgrade.rs create mode 100644 crates/resolver/src/normalize.rs diff --git a/Cargo.lock b/Cargo.lock index 1c25922..862bd51 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -330,7 +330,6 @@ dependencies = [ "bobbin-types", "jacquard-common", "lasso", - "roaring", "scc", "thiserror 2.0.18", "tokio", @@ -342,6 +341,7 @@ version = "0.0.1" dependencies = [ "bobbin-edge-index", "bobbin-record-lru", + "bobbin-resolver", "bobbin-runtime", "bobbin-slingshot-client", "bobbin-types", @@ -388,6 +388,22 @@ dependencies = [ "quick_cache", ] +[[package]] +name = "bobbin-resolver" +version = "0.0.1" +dependencies = [ + "bobbin-runtime", + "bobbin-slingshot-client", + "bobbin-types", + "jacquard-common", + "scc", + "serde_json", + "tokio", + "tracing", + "url", + "wiremock", +] + [[package]] name = "bobbin-runtime" version = "0.0.1" @@ -489,6 +505,7 @@ dependencies = [ "bobbin-edge-index", "bobbin-knot-proxy", "bobbin-record-lru", + "bobbin-resolver", "bobbin-runtime", "bobbin-search", "bobbin-slingshot-client", @@ -555,12 +572,6 @@ version = "3.20.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5d20789868f4b01b2f2caec9f5c4e0213b41e3e5702a50157d699ae31ced2fcb" -[[package]] -name = "bytemuck" -version = "1.25.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c8efb64bd706a16a1bdde310ae86b351e4d21550d98d056f22f8a7f7a2183fec" - [[package]] name = "byteorder" version = "1.5.0" @@ -2986,16 +2997,6 @@ dependencies = [ "windows-sys 0.52.0", ] -[[package]] -name = "roaring" -version = "0.11.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1dedc5658c6ecb3bdb5ef5f3295bb9253f42dcf3fd1402c03f6b1f7659c3c4a9" -dependencies = [ - "bytemuck", - "byteorder", -] - [[package]] name = "rust-stemmers" version = "1.2.0" diff --git a/Cargo.toml b/Cargo.toml index 9308b74..edbe6ec 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,6 +6,7 @@ members = [ "crates/types", "crates/edge-index", "crates/ingest", + "crates/resolver", "crates/slingshot-client", "crates/record-lru", "crates/knot-proxy", @@ -17,13 +18,14 @@ members = [ [workspace.package] version = "0.0.1" edition = "2024" -license = "AGPL-3.0-or-later" +license = "MIT" rust-version = "1.95" [workspace.dependencies] bobbin-types = { path = "crates/types" } bobbin-edge-index = { path = "crates/edge-index" } bobbin-ingest = { path = "crates/ingest" } +bobbin-resolver = { path = "crates/resolver" } bobbin-slingshot-client = { path = "crates/slingshot-client" } bobbin-record-lru = { path = "crates/record-lru" } bobbin-knot-proxy = { path = "crates/knot-proxy" } diff --git a/crates/resolver/Cargo.toml b/crates/resolver/Cargo.toml new file mode 100644 index 0000000..f0e2fbf --- /dev/null +++ b/crates/resolver/Cargo.toml @@ -0,0 +1,23 @@ +[package] +name = "bobbin-resolver" +version.workspace = true +edition.workspace = true +license.workspace = true +rust-version.workspace = true + +[dependencies] +bobbin-types = { workspace = true } +bobbin-runtime = { workspace = true } +bobbin-slingshot-client = { workspace = true } +jacquard-common = { workspace = true } + +scc = { workspace = true } +serde_json = { workspace = true } +tokio = { workspace = true } +tracing = { workspace = true } + +[dev-dependencies] +tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } +wiremock = { workspace = true } +url = { workspace = true } +serde_json = { workspace = true } diff --git a/crates/resolver/src/legacy_upgrade.rs b/crates/resolver/src/legacy_upgrade.rs new file mode 100644 index 0000000..fb2374a --- /dev/null +++ b/crates/resolver/src/legacy_upgrade.rs @@ -0,0 +1,466 @@ +use bobbin_types::edges::{ExtractError, Record}; +use bobbin_types::legacy::{ + LegacyCollaborator, LegacyIssue, LegacyPull, LegacyRecord, LegacyRefUpdate, LegacySource, + LegacyStar, LegacyTarget, +}; +use bobbin_types::sh_tangled::feed::star::{Repo as StarRepo, Star, StarString, StarSubject}; +use bobbin_types::sh_tangled::git::ref_update::RefUpdate; +use bobbin_types::sh_tangled::repo::collaborator::Collaborator; +use bobbin_types::sh_tangled::repo::issue::Issue; +use bobbin_types::sh_tangled::repo::pull::{Pull, Source, Target}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::string::AtUri; + +use crate::{RepoIdResolver, Resolution}; +use crate::normalize::{is_repo_at_uri, resolve_repo_uri}; +use jacquard_common::IntoStatic; +use jacquard_common::types::ident::AtIdentifier; +use jacquard_common::types::recordkey::Rkey; + +#[derive(Debug)] +pub enum DecodedRecord { + Canon(Record), + Legacy(LegacyRecord), +} + +impl DecodedRecord { + pub fn try_decode(nsid: &str, bytes: &[u8]) -> Result { + match Record::from_json_bytes(nsid, bytes) { + Ok(record) => Ok(Self::Canon(record)), + Err(canon_err) => match LegacyRecord::from_json_bytes(nsid, bytes) { + Ok(legacy) => Ok(Self::Legacy(legacy)), + Err(_) => Err(canon_err), + }, + } + } +} + +async fn upgrade_repo_did( + resolver: &RepoIdResolver, + at_uri: Option>, + explicit_did: Option>, +) -> Option> { + if let Some(d) = explicit_did { + return Some(d); + } + let uri = at_uri?; + resolve_repo_uri(resolver, &uri).await +} + +pub async fn upgrade_wire_bytes( + nsid: &str, + bytes: &[u8], + resolver: &RepoIdResolver, +) -> Result, ExtractError> { + let legacy = LegacyRecord::from_json_bytes(nsid, bytes)?; + let Some(canon) = upgrade(legacy, resolver).await else { + return Err(ExtractError::UnknownCollection(alloc::format!( + "{nsid}: legacy upgrade failed" + ))); + }; + serialize_canon_variant(&canon).map_err(ExtractError::DecodeJson) +} + +pub async fn decode_canon_or_upgrade_bytes<'a>( + nsid: &str, + bytes: &'a [u8], + resolver: &RepoIdResolver, +) -> Result<(Record, alloc::borrow::Cow<'a, [u8]>), ExtractError> { + let decoded = DecodedRecord::try_decode(nsid, bytes)?; + match decoded { + DecodedRecord::Canon(r) => Ok((r, alloc::borrow::Cow::Borrowed(bytes))), + DecodedRecord::Legacy(legacy) => { + let Some(canon) = upgrade(legacy, resolver).await else { + return Err(ExtractError::UnknownCollection(alloc::format!( + "{nsid}: legacy upgrade failed" + ))); + }; + let canon_bytes = + serialize_canon_variant(&canon).map_err(ExtractError::DecodeJson)?; + Ok((canon, alloc::borrow::Cow::Owned(canon_bytes))) + } + } +} + +fn serialize_canon_variant(record: &Record) -> Result, serde_json::Error> { + match record { + Record::Issue(r) => serde_json::to_vec(r), + Record::Pull(r) => serde_json::to_vec(r), + Record::Collaborator(r) => serde_json::to_vec(r), + Record::RefUpdate(r) => serde_json::to_vec(r), + Record::Star(r) => serde_json::to_vec(r), + _ => unreachable!("upgrade only produces Issue/Pull/Collaborator/RefUpdate/Star"), + } +} + +pub async fn upgrade( + legacy: LegacyRecord, + resolver: &RepoIdResolver, +) -> Option { + match legacy { + LegacyRecord::Issue(l) => upgrade_issue(l, resolver).await.map(Record::Issue), + LegacyRecord::Pull(l) => upgrade_pull(l, resolver).await.map(Record::Pull), + LegacyRecord::Collaborator(l) => upgrade_collaborator(l, resolver) + .await + .map(Record::Collaborator), + LegacyRecord::RefUpdate(l) => Some(Record::RefUpdate(upgrade_ref_update(l))), + LegacyRecord::Star(l) => upgrade_star(l, resolver).await.map(Record::Star), + } +} + +async fn upgrade_issue( + l: LegacyIssue, + resolver: &RepoIdResolver, +) -> Option> { + let repo = upgrade_repo_did(resolver, l.repo, l.repo_did).await?; + Some(Issue { + created_at: l.created_at, + body: l.body, + mentions: l.mentions, + references: l.references, + repo, + title: l.title, + extra_data: l.extra_data, + }) +} + +async fn upgrade_target( + l: LegacyTarget, + resolver: &RepoIdResolver, +) -> Option> { + let repo = upgrade_repo_did(resolver, l.repo, l.repo_did).await?; + Some(Target { + branch: l.branch, + repo, + extra_data: None, + }) +} + +async fn upgrade_source( + l: LegacySource, + resolver: &RepoIdResolver, +) -> Source { + let repo = upgrade_repo_did(resolver, l.repo, l.repo_did).await; + Source { + branch: l.branch, + repo, + extra_data: None, + } +} + +async fn upgrade_pull( + l: LegacyPull, + resolver: &RepoIdResolver, +) -> Option> { + let target = upgrade_target(l.target, resolver).await?; + let source = match l.source { + Some(s) => Some(upgrade_source(s, resolver).await), + None => None, + }; + Some(Pull { + created_at: l.created_at, + body: l.body, + dependent_on: l.dependent_on, + mentions: l.mentions, + references: l.references, + rounds: l.rounds, + source, + target, + title: l.title, + extra_data: l.extra_data, + }) +} + +async fn upgrade_collaborator( + l: LegacyCollaborator, + resolver: &RepoIdResolver, +) -> Option> { + let repo = upgrade_repo_did(resolver, l.repo, l.repo_did).await?; + Some(Collaborator { + created_at: l.created_at, + repo, + subject: l.subject, + extra_data: l.extra_data, + }) +} + +fn upgrade_ref_update(l: LegacyRefUpdate) -> RefUpdate { + RefUpdate { + committer_did: l.committer_did, + meta: l.meta, + new_sha: l.new_sha, + old_sha: l.old_sha, + owner_did: l.owner_did, + r#ref: l.r#ref, + repo: l.repo_did, + extra_data: l.extra_data, + } +} + +async fn upgrade_star( + l: LegacyStar, + resolver: &RepoIdResolver, +) -> Option> { + let subject = if let Some(did) = l.subject_did { + StarSubject::Repo(alloc::boxed::Box::new(StarRepo { + did, + extra_data: None, + })) + } else { + let uri = l.subject?; + let resolved = if is_repo_at_uri(&uri) { + cached_repo_did(resolver, &uri).await + } else { + None + }; + match resolved { + Some(did) => StarSubject::Repo(alloc::boxed::Box::new(StarRepo { + did, + extra_data: None, + })), + None => StarSubject::String(alloc::boxed::Box::new(StarString { + uri, + extra_data: None, + })), + } + }; + Some(Star { + created_at: l.created_at, + subject, + extra_data: l.extra_data, + }) +} + +async fn cached_repo_did( + resolver: &RepoIdResolver, + uri: &jacquard_common::types::string::AtUri, +) -> Option> { + let owner = match uri.authority() { + AtIdentifier::Did(d) => d.clone().into_static(), + AtIdentifier::Handle(_) => return None, + }; + let rkey: Rkey = uri.rkey()?.clone().into_static(); + match resolver.cached_resolution(&owner, &rkey).await? { + Resolution::Mapped(did) => Some(did), + Resolution::NoRepoDid | Resolution::Unresolvable => None, + } +} + +extern crate alloc; + +#[cfg(test)] +mod tests { + use super::*; + use crate::RepoIdResolver; + use bobbin_runtime::RuntimeHasher; + use bobbin_types::edges::Record; + use jacquard_common::DefaultStr; + 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() + } + + #[test] + fn legacy_decode_routes_through_try_decode_for_known_nsids() { + let json = br#"{"$type":"sh.tangled.repo.issue","repoDid":"did:plc:abalone","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; + let decoded = DecodedRecord::try_decode("sh.tangled.repo.issue", json) + .expect("legacy issue must decode"); + assert!(matches!(decoded, DecodedRecord::Legacy(LegacyRecord::Issue(_)))); + } + + #[test] + fn canon_decode_wins_when_wire_matches_new_shape() { + let json = br#"{"$type":"sh.tangled.repo.issue","repo":"did:plc:abalone","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; + let decoded = DecodedRecord::try_decode("sh.tangled.repo.issue", json) + .expect("canon issue must decode"); + match decoded { + DecodedRecord::Canon(Record::Issue(i)) => { + assert_eq!(i.repo.as_ref(), "did:plc:abalone") + } + other => panic!("expected canon issue, got {other:?}"), + } + } + + #[test] + fn legacy_decode_passes_through_for_unaffected_nsids() { + let json = br#"{"$type":"sh.tangled.graph.follow","subject":"did:plc:bailey","createdAt":"2026-05-01T00:00:00Z"}"#; + let decoded = DecodedRecord::try_decode("sh.tangled.graph.follow", json) + .expect("follow has no legacy form, must decode canon"); + assert!(matches!(decoded, DecodedRecord::Canon(Record::Follow(_)))); + } + + #[tokio::test] + async fn upgrade_issue_uses_repo_did_directly() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.repo.issue","repoDid":"did:plc:abalone","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; + let legacy = + LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::Issue(i) => assert_eq!(i.repo, did("did:plc:abalone")), + other => panic!("expected canon issue, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_issue_resolves_repo_uri_via_observed_resolver() { + 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:abalone"))) + .await; + let json = br#"{"$type":"sh.tangled.repo.issue","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; + let legacy = + LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::Issue(i) => assert_eq!(i.repo, did("did:plc:abalone")), + other => panic!("expected canon issue, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_issue_drops_when_resolver_cannot_map() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.repo.issue","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; + let legacy = + LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); + assert!( + upgrade(legacy, &resolver).await.is_none(), + "no resolver entry and no repoDid means the canon Did cannot be constructed", + ); + } + + #[tokio::test] + async fn upgrade_pull_propagates_target_resolution() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repoDid":"did:plc:abalone"}}"#; + let legacy = + LegacyRecord::from_json_bytes("sh.tangled.repo.pull", json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::Pull(p) => { + assert_eq!(p.target.repo, did("did:plc:abalone")); + assert!(p.source.is_none()); + } + other => panic!("expected canon pull, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_pull_source_repo_resolution_is_independent_of_target() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repoDid":"did:plc:abalone"},"source":{"branch":"feat","repo":"at://did:plc:nel/sh.tangled.repo/missing"}}"#; + let legacy = + LegacyRecord::from_json_bytes("sh.tangled.repo.pull", json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::Pull(p) => { + assert_eq!(p.target.repo, did("did:plc:abalone")); + let source = p.source.expect("source struct retained"); + assert_eq!(source.branch.as_str(), "feat"); + assert!( + source.repo.is_none(), + "unresolvable source repo at-uri leaves the source.repo None rather than dropping the whole pull", + ); + } + other => panic!("expected canon pull, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_ref_update_renames_repo_did_to_repo() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.git.refUpdate","ref":"refs/heads/main","committerDid":"did:plc:olaren","repoDid":"did:plc:abalone","oldSha":"0000000000000000000000000000000000000000","newSha":"1111111111111111111111111111111111111111","meta":{"isDefaultRef":true,"commitCount":{}}}"#; + let legacy = + LegacyRecord::from_json_bytes("sh.tangled.git.refUpdate", json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::RefUpdate(r) => assert_eq!(r.repo, did("did:plc:abalone")), + other => panic!("expected canon ref update, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_star_prefers_subject_did_over_subject_uri() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.feed.star","createdAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:nel/sh.tangled.string/k1","subjectDid":"did:plc:abalone"}"#; + let legacy = LegacyRecord::from_json_bytes("sh.tangled.feed.star", json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::Star(s) => match s.subject { + StarSubject::Repo(r) => assert_eq!(r.did, did("did:plc:abalone")), + StarSubject::String(_) => panic!("subjectDid must win"), + }, + other => panic!("expected canon star, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_star_falls_back_to_string_when_repo_uri_not_in_cache() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.feed.star","createdAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"}"#; + let legacy = LegacyRecord::from_json_bytes("sh.tangled.feed.star", json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::Star(s) => match s.subject { + StarSubject::String(s) => assert_eq!( + s.uri.as_ref(), + "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", + "cache-miss on repo uri preserves the uri under the #string variant for later normalization", + ), + StarSubject::Repo(_) => panic!("cold cache must not upgrade to Repo variant"), + }, + other => panic!("expected canon star, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_star_uses_cached_repo_did_when_observed() { + 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:abalone"))) + .await; + let json = br#"{"$type":"sh.tangled.feed.star","createdAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"}"#; + let legacy = LegacyRecord::from_json_bytes("sh.tangled.feed.star", json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::Star(s) => match s.subject { + StarSubject::Repo(r) => assert_eq!(r.did, did("did:plc:abalone")), + StarSubject::String(_) => panic!("observed cache must upgrade to Repo variant"), + }, + other => panic!("expected canon star, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_collaborator_requires_repo_did() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let with_did = br#"{"$type":"sh.tangled.repo.collaborator","createdAt":"2026-05-01T00:00:00Z","subject":"did:plc:lyna","repoDid":"did:plc:abalone"}"#; + let canon = upgrade( + LegacyRecord::from_json_bytes("sh.tangled.repo.collaborator", with_did).expect("decode"), + &resolver, + ) + .await + .expect("upgrade"); + match canon { + Record::Collaborator(c) => assert_eq!(c.repo, did("did:plc:abalone")), + other => panic!("expected canon collaborator, got {other:?}"), + } + + let no_resolution = br#"{"$type":"sh.tangled.repo.collaborator","createdAt":"2026-05-01T00:00:00Z","subject":"did:plc:lyna","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"}"#; + let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.collaborator", no_resolution) + .expect("decode"); + assert!(upgrade(legacy, &resolver).await.is_none()); + } +} diff --git a/crates/resolver/src/normalize.rs b/crates/resolver/src/normalize.rs new file mode 100644 index 0000000..9a7ad30 --- /dev/null +++ b/crates/resolver/src/normalize.rs @@ -0,0 +1,155 @@ +use crate::{RepoIdResolver, Resolution}; +use bobbin_types::search::SearchableRecord; +use bobbin_types::sh_tangled::repo::artifact::Artifact; +use jacquard_common::DefaultStr; +use jacquard_common::IntoStatic; +use jacquard_common::types::did::Did; +use jacquard_common::types::ident::AtIdentifier; +use jacquard_common::types::recordkey::Rkey; +use jacquard_common::types::string::AtUri; + +pub(crate) const REPO_COLLECTION: &str = "sh.tangled.repo"; + +pub trait NormalizeRepoRefs: Sized { + fn normalize( + self, + resolver: &RepoIdResolver, + ) -> impl std::future::Future> + Send; +} + +pub(crate) fn is_repo_at_uri(uri: &AtUri) -> bool { + uri.collection() + .map(|c| c.as_ref() == REPO_COLLECTION) + .unwrap_or(false) +} + +pub(crate) async fn resolve_repo_uri( + resolver: &RepoIdResolver, + uri: &AtUri, +) -> Option> { + if !is_repo_at_uri(uri) { + return None; + } + let owner = match uri.authority() { + AtIdentifier::Did(d) => d.clone().into_static(), + AtIdentifier::Handle(_) => return None, + }; + let rkey: Rkey = uri.rkey()?.clone().into_static(); + match resolver.resolve(&owner, &rkey).await { + Resolution::Mapped(did) => Some(did), + Resolution::NoRepoDid | Resolution::Unresolvable => None, + } +} + +async fn fill_repo_did( + resolver: &RepoIdResolver, + uri_field: &mut Option>, + did_field: &mut Option>, +) -> bool { + if did_field.is_some() { + *uri_field = None; + return true; + } + let Some(uri) = uri_field.as_ref() else { + return false; + }; + let Some(did) = resolve_repo_uri(resolver, uri).await else { + return false; + }; + *did_field = Some(did); + *uri_field = None; + true +} + +impl NormalizeRepoRefs for Artifact { + async fn normalize(mut self, resolver: &RepoIdResolver) -> Option { + if !fill_repo_did(resolver, &mut self.repo, &mut self.repo_did).await { + return None; + } + Some(self) + } +} + +impl NormalizeRepoRefs for bobbin_types::sh_tangled::pipeline::Pipeline { + async fn normalize(mut self, resolver: &RepoIdResolver) -> Option { + let trig = &mut self.trigger_metadata.repo; + if trig.repo_did.is_some() { + trig.repo = None; + return Some(self); + } + let raw = trig.repo.as_deref()?; + let parsed = AtUri::::new_owned(raw).ok()?; + let did = resolve_repo_uri(resolver, &parsed).await?; + trig.repo_did = Some(did); + trig.repo = None; + Some(self) + } +} + +impl NormalizeRepoRefs for SearchableRecord { + async fn normalize(self, _resolver: &RepoIdResolver) -> Option { + Some(self) + } +} + +macro_rules! identity_normalize { + ($($t:ty),+ $(,)?) => { + $( + impl NormalizeRepoRefs for $t { + async fn normalize(self, _resolver: &RepoIdResolver) -> Option { + Some(self) + } + } + )+ + }; +} + +use bobbin_types::sh_tangled::actor::profile::Profile; +use bobbin_types::sh_tangled::feed::reaction::Reaction; +use bobbin_types::sh_tangled::feed::star::Star; +use bobbin_types::sh_tangled::git::ref_update::RefUpdate; +use bobbin_types::sh_tangled::graph::follow::Follow; +use bobbin_types::sh_tangled::graph::vouch::Vouch; +use bobbin_types::sh_tangled::knot::Knot; +use bobbin_types::sh_tangled::knot::member::Member as KnotMember; +use bobbin_types::sh_tangled::label::definition::Definition as LabelDefinition; +use bobbin_types::sh_tangled::label::op::Op as LabelOp; +use bobbin_types::sh_tangled::pipeline::status::Status as PipelineStatus; +use bobbin_types::sh_tangled::public_key::PublicKey; +use bobbin_types::sh_tangled::repo::Repo; +use bobbin_types::sh_tangled::repo::collaborator::Collaborator; +use bobbin_types::sh_tangled::repo::issue::Issue; +use bobbin_types::sh_tangled::repo::issue::comment::Comment as IssueComment; +use bobbin_types::sh_tangled::repo::issue::state::State as IssueState; +use bobbin_types::sh_tangled::repo::pull::Pull; +use bobbin_types::sh_tangled::repo::pull::comment::Comment as PullComment; +use bobbin_types::sh_tangled::repo::pull::status::Status as PullStatus; +use bobbin_types::sh_tangled::spindle::Spindle; +use bobbin_types::sh_tangled::spindle::member::Member as SpindleMember; +use bobbin_types::sh_tangled::string::TangledString; + +identity_normalize!( + Profile, + Reaction, + Star, + RefUpdate, + Follow, + Vouch, + Knot, + KnotMember, + LabelDefinition, + LabelOp, + PipelineStatus, + PublicKey, + Repo, + Collaborator, + Issue, + IssueComment, + IssueState, + Pull, + PullComment, + PullStatus, + Spindle, + SpindleMember, + TangledString, +); -- 2.51.2