diff --git a/crates/bobbin/src/main.rs b/crates/bobbin/src/main.rs index e50c17b..d14898f 100644 --- a/crates/bobbin/src/main.rs +++ b/crates/bobbin/src/main.rs @@ -4,23 +4,15 @@ use std::sync::Arc; use anyhow::{Context, anyhow}; use bobbin_edge_index::{CoverageWatch, EdgeStore, HydrantCursor}; -use bobbin_ingest::{IngestConfig, RepoDidResolver, ResolveError, run as run_ingest}; +use bobbin_ingest::{IngestConfig, run as run_ingest}; use bobbin_knot_proxy::{KnotProxy, KnotProxyConfig}; use bobbin_record_lru::{CacheCapacity, LruRecordStore, RecordStore}; use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex}; -use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; -use bobbin_types::record::RecordBody; +use bobbin_slingshot_client::SlingshotClient; use bobbin_xrpc::{AppState, router}; -use jacquard_common::DefaultStr; -use jacquard_common::types::did::Did; -use jacquard_common::types::ident::AtIdentifier; -use jacquard_common::types::nsid::Nsid; -use jacquard_common::types::string::AtUri; use tracing_subscriber::EnvFilter; use url::Url; -const REPO_NSID: &str = "sh.tangled.repo"; - #[tokio::main] async fn main() -> anyhow::Result<()> { tracing_subscriber::fmt() @@ -61,13 +53,8 @@ async fn main() -> anyhow::Result<()> { hydrant_base: Url::parse(&hydrant_url)?, start_cursor: HydrantCursor::new(start_cursor), }; - let resolver = Arc::new(SlingshotRepoDidResolver { - slingshot: slingshot.clone(), - records: records.clone(), - }); let ingest_edges = edges.clone(); let ingest_coverage = coverage.clone(); - let ingest_resolver = resolver.clone(); let ingest_search = search.clone(); let ingest_records = records.clone(); let ingest_handle = tokio::spawn(async move { @@ -75,7 +62,6 @@ async fn main() -> anyhow::Result<()> { ingest_cfg, ingest_edges, ingest_coverage, - ingest_resolver, ingest_search, ingest_records, ) @@ -102,107 +88,3 @@ async fn main() -> anyhow::Result<()> { }, } } - -struct SlingshotRepoDidResolver { - slingshot: SlingshotClient, - records: Arc, -} - -impl RepoDidResolver for SlingshotRepoDidResolver { - async fn resolve( - &self, - repo_aturi: &AtUri, - ) -> Result>, ResolveError> { - let body = match self.records.get(repo_aturi) { - Some(b) => b, - None => match fetch_repo_record(&self.slingshot, repo_aturi).await? { - Some(b) => { - self.records.put(repo_aturi.clone(), b.clone()); - b - } - None => return Ok(None), - }, - }; - Ok(parse_repo_did(&body)) - } -} - -async fn fetch_repo_record( - slingshot: &SlingshotClient, - repo_aturi: &AtUri, -) -> Result>, ResolveError> { - let owner = match repo_aturi.authority() { - AtIdentifier::Did(d) => d, - AtIdentifier::Handle(_) => return Ok(None), - }; - let Some(rkey) = repo_aturi.rkey() else { - return Ok(None); - }; - let collection = Nsid::<&str>::new(REPO_NSID).expect("REPO_NSID literal validates as NSID"); - match slingshot.get_record(&owner, &collection, &rkey).await { - Ok(body) => Ok(Some(body)), - Err(SlingshotError::NotFound) => Ok(None), - Err(e) => Err(ResolveError::Transport(e.to_string())), - } -} - -fn parse_repo_did(body: &RecordBody) -> Option> { - let v: serde_json::Value = serde_json::from_slice(&body.value).ok()?; - let s = v.get("repoDid")?.as_str()?; - Did::::new(jacquard_common::deps::smol_str::SmolStr::new(s)).ok() -} - -#[cfg(test)] -mod tests { - use super::*; - use bytes::Bytes; - use jacquard_common::types::string::Cid; - - fn body_with_value(value: &str) -> RecordBody { - let cid: Cid = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i" - .parse() - .unwrap(); - RecordBody { - uri: AtUri::new_owned("at://did:plc:abalone/sh.tangled.repo/r1").unwrap(), - cid, - value: Bytes::copy_from_slice(value.as_bytes()), - } - } - - #[test] - fn parses_repo_did_when_field_present() { - let body = body_with_value( - r#"{"$type":"sh.tangled.repo","repoDid":"did:plc:limpet","name":"abalone","knot":"oyster.cafe","createdAt":"2026-05-01T00:00:00Z"}"#, - ); - assert_eq!( - parse_repo_did(&body).map(|d| d.as_ref().to_owned()), - Some("did:plc:limpet".into()), - ); - } - - #[test] - fn returns_none_when_repo_did_missing() { - let body = body_with_value( - r#"{"$type":"sh.tangled.repo","name":"abalone","knot":"oyster.cafe","createdAt":"2026-05-01T00:00:00Z"}"#, - ); - assert!(parse_repo_did(&body).is_none()); - } - - #[test] - fn returns_none_when_repo_did_is_not_a_string() { - let body = body_with_value(r#"{"$type":"sh.tangled.repo","repoDid":42}"#); - assert!(parse_repo_did(&body).is_none()); - } - - #[test] - fn returns_none_when_value_is_not_json() { - let body = body_with_value("not json at all"); - assert!(parse_repo_did(&body).is_none()); - } - - #[test] - fn returns_none_when_repo_did_fails_did_validation() { - let body = body_with_value(r#"{"repoDid":"definitely not a did"}"#); - assert!(parse_repo_did(&body).is_none()); - } -} diff --git a/crates/ingest/examples/smoke.rs b/crates/ingest/examples/smoke.rs index e7584a3..f1a78a7 100644 --- a/crates/ingest/examples/smoke.rs +++ b/crates/ingest/examples/smoke.rs @@ -2,7 +2,7 @@ use std::sync::Arc; use std::time::Duration; use bobbin_edge_index::{CoverageWatch, EdgeStore}; -use bobbin_ingest::{IngestConfig, NoopResolver, run}; +use bobbin_ingest::{IngestConfig, run}; use bobbin_record_lru::{NoopRecordStore, RecordStore}; use bobbin_types::search::NoopSearchSink; use futures::stream::{self, StreamExt}; @@ -29,11 +29,10 @@ async fn main() { let cfg = IngestConfig::new(url); let store_runner = store.clone(); let coverage_runner = coverage.clone(); - let resolver = Arc::new(NoopResolver); let search = Arc::new(NoopSearchSink); let records: Arc = Arc::new(NoopRecordStore); let task = tokio::spawn(async move { - let _ = run(cfg, store_runner, coverage_runner, resolver, search, records).await; + let _ = run(cfg, store_runner, coverage_runner, search, records).await; }); stream::iter(0..ticks) diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index a6204b0..6ba2f46 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -3,13 +3,10 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH}; use bobbin_edge_index::{Coverage, CoverageWatch, EdgeStore, HydrantCursor, PromotionSignal}; use bobbin_record_lru::RecordStore; -use bobbin_types::edges::{Edge, ExtractError, Record}; +use bobbin_types::edges::{ExtractError, Record}; use bobbin_types::search::{SearchSink, SearchableRecord}; -use futures::SinkExt; -use futures::stream::{self, StreamExt, TryStreamExt}; +use futures::{SinkExt, StreamExt}; use jacquard_common::DefaultStr; -use jacquard_common::types::did::Did; -use jacquard_common::types::ident::AtIdentifier; use jacquard_common::types::string::{AtStrError, AtUri}; use thiserror::Error; use tokio::time::{Instant, MissedTickBehavior, interval}; @@ -20,34 +17,6 @@ use url::Url; mod frame; pub use frame::{FrameKind, HydrantFrame, RecordAction, RecordFrame}; -const REPO_NSID: &str = "sh.tangled.repo"; - -#[derive(Debug, Error)] -pub enum ResolveError { - #[error("transport: {0}")] - Transport(String), - #[error("invalid did from upstream: {0}")] - InvalidDid(#[from] AtStrError), -} - -pub trait RepoDidResolver: Send + Sync { - fn resolve( - &self, - repo_aturi: &AtUri, - ) -> impl std::future::Future>, ResolveError>> + Send; -} - -pub struct NoopResolver; - -impl RepoDidResolver for NoopResolver { - async fn resolve( - &self, - _repo_aturi: &AtUri, - ) -> Result>, ResolveError> { - Ok(None) - } -} - const TANGLED_PREFIX: &str = "sh.tangled."; const RECONNECT_INITIAL_DELAY: Duration = Duration::from_millis(500); const RECONNECT_MAX_DELAY: Duration = Duration::from_secs(30); @@ -104,8 +73,6 @@ pub enum IngestError { InvalidAtUri(#[from] AtStrError), #[error("record extraction: {0}")] Extract(#[from] ExtractError), - #[error("repo-did resolve: {0}")] - Resolve(#[from] ResolveError), #[error("hydrant did not respond to ping within {0:?}")] PongTimeout(Duration), } @@ -122,11 +89,10 @@ struct SessionEnd { error: Option, } -pub async fn run( +pub async fn run( config: IngestConfig, store: Arc, coverage: Arc, - resolver: Arc, search: Arc, records: Arc, ) -> Result<(), IngestError> { @@ -138,7 +104,6 @@ pub async fn run( cursor, store.clone(), coverage.clone(), - resolver.clone(), search.clone(), records.clone(), ) @@ -179,12 +144,11 @@ fn jittered(base: Duration) -> Duration { base + Duration::from_millis(entropy % cap_ms) } -async fn run_session( +async fn run_session( config: &IngestConfig, cursor: HydrantCursor, store: Arc, coverage: Arc, - resolver: Arc, search: Arc, records: Arc, ) -> SessionEnd { @@ -216,7 +180,6 @@ async fn run_session( let (frame_tx, mut frame_rx) = tokio::sync::mpsc::channel::(FRAME_CHANNEL_DEPTH); let processor_store = store.clone(); let processor_coverage = coverage.clone(); - let processor_resolver = resolver.clone(); let processor_search = search.clone(); let processor_records = records.clone(); let processor = tokio::spawn(async move { @@ -225,7 +188,6 @@ async fn run_session( frame, &processor_store, &processor_coverage, - &*processor_resolver, &*processor_search, &*processor_records, ) @@ -301,11 +263,10 @@ async fn wait_until(deadline: Option) { } } -async fn handle_frame( +async fn handle_frame( frame: HydrantFrame, store: &EdgeStore, coverage: &CoverageWatch, - resolver: &R, search: &S, records: &dyn RecordStore, ) { @@ -325,7 +286,7 @@ async fn handle_frame( if !record.collection.as_ref().starts_with(TANGLED_PREFIX) { return; } - if let Err(err) = apply_record(record, store, resolver, search, records).await { + if let Err(err) = apply_record(record, store, search, records).await { warn!(?err, "dropped record"); } } @@ -352,10 +313,9 @@ fn now_micros() -> u64 { .as_micros() as u64 } -async fn apply_record( +async fn apply_record( record: RecordFrame, store: &EdgeStore, - resolver: &R, search: &S, records: &dyn RecordStore, ) -> Result<(), IngestError> { @@ -376,12 +336,7 @@ async fn apply_record( Err(e) => return Err(e.into()), }; let edges = parsed.extract_edges(&source)?; - let normalized: Vec = stream::iter(edges) - .map(Ok::<_, IngestError>) - .try_filter_map(|edge| normalize_repo_subject(edge, resolver)) - .try_collect() - .await?; - store.upsert_source(&source, normalized); + store.upsert_source(&source, edges); if let Some(searchable) = SearchableRecord::try_from_record(parsed) { search.upsert(searchable.to_search_doc(&source)).await; } @@ -398,37 +353,6 @@ async fn apply_record( Ok(()) } -async fn normalize_repo_subject( - edge: Edge, - resolver: &R, -) -> Result, IngestError> { - if !is_repo_aturi(&edge.subject) { - return Ok(Some(edge)); - } - match resolver.resolve(&edge.subject).await? { - Some(did) => Ok(Some(Edge { - subject: AtUri::new_owned(format!("at://{}", did.as_ref()))?, - ..edge - })), - None => { - debug!( - subject = %edge.subject.as_ref(), - "dropping edge: repoDid resolution returned no record", - ); - Ok(None) - } - } -} - -fn is_repo_aturi(uri: &AtUri) -> bool { - let Some(collection) = uri.collection() else { - return false; - }; - collection.as_ref() == REPO_NSID - && matches!(uri.authority(), AtIdentifier::Did(_)) - && uri.rkey().is_some() -} - fn build_source_uri(r: &RecordFrame) -> Result, IngestError> { Ok(AtUri::from_parts_owned( r.did.as_ref(), @@ -476,7 +400,7 @@ mod tests { "record": {"$type": "app.bsky.feed.post", "text": "hi"} } })); - handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; + handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().events_processed(), 1); } @@ -501,7 +425,7 @@ mod tests { } } })); - handle_frame(create, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; + handle_frame(create, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; let key = bobbin_types::ids::EdgeKey::new( Nsid::new_static("sh.tangled.feed.star").unwrap(), AtUri::new_owned("at://did:plc:abalone").unwrap(), @@ -521,7 +445,7 @@ mod tests { "record": null } })); - handle_frame(delete, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; + handle_frame(delete, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; assert_eq!(store.count(&key), 0); } @@ -551,7 +475,6 @@ mod tests { mk("did:plc:abalone", 1), &store, &cov, - &NoopResolver, &NoopSearchSink, &NoopRecordStore, ) @@ -560,7 +483,6 @@ mod tests { mk("did:plc:uni", 2), &store, &cov, - &NoopResolver, &NoopSearchSink, &NoopRecordStore, ) @@ -598,7 +520,7 @@ mod tests { } } })); - handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; + handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; assert!(cov.snapshot().is_ready()); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(99)); } @@ -624,7 +546,7 @@ mod tests { } } })); - handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; + handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; assert!(!cov.snapshot().is_ready()); assert!(matches!(cov.snapshot(), Coverage::Warming { .. })); } @@ -684,7 +606,7 @@ mod tests { "handle": "olaren.dev" } })); - handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; + handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(5)); assert!(!cov.snapshot().is_ready()); @@ -698,7 +620,7 @@ mod tests { "type": "account", "account": {"did": "did:plc:olaren", "active": true} })); - handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; + handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(6)); } @@ -708,7 +630,7 @@ mod tests { let (store, cov) = fresh(); let frame: HydrantFrame = parse_frame(json!({"id": 8, "type": "future_event"})); assert_eq!(frame.kind, FrameKind::Other); - handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; + handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(8)); } diff --git a/crates/types/src/edges.rs b/crates/types/src/edges.rs index 4ee0919..a2bf5a9 100644 --- a/crates/types/src/edges.rs +++ b/crates/types/src/edges.rs @@ -195,19 +195,42 @@ fn aturi_to_owned>(uri: &AtUri) -> Result>( - repo: &Option>, - repo_did: &Option>, +fn did_subject>( + uri: &Option>, + did: &Option>, ) -> Result>, AtStrError> { - if let Some(did) = repo_did { - Ok(Some(did_to_aturi(did)?)) - } else if let Some(uri) = repo { - Ok(Some(aturi_to_owned(uri)?)) - } else { - Ok(None) + match (did.as_ref(), uri.as_ref()) { + (Some(d), _) => did_to_aturi(d).map(Some), + (None, Some(u)) => uri_authority_did_aturi(u), + (None, None) => Ok(None), } } +fn star_subject>( + uri: &Option>, + did: &Option>, +) -> Result>, AtStrError> { + match (did.as_ref(), uri.as_ref()) { + (Some(d), _) => did_to_aturi(d).map(Some), + (None, Some(u)) if is_string_at_uri(u) => aturi_to_owned(u).map(Some), + (None, Some(u)) => uri_authority_did_aturi(u), + (None, None) => Ok(None), + } +} + +fn uri_authority_did_aturi>( + uri: &AtUri, +) -> Result>, AtStrError> { + crate::ids::did_from_aturi(uri.as_ref()) + .map(|d| did_to_aturi(&d)) + .transpose() +} + +fn is_string_at_uri>(uri: &AtUri) -> bool { + uri.collection() + .is_some_and(|c| c.as_ref() == "sh.tangled.string") +} + fn one_edge( kind: &'static str, subject: AtUri, @@ -224,7 +247,7 @@ fn star_edges( source: &AtUri, record: &Star, ) -> Result, ExtractError> { - let Some(subject) = repo_subject(&record.subject, &record.subject_did)? else { + let Some(subject) = star_subject(&record.subject, &record.subject_did)? else { return Ok(Vec::new()); }; Ok(one_edge("sh.tangled.feed.star", subject, source)) @@ -303,7 +326,7 @@ fn artifact_edges( source: &AtUri, record: &Artifact, ) -> Result, ExtractError> { - let Some(subject) = repo_subject(&record.repo, &record.repo_did)? else { + let Some(subject) = did_subject(&record.repo, &record.repo_did)? else { return Ok(Vec::new()); }; Ok(one_edge("sh.tangled.repo.artifact", subject, source)) @@ -313,7 +336,7 @@ fn collaborator_edges( source: &AtUri, record: &Collaborator, ) -> Result, ExtractError> { - let Some(subject) = repo_subject(&record.repo, &record.repo_did)? else { + let Some(subject) = did_subject(&record.repo, &record.repo_did)? else { return Ok(Vec::new()); }; Ok(one_edge("sh.tangled.repo.collaborator", subject, source)) @@ -323,7 +346,7 @@ fn issue_edges( source: &AtUri, record: &Issue, ) -> Result, ExtractError> { - let Some(subject) = repo_subject(&record.repo, &record.repo_did)? else { + let Some(subject) = did_subject(&record.repo, &record.repo_did)? else { return Ok(Vec::new()); }; Ok(one_edge("sh.tangled.repo.issue", subject, source)) @@ -355,7 +378,7 @@ fn pull_edges( source: &AtUri, record: &Pull, ) -> Result, ExtractError> { - let Some(subject) = repo_subject(&record.target.repo, &record.target.repo_did)? else { + let Some(subject) = did_subject(&record.target.repo, &record.target.repo_did)? else { return Ok(Vec::new()); }; Ok(one_edge("sh.tangled.repo.pull", subject, source)) @@ -459,7 +482,7 @@ mod tests { } #[test] - fn star_flat_subject_uri_falls_back_when_no_did() { + fn star_string_subject_uri_keeps_full_path() { let edges = extract( "sh.tangled.feed.star", "at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz", @@ -476,6 +499,66 @@ mod tests { ); } + #[test] + fn star_repo_subject_uri_collapses_to_authority_did() { + let edges = extract( + "sh.tangled.feed.star", + "at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz", + json!({ + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://did:plc:abalone/sh.tangled.repo/r1" + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].subject, at("at://did:plc:abalone")); + } + + #[test] + fn star_bare_did_subject_uri_keeps_did_aturi() { + let edges = extract( + "sh.tangled.feed.star", + "at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz", + json!({ + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://did:plc:abalone" + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].subject, at("at://did:plc:abalone")); + } + + #[test] + fn star_handle_authority_subject_emits_no_edge() { + let edges = extract( + "sh.tangled.feed.star", + "at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz", + json!({ + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://oyster.cafe/sh.tangled.repo/r1" + }), + ); + assert!(edges.is_empty()); + } + + #[test] + fn star_subject_did_wins_over_subject_uri() { + let edges = extract( + "sh.tangled.feed.star", + "at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz", + json!({ + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://did:plc:teq/sh.tangled.string/k1", + "subjectDid": "did:plc:abalone" + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].subject, at("at://did:plc:abalone")); + } + #[test] fn follow_uses_subject_did() { let edges = extract( @@ -602,6 +685,203 @@ mod tests { assert!(edges.is_empty()); } + #[test] + fn artifact_uri_only_collapses_to_authority_did() { + let edges = extract( + "sh.tangled.repo.artifact", + "at://did:plc:nel/sh.tangled.repo.artifact/abcabcabcabcz", + artifact_body(Some("at://did:plc:abalone/sh.tangled.repo/r1"), None), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].subject, at("at://did:plc:abalone")); + } + + #[test] + fn artifact_uri_with_handle_authority_emits_no_edge() { + let edges = extract( + "sh.tangled.repo.artifact", + "at://did:plc:nel/sh.tangled.repo.artifact/abcabcabcabcz", + artifact_body(Some("at://oyster.cafe/sh.tangled.repo/r1"), None), + ); + assert!(edges.is_empty()); + } + + #[test] + fn issue_uri_only_collapses_to_authority_did() { + let edges = extract( + "sh.tangled.repo.issue", + "at://did:plc:nel/sh.tangled.repo.issue/abcabcabcabcz", + json!({ + "$type": "sh.tangled.repo.issue", + "repo": "at://did:plc:abalone/sh.tangled.repo/r1", + "title": "bug", + "createdAt": "2026-05-01T00:00:00Z" + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].subject, at("at://did:plc:abalone")); + } + + #[test] + fn pull_uri_only_collapses_to_authority_did() { + let edges = extract( + "sh.tangled.repo.pull", + "at://did:plc:nel/sh.tangled.repo.pull/abcabcabcabcz", + json!({ + "$type": "sh.tangled.repo.pull", + "title": "feature", + "createdAt": "2026-05-01T00:00:00Z", + "rounds": [], + "target": { + "repo": "at://did:plc:abalone/sh.tangled.repo/r1", + "branch": "main" + } + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].subject, at("at://did:plc:abalone")); + } + + #[test] + fn label_op_keys_on_subject_at_uri() { + let issue_uri = "at://did:plc:abalone/sh.tangled.repo.issue/i1"; + let edges = extract( + "sh.tangled.label.op", + "at://did:plc:nel/sh.tangled.label.op/abcabcabcabcz", + json!({ + "$type": "sh.tangled.label.op", + "performedAt": "2026-05-01T00:00:00Z", + "subject": issue_uri, + "add": [{ + "key": "at://did:plc:abalone/sh.tangled.label.definition/bug", + "value": "true" + }], + "delete": [] + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].kind, nsid("sh.tangled.label.op")); + assert_eq!(edges[0].subject, at(issue_uri)); + } + + #[test] + fn label_op_pull_subject_keys_on_pull_uri() { + let pull_uri = "at://did:plc:abalone/sh.tangled.repo.pull/p1"; + let edges = extract( + "sh.tangled.label.op", + "at://did:plc:nel/sh.tangled.label.op/abcabcabcabcz", + json!({ + "$type": "sh.tangled.label.op", + "performedAt": "2026-05-01T00:00:00Z", + "subject": pull_uri, + "add": [], + "delete": [] + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].subject, at(pull_uri)); + } + + #[test] + fn pipeline_keys_on_repo_did_when_present() { + let edges = extract( + "sh.tangled.pipeline", + "at://did:plc:lyna/sh.tangled.pipeline/pl1", + json!({ + "$type": "sh.tangled.pipeline", + "workflows": [], + "triggerMetadata": { + "kind": "manual", + "repo": { + "did": "did:plc:nel", + "repoDid": "did:plc:abalone", + "knot": "oyster.cafe", + "defaultBranch": "main" + } + } + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].kind, nsid("sh.tangled.pipeline")); + assert_eq!(edges[0].subject, at("at://did:plc:abalone")); + } + + #[test] + fn pipeline_falls_back_to_owner_did_when_repo_did_absent() { + let edges = extract( + "sh.tangled.pipeline", + "at://did:plc:lyna/sh.tangled.pipeline/pl1", + json!({ + "$type": "sh.tangled.pipeline", + "workflows": [], + "triggerMetadata": { + "kind": "manual", + "repo": { + "did": "did:plc:nel", + "knot": "oyster.cafe", + "defaultBranch": "main" + } + } + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].kind, nsid("sh.tangled.pipeline")); + assert_eq!(edges[0].subject, at("at://did:plc:nel")); + } + + #[test] + fn pipeline_status_keys_on_pipeline_at_uri() { + let pipeline_uri = "at://did:plc:lyna/sh.tangled.pipeline/pl1"; + let edges = extract( + "sh.tangled.pipeline.status", + "at://did:plc:bailey/sh.tangled.pipeline.status/abcabcabcabcz", + json!({ + "$type": "sh.tangled.pipeline.status", + "createdAt": "2026-05-01T00:00:00Z", + "pipeline": pipeline_uri, + "workflow": pipeline_uri, + "status": "success" + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].kind, nsid("sh.tangled.pipeline.status")); + assert_eq!(edges[0].subject, at(pipeline_uri)); + } + + #[test] + fn knot_member_keys_on_subject_did() { + let edges = extract( + "sh.tangled.knot.member", + "at://did:plc:teq/sh.tangled.knot.member/abcabcabcabcz", + json!({ + "$type": "sh.tangled.knot.member", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "did:plc:nel", + "domain": "oyster.cafe" + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].kind, nsid("sh.tangled.knot.member")); + assert_eq!(edges[0].subject, at("at://did:plc:nel")); + } + + #[test] + fn spindle_member_keys_on_subject_did() { + let edges = extract( + "sh.tangled.spindle.member", + "at://did:plc:teq/sh.tangled.spindle.member/abcabcabcabcz", + json!({ + "$type": "sh.tangled.spindle.member", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "did:plc:olaren", + "instance": "spin.nel.pet" + }), + ); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].kind, nsid("sh.tangled.spindle.member")); + assert_eq!(edges[0].subject, at("at://did:plc:olaren")); + } + #[test] fn repo_self_edge_keys_on_owner_did() { let edges = extract(