diff --git a/crates/xrpc/src/lib.rs b/crates/xrpc/src/lib.rs index b0c8bcc..d973479 100644 --- a/crates/xrpc/src/lib.rs +++ b/crates/xrpc/src/lib.rs @@ -30,12 +30,28 @@ use bobbin_types::search::SearchableRecord; use bobbin_types::sh_tangled::actor::profile::{Profile, ProfileGetRecordOutput, ProfileRecord}; use bobbin_types::sh_tangled::feed::star::{Star, StarRecord}; use bobbin_types::sh_tangled::graph::follow::{Follow, FollowRecord}; +use bobbin_types::sh_tangled::knot::member::{ + Member as KnotMember, MemberRecord as KnotMemberRecord, +}; +use bobbin_types::sh_tangled::label::definition::{ + Definition as LabelDefinition, DefinitionRecord as LabelDefinitionRecord, +}; +use bobbin_types::sh_tangled::label::op::{Op as LabelOp, OpRecord as LabelOpRecord}; +use bobbin_types::sh_tangled::pipeline::status::{ + Status as PipelineStatus, StatusRecord as PipelineStatusRecord, +}; +use bobbin_types::sh_tangled::pipeline::{Pipeline, PipelineRecord}; +use bobbin_types::sh_tangled::repo::artifact::{Artifact, ArtifactRecord}; use bobbin_types::sh_tangled::repo::issue::comment::{ Comment as IssueComment, CommentRecord as IssueCommentRecord, }; use bobbin_types::sh_tangled::repo::issue::{Issue, IssueGetRecordOutput, IssueRecord}; use bobbin_types::sh_tangled::repo::pull::{Pull, PullGetRecordOutput, PullRecord}; use bobbin_types::sh_tangled::repo::{Repo, RepoGetRecordOutput, RepoRecord}; +use bobbin_types::sh_tangled::spindle::member::{ + Member as SpindleMember, MemberRecord as SpindleMemberRecord, +}; +use bobbin_types::sh_tangled::string::{TangledString, TangledStringRecord}; use futures::stream::{self, StreamExt, TryStreamExt}; use jacquard_common::types::did::Did; use jacquard_common::types::ident::AtIdentifier; @@ -101,6 +117,46 @@ pub fn router(state: AppState) -> Router { "/xrpc/sh.tangled.repo.issue.countComments", get(count_issue_comments), ) + .route( + "/xrpc/sh.tangled.label.listDefinitions", + get(list_label_definitions), + ) + .route( + "/xrpc/sh.tangled.label.countDefinitions", + get(count_label_definitions), + ) + .route("/xrpc/sh.tangled.label.listOps", get(list_label_ops)) + .route("/xrpc/sh.tangled.label.countOps", get(count_label_ops)) + .route( + "/xrpc/sh.tangled.pipeline.listPipelines", + get(list_pipelines), + ) + .route( + "/xrpc/sh.tangled.pipeline.countPipelines", + get(count_pipelines), + ) + .route( + "/xrpc/sh.tangled.pipeline.listStatuses", + get(list_pipeline_statuses), + ) + .route( + "/xrpc/sh.tangled.pipeline.countStatuses", + get(count_pipeline_statuses), + ) + .route("/xrpc/sh.tangled.repo.listArtifacts", get(list_artifacts)) + .route("/xrpc/sh.tangled.repo.countArtifacts", get(count_artifacts)) + .route("/xrpc/sh.tangled.knot.listMembers", get(list_knot_members)) + .route("/xrpc/sh.tangled.knot.countMembers", get(count_knot_members)) + .route( + "/xrpc/sh.tangled.spindle.listMembers", + get(list_spindle_members), + ) + .route( + "/xrpc/sh.tangled.spindle.countMembers", + get(count_spindle_members), + ) + .route("/xrpc/sh.tangled.string.listStrings", get(list_strings)) + .route("/xrpc/sh.tangled.string.countStrings", get(count_strings)) .route("/xrpc/sh.tangled.search.query", get(search_query)) .merge(knot_proxied_routes()) .with_state(state) @@ -382,6 +438,7 @@ pub enum SubjectShape { BareDid, Collection(&'static str), BareDidOrCollection(&'static str), + OneOfCollections(&'static [&'static str]), } pub trait HasSubject { @@ -403,6 +460,31 @@ impl HasSubject for PullRecord { impl HasSubject for IssueCommentRecord { const SHAPE: SubjectShape = SubjectShape::Collection("sh.tangled.repo.issue"); } +impl HasSubject for LabelDefinitionRecord { + const SHAPE: SubjectShape = SubjectShape::BareDid; +} +impl HasSubject for LabelOpRecord { + const SHAPE: SubjectShape = + SubjectShape::OneOfCollections(&["sh.tangled.repo.issue", "sh.tangled.repo.pull"]); +} +impl HasSubject for PipelineRecord { + const SHAPE: SubjectShape = SubjectShape::BareDid; +} +impl HasSubject for PipelineStatusRecord { + const SHAPE: SubjectShape = SubjectShape::Collection("sh.tangled.pipeline"); +} +impl HasSubject for ArtifactRecord { + const SHAPE: SubjectShape = SubjectShape::BareDid; +} +impl HasSubject for KnotMemberRecord { + const SHAPE: SubjectShape = SubjectShape::BareDid; +} +impl HasSubject for SpindleMemberRecord { + const SHAPE: SubjectShape = SubjectShape::BareDid; +} +impl HasSubject for TangledStringRecord { + const SHAPE: SubjectShape = SubjectShape::BareDid; +} fn parse_subject(raw: &RawAtUriParam, shape: SubjectShape) -> Result, XrpcError> { let uri = parse_uri(raw.as_str())?; @@ -434,6 +516,20 @@ fn parse_subject(raw: &RawAtUriParam, shape: SubjectShape) -> Result or at:///{expected}/; got collection {c}" ))) } + (SubjectShape::OneOfCollections(allowed), Some(c)) if allowed.contains(&c) => { + require_rkey(&uri, c)?; + Ok(uri) + } + (SubjectShape::OneOfCollections(allowed), Some(c)) => { + Err(XrpcError::InvalidParams(format!( + "subject must be at://// with nsid in [{}]; got collection {c}", + allowed.join(", "), + ))) + } + (SubjectShape::OneOfCollections(allowed), None) => Err(XrpcError::InvalidParams(format!( + "subject must be at://// with nsid in [{}]", + allowed.join(", "), + ))), } } @@ -711,6 +807,128 @@ async fn count_issue_comments( count_for::(&state, q).map(Json) } +async fn list_label_definitions( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q) + .await + .map(Json) +} + +async fn count_label_definitions( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_label_ops( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) +} + +async fn count_label_ops( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_pipelines( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) +} + +async fn count_pipelines( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_pipeline_statuses( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q) + .await + .map(Json) +} + +async fn count_pipeline_statuses( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_artifacts( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) +} + +async fn count_artifacts( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_knot_members( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q) + .await + .map(Json) +} + +async fn count_knot_members( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_spindle_members( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q) + .await + .map(Json) +} + +async fn count_spindle_members( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_strings( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q) + .await + .map(Json) +} + +async fn count_strings( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + async fn search_query( State(state): State, XrpcQuery(q): XrpcQuery, diff --git a/crates/xrpc/tests/extended.rs b/crates/xrpc/tests/extended.rs new file mode 100644 index 0000000..a91264b --- /dev/null +++ b/crates/xrpc/tests/extended.rs @@ -0,0 +1,774 @@ +use std::sync::Arc; + +use axum::body::{Body, to_bytes}; +use bobbin_edge_index::{CoverageWatch, EdgeStore}; +use bobbin_knot_proxy::{KnotProxy, KnotProxyConfig}; +use bobbin_record_lru::{CacheCapacity, LruRecordStore}; +use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex}; +use bobbin_slingshot_client::SlingshotClient; +use bobbin_types::edges::Edge; +use bobbin_xrpc::{AppState, router}; +use http::{Request, StatusCode}; +use jacquard_common::DefaultStr; +use jacquard_common::types::nsid::Nsid; +use jacquard_common::types::string::AtUri; +use serde_json::{Value, json}; +use tower::ServiceExt; +use url::Url; +use url::form_urlencoded::byte_serialize; +use wiremock::matchers::{method, path, query_param}; +use wiremock::{Mock, MockServer, ResponseTemplate}; + +const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; +const TAG_BYTES: &str = "AAAAAAAAAAAAAAAAAAAAAAAAAAA="; +const ARTIFACT_LINK: &str = "bafkreigh2akiscaildc7gnvtklbsfhdgwz72eolmpckbqr5ej26byp3uli"; + +fn at(s: &str) -> AtUri { + AtUri::new_owned(s).unwrap() +} + +fn nsid(s: &'static str) -> Nsid { + Nsid::new_static(s).unwrap() +} + +struct Harness { + server: MockServer, + edges: Arc, + state: AppState, +} + +impl Harness { + async fn new() -> Self { + let server = MockServer::start().await; + let edges = Arc::new(EdgeStore::new()); + let coverage = Arc::new(CoverageWatch::new()); + let state = AppState::new( + Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), + SlingshotClient::new(Url::parse(&server.uri()).unwrap()).unwrap(), + edges.clone(), + coverage, + Arc::new(KnotProxy::new(KnotProxyConfig::default()).unwrap()), + Arc::new(SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES).unwrap()), + ); + Self { + server, + edges, + state, + } + } + + fn add_edge(&self, kind: &'static str, subject: &str, source: &str) { + self.edges.add(Edge { + kind: nsid(kind), + subject: at(subject), + source: at(source), + }); + } + + async fn mount(&self, did: &str, collection: &str, rkey: &str, value: Value) { + let uri = format!("at://{did}/{collection}/{rkey}"); + Mock::given(method("GET")) + .and(path("/xrpc/com.atproto.repo.getRecord")) + .and(query_param("repo", did)) + .and(query_param("collection", collection)) + .and(query_param("rkey", rkey)) + .respond_with( + ResponseTemplate::new(200).set_body_json(json!({ + "uri": uri, + "cid": CID, + "value": value, + })), + ) + .mount(&self.server) + .await; + } +} + +fn list_request(endpoint: &str, subject: &str, extras: &[(&str, &str)]) -> Request { + let mut qs = format!("subject={}", encode(subject)); + extras.iter().for_each(|(k, v)| { + qs.push('&'); + qs.push_str(k); + qs.push('='); + qs.push_str(&encode(v)); + }); + Request::builder() + .uri(format!("/xrpc/{endpoint}?{qs}")) + .body(Body::empty()) + .unwrap() +} + +fn encode(s: &str) -> String { + byte_serialize(s.as_bytes()).collect() +} + +async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) { + let status = resp.status(); + let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap(); + let parsed: Value = serde_json::from_slice(&bytes).expect("JSON body"); + (status, parsed) +} + +fn label_definition_body(name: &str) -> Value { + json!({ + "$type": "sh.tangled.label.definition", + "createdAt": "2026-05-01T00:00:00Z", + "name": name, + "scope": ["sh.tangled.repo.issue"], + "valueType": {"type": "boolean", "format": "any"} + }) +} + +fn label_op_body(subject: &str, def_uri: &str, value: &str) -> Value { + json!({ + "$type": "sh.tangled.label.op", + "performedAt": "2026-05-01T00:00:00Z", + "subject": subject, + "add": [{"key": def_uri, "value": value}], + "delete": [] + }) +} + +fn pipeline_body(repo_did: &str) -> Value { + json!({ + "$type": "sh.tangled.pipeline", + "workflows": [], + "triggerMetadata": { + "kind": "manual", + "repo": { + "did": "did:plc:teq", + "repoDid": repo_did, + "knot": "nel.pet", + "defaultBranch": "main" + } + } + }) +} + +fn pipeline_body_owner_only(owner_did: &str) -> Value { + json!({ + "$type": "sh.tangled.pipeline", + "workflows": [], + "triggerMetadata": { + "kind": "manual", + "repo": { + "did": owner_did, + "knot": "nel.pet", + "defaultBranch": "main" + } + } + }) +} + +fn pipeline_status_body(pipeline_uri: &str) -> Value { + json!({ + "$type": "sh.tangled.pipeline.status", + "createdAt": "2026-05-01T00:00:00Z", + "pipeline": pipeline_uri, + "workflow": pipeline_uri, + "status": "success" + }) +} + +fn artifact_body(repo_did: &str, name: &str) -> Value { + json!({ + "$type": "sh.tangled.repo.artifact", + "createdAt": "2026-05-01T00:00:00Z", + "name": name, + "repoDid": repo_did, + "tag": {"$bytes": TAG_BYTES}, + "artifact": { + "$type": "blob", + "ref": {"$link": ARTIFACT_LINK}, + "mimeType": "application/octet-stream", + "size": 12 + } + }) +} + +fn knot_member_body(subject_did: &str) -> Value { + json!({ + "$type": "sh.tangled.knot.member", + "createdAt": "2026-05-01T00:00:00Z", + "subject": subject_did, + "domain": "oyster.cafe" + }) +} + +fn spindle_member_body(subject_did: &str) -> Value { + json!({ + "$type": "sh.tangled.spindle.member", + "createdAt": "2026-05-01T00:00:00Z", + "subject": subject_did, + "instance": "spin.nel.pet" + }) +} + +fn string_body(filename: &str, contents: &str) -> Value { + json!({ + "$type": "sh.tangled.string", + "createdAt": "2026-05-01T00:00:00Z", + "filename": filename, + "description": "test fixture", + "contents": contents + }) +} + +#[tokio::test] +async fn list_label_definitions_keys_on_owner_did() { + let h = Harness::new().await; + let owner = "did:plc:abalone"; + let rkey = "bug"; + let source = format!("at://{owner}/sh.tangled.label.definition/{rkey}"); + h.add_edge( + "sh.tangled.label.definition", + &format!("at://{owner}"), + &source, + ); + h.mount( + owner, + "sh.tangled.label.definition", + rkey, + label_definition_body("bug"), + ) + .await; + + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(list_request( + "sh.tangled.label.listDefinitions", + &format!("at://{owner}"), + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK); + let items = body["items"].as_array().unwrap(); + assert_eq!(items.len(), 1); + assert_eq!(items[0]["value"]["name"], json!("bug")); + assert_eq!(items[0]["value"]["scope"][0], json!("sh.tangled.repo.issue")); +} + +#[tokio::test] +async fn count_label_definitions_dedupes_per_author() { + let h = Harness::new().await; + let owner = "did:plc:abalone"; + let subject = format!("at://{owner}"); + h.add_edge( + "sh.tangled.label.definition", + &subject, + &format!("at://{owner}/sh.tangled.label.definition/bug"), + ); + h.add_edge( + "sh.tangled.label.definition", + &subject, + &format!("at://{owner}/sh.tangled.label.definition/wontfix"), + ); + + let app = router(h.state.clone()); + let (_, body) = json_response( + app.oneshot(list_request( + "sh.tangled.label.countDefinitions", + &subject, + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(body["count"], json!(2)); + assert_eq!(body["distinctAuthors"], json!(1)); +} + +#[tokio::test] +async fn list_label_ops_accepts_issue_subject() { + let h = Harness::new().await; + let issue_uri = "at://did:plc:abalone/sh.tangled.repo.issue/i1"; + let author = "did:plc:nel"; + let rkey = "op1"; + let def_uri = "at://did:plc:abalone/sh.tangled.label.definition/bug"; + h.add_edge( + "sh.tangled.label.op", + issue_uri, + &format!("at://{author}/sh.tangled.label.op/{rkey}"), + ); + h.mount( + author, + "sh.tangled.label.op", + rkey, + label_op_body(issue_uri, def_uri, "true"), + ) + .await; + + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(list_request("sh.tangled.label.listOps", issue_uri, &[])) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK); + let items = body["items"].as_array().unwrap(); + assert_eq!(items.len(), 1); + assert_eq!(items[0]["value"]["subject"], json!(issue_uri)); + assert_eq!(items[0]["value"]["add"][0]["key"], json!(def_uri)); +} + +#[tokio::test] +async fn list_label_ops_pull_subject_round_trip() { + let h = Harness::new().await; + let pull_uri = "at://did:plc:abalone/sh.tangled.repo.pull/p1"; + let author = "did:plc:bailey"; + let rkey = "op1"; + let def_uri = "at://did:plc:abalone/sh.tangled.label.definition/wontfix"; + let source = format!("at://{author}/sh.tangled.label.op/{rkey}"); + let body = label_op_body(pull_uri, def_uri, "true"); + let parsed = + bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.label.op"), body.clone()) + .expect("parse label.op record"); + parsed + .extract_edges(&at(&source)) + .expect("extract") + .into_iter() + .for_each(|e| h.edges.add(e)); + h.mount(author, "sh.tangled.label.op", rkey, body).await; + + let app = router(h.state.clone()); + let (status, json) = json_response( + app.oneshot(list_request("sh.tangled.label.listOps", pull_uri, &[])) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK); + let items = json["items"].as_array().unwrap(); + assert_eq!(items.len(), 1); + assert_eq!(items[0]["value"]["subject"], json!(pull_uri)); + assert_eq!(items[0]["value"]["add"][0]["key"], json!(def_uri)); +} + +#[tokio::test] +async fn list_label_ops_rejects_bare_did_subject() { + let h = Harness::new().await; + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(list_request( + "sh.tangled.label.listOps", + "at://did:plc:abalone", + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::BAD_REQUEST); + let msg = body["message"].as_str().unwrap_or_default(); + assert!( + msg.contains("sh.tangled.repo.issue") && msg.contains("sh.tangled.repo.pull"), + "message must list allowed collections, got {msg}" + ); +} + +#[tokio::test] +async fn list_label_ops_rejects_unrelated_collection() { + let h = Harness::new().await; + let app = router(h.state.clone()); + let (status, _) = json_response( + app.oneshot(list_request( + "sh.tangled.label.listOps", + "at://did:plc:abalone/sh.tangled.repo/r1", + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::BAD_REQUEST); +} + +#[tokio::test] +async fn list_pipelines_keys_on_repo_did() { + let h = Harness::new().await; + let repo_did = "did:plc:abalone"; + let subject = format!("at://{repo_did}"); + let spindle_did = "did:plc:lyna"; + let rkey = "pl1"; + h.add_edge( + "sh.tangled.pipeline", + &subject, + &format!("at://{spindle_did}/sh.tangled.pipeline/{rkey}"), + ); + h.mount( + spindle_did, + "sh.tangled.pipeline", + rkey, + pipeline_body(repo_did), + ) + .await; + + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(list_request( + "sh.tangled.pipeline.listPipelines", + &subject, + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK); + let items = body["items"].as_array().unwrap(); + assert_eq!(items.len(), 1); + assert_eq!( + items[0]["value"]["triggerMetadata"]["repo"]["repoDid"], + json!(repo_did) + ); +} + +#[tokio::test] +async fn count_pipelines_envelope_present() { + let h = Harness::new().await; + let app = router(h.state.clone()); + let (_, body) = json_response( + app.oneshot(list_request( + "sh.tangled.pipeline.countPipelines", + "at://did:plc:abalone", + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(body["count"], json!(0)); + assert_eq!(body["coverage"]["ready"], json!(false)); +} + +#[tokio::test] +async fn list_pipeline_statuses_keys_on_pipeline_uri() { + let h = Harness::new().await; + let pipeline_uri = "at://did:plc:lyna/sh.tangled.pipeline/pl1"; + let author = "did:plc:bailey"; + let rkey = "s1"; + h.add_edge( + "sh.tangled.pipeline.status", + pipeline_uri, + &format!("at://{author}/sh.tangled.pipeline.status/{rkey}"), + ); + h.mount( + author, + "sh.tangled.pipeline.status", + rkey, + pipeline_status_body(pipeline_uri), + ) + .await; + + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(list_request( + "sh.tangled.pipeline.listStatuses", + pipeline_uri, + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK); + let items = body["items"].as_array().unwrap(); + assert_eq!(items.len(), 1); + assert_eq!(items[0]["value"]["pipeline"], json!(pipeline_uri)); + assert_eq!(items[0]["value"]["status"], json!("success")); +} + +#[tokio::test] +async fn pipeline_status_endpoint_rejects_bare_did_subject() { + let h = Harness::new().await; + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(list_request( + "sh.tangled.pipeline.listStatuses", + "at://did:plc:lyna", + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert!( + body["message"] + .as_str() + .unwrap_or_default() + .contains("sh.tangled.pipeline/"), + "{body}" + ); +} + +#[tokio::test] +async fn list_artifacts_keys_on_repo_did() { + let h = Harness::new().await; + let repo_did = "did:plc:abalone"; + let subject = format!("at://{repo_did}"); + let owner = "did:plc:nel"; + let rkey = "a1"; + h.add_edge( + "sh.tangled.repo.artifact", + &subject, + &format!("at://{owner}/sh.tangled.repo.artifact/{rkey}"), + ); + h.mount( + owner, + "sh.tangled.repo.artifact", + rkey, + artifact_body(repo_did, "out.bin"), + ) + .await; + + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(list_request( + "sh.tangled.repo.listArtifacts", + &subject, + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK); + let items = body["items"].as_array().unwrap(); + assert_eq!(items.len(), 1); + assert_eq!(items[0]["value"]["name"], json!("out.bin")); + assert_eq!(items[0]["value"]["repoDid"], json!(repo_did)); +} + +#[tokio::test] +async fn list_knot_members_keys_on_subject_did() { + let h = Harness::new().await; + let subject_did = "did:plc:nel"; + let subject = format!("at://{subject_did}"); + let admin = "did:plc:teq"; + let rkey = "m1"; + h.add_edge( + "sh.tangled.knot.member", + &subject, + &format!("at://{admin}/sh.tangled.knot.member/{rkey}"), + ); + h.mount( + admin, + "sh.tangled.knot.member", + rkey, + knot_member_body(subject_did), + ) + .await; + + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(list_request( + "sh.tangled.knot.listMembers", + &subject, + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK); + let items = body["items"].as_array().unwrap(); + assert_eq!(items.len(), 1); + assert_eq!(items[0]["value"]["subject"], json!(subject_did)); + assert_eq!(items[0]["value"]["domain"], json!("oyster.cafe")); +} + +#[tokio::test] +async fn list_spindle_members_keys_on_subject_did() { + let h = Harness::new().await; + let subject_did = "did:plc:olaren"; + let subject = format!("at://{subject_did}"); + let admin = "did:plc:teq"; + let rkey = "m1"; + h.add_edge( + "sh.tangled.spindle.member", + &subject, + &format!("at://{admin}/sh.tangled.spindle.member/{rkey}"), + ); + h.mount( + admin, + "sh.tangled.spindle.member", + rkey, + spindle_member_body(subject_did), + ) + .await; + + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(list_request( + "sh.tangled.spindle.listMembers", + &subject, + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK); + let items = body["items"].as_array().unwrap(); + assert_eq!(items.len(), 1); + assert_eq!(items[0]["value"]["subject"], json!(subject_did)); + assert_eq!(items[0]["value"]["instance"], json!("spin.nel.pet")); +} + +#[tokio::test] +async fn list_strings_keys_on_owner_did() { + let h = Harness::new().await; + let owner = "did:plc:abalone"; + let subject = format!("at://{owner}"); + let rkey = "k1"; + h.add_edge( + "sh.tangled.string", + &subject, + &format!("at://{owner}/sh.tangled.string/{rkey}"), + ); + h.mount( + owner, + "sh.tangled.string", + rkey, + string_body("snippet.rs", "fn main() {}"), + ) + .await; + + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(list_request( + "sh.tangled.string.listStrings", + &subject, + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK); + let items = body["items"].as_array().unwrap(); + assert_eq!(items.len(), 1); + assert_eq!(items[0]["value"]["filename"], json!("snippet.rs")); + assert_eq!(items[0]["value"]["contents"], json!("fn main() {}")); +} + +#[tokio::test] +async fn count_strings_dedupes_per_owner() { + let h = Harness::new().await; + let owner = "did:plc:abalone"; + let subject = format!("at://{owner}"); + ["k1", "k2", "k3"].iter().for_each(|rkey| { + h.add_edge( + "sh.tangled.string", + &subject, + &format!("at://{owner}/sh.tangled.string/{rkey}"), + ); + }); + + let app = router(h.state.clone()); + let (_, body) = json_response( + app.oneshot(list_request( + "sh.tangled.string.countStrings", + &subject, + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(body["count"], json!(3)); + assert_eq!(body["distinctAuthors"], json!(1)); +} + +#[tokio::test] +async fn extractor_to_xrpc_round_trip_for_pipeline() { + let h = Harness::new().await; + let repo_did = "did:plc:abalone"; + let spindle_did = "did:plc:lyna"; + let rkey = "pl1"; + let source = format!("at://{spindle_did}/sh.tangled.pipeline/{rkey}"); + let body = pipeline_body(repo_did); + let parsed = + bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.pipeline"), body.clone()) + .expect("parse pipeline record"); + parsed + .extract_edges(&at(&source)) + .expect("extract") + .into_iter() + .for_each(|e| h.edges.add(e)); + h.mount(spindle_did, "sh.tangled.pipeline", rkey, body).await; + + let app = router(h.state.clone()); + let (status, json) = json_response( + app.oneshot(list_request( + "sh.tangled.pipeline.listPipelines", + &format!("at://{repo_did}"), + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!( + status, + StatusCode::OK, + "extractor key must match handler subject; body: {json}" + ); + let items = json["items"].as_array().unwrap(); + assert_eq!(items.len(), 1, "expected exactly one pipeline edge"); +} + +#[tokio::test] +async fn list_pipelines_falls_back_to_owner_did_for_legacy_records() { + let h = Harness::new().await; + let owner_did = "did:plc:nel"; + let spindle_did = "did:plc:lyna"; + let rkey = "pl1"; + let source = format!("at://{spindle_did}/sh.tangled.pipeline/{rkey}"); + let body = pipeline_body_owner_only(owner_did); + let parsed = + bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.pipeline"), body.clone()) + .expect("parse pipeline record"); + parsed + .extract_edges(&at(&source)) + .expect("extract") + .into_iter() + .for_each(|e| h.edges.add(e)); + h.mount(spindle_did, "sh.tangled.pipeline", rkey, body).await; + + let app = router(h.state.clone()); + let (status, json) = json_response( + app.oneshot(list_request( + "sh.tangled.pipeline.listPipelines", + &format!("at://{owner_did}"), + &[], + )) + .await + .unwrap(), + ) + .await; + assert_eq!( + status, + StatusCode::OK, + "owner-DID fallback must hydrate; body: {json}" + ); + let items = json["items"].as_array().unwrap(); + assert_eq!(items.len(), 1); + assert_eq!( + items[0]["value"]["triggerMetadata"]["repo"]["did"], + json!(owner_did) + ); + assert!( + items[0]["value"]["triggerMetadata"]["repo"] + .get("repoDid") + .is_none(), + "fixture must omit repoDid; got {}", + items[0]["value"]["triggerMetadata"]["repo"] + ); +}