Monorepo for Tangled
tangled.org
1use std::sync::Arc;
2
3use axum::body::{Body, to_bytes};
4use bobbin_edge_index::{CoverageWatch, EdgeStore, StateIndex};
5use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig};
6use bobbin_record_lru::{CacheCapacity, LruRecordStore};
7use bobbin_resolver::RepoIdResolver;
8use bobbin_runtime::{RuntimeHasher, SystemClock};
9use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader};
10use bobbin_slingshot_client::SlingshotClient;
11use bobbin_types::edges::Edge;
12use bobbin_types::ids::SubjectRef;
13use bobbin_xrpc::{AppState, router};
14use http::{Request, StatusCode};
15use jacquard_common::DefaultStr;
16use jacquard_common::types::did::Did;
17use jacquard_common::types::nsid::Nsid;
18use jacquard_common::types::recordkey::Rkey;
19use jacquard_common::types::string::AtUri;
20use serde_json::{Value, json};
21use tower::ServiceExt;
22use url::Url;
23use url::form_urlencoded::byte_serialize;
24use wiremock::matchers::{method, path, query_param};
25use wiremock::{Mock, MockServer, ResponseTemplate};
26
27const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i";
28const TAG_BYTES: &str = "AAAAAAAAAAAAAAAAAAAAAAAAAAA=";
29const ARTIFACT_LINK: &str = "bafkreigh2akiscaildc7gnvtklbsfhdgwz72eolmpckbqr5ej26byp3uli";
30
31fn at(s: &str) -> AtUri<DefaultStr> {
32 AtUri::new_owned(s).unwrap()
33}
34
35fn did(s: &str) -> Did<DefaultStr> {
36 Did::new_owned(s).unwrap()
37}
38
39fn rkey(s: &str) -> Rkey<DefaultStr> {
40 Rkey::new_owned(s).unwrap()
41}
42
43fn nsid(s: &'static str) -> Nsid<DefaultStr> {
44 Nsid::new_static(s).unwrap()
45}
46
47fn subj(s: &str) -> SubjectRef {
48 Did::<DefaultStr>::new_owned(s)
49 .map(SubjectRef::Did)
50 .unwrap_or_else(|_| SubjectRef::Uri(AtUri::new_owned(s).unwrap()))
51}
52
53struct Harness {
54 server: MockServer,
55 edges: Arc<EdgeStore>,
56 state: AppState,
57}
58
59impl Harness {
60 async fn new() -> Self {
61 let server = MockServer::start().await;
62 let edges = Arc::new(EdgeStore::new(RuntimeHasher::default()));
63 let coverage = Arc::new(CoverageWatch::new());
64 let state = AppState::new(
65 Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))),
66 SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(),
67 edges.clone(),
68 Arc::new(StateIndex::new(RuntimeHasher::default())),
69 Arc::new(StateIndex::new(RuntimeHasher::default())),
70 coverage,
71 Arc::new(
72 KnotProxy::new(
73 KnotProxyConfig::default(),
74 KnotHttpConfig::default(),
75 Arc::new(SystemClock::new()),
76 RuntimeHasher::default(),
77 )
78 .unwrap(),
79 ),
80 Arc::new(
81 SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES, Arc::new(SystemClock::new())).unwrap(),
82 ) as Arc<dyn SearchReader>,
83 Arc::new(RepoIdResolver::detached(RuntimeHasher::default())),
84 Arc::new(jacquard_identity::JacquardResolver::default()),
85 );
86 Self {
87 server,
88 edges,
89 state,
90 }
91 }
92
93 fn add_edge(
94 &self,
95 kind: &Nsid<DefaultStr>,
96 subject: &AtUri<DefaultStr>,
97 source: &AtUri<DefaultStr>,
98 ) {
99 static EDGE_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
100 self.edges.add(Edge {
101 kind: kind.clone(),
102 subject: subj(subject.as_ref()),
103 source: source.clone(),
104 sort_micros: EDGE_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed),
105 });
106 }
107
108 async fn mount(
109 &self,
110 did: &Did<DefaultStr>,
111 collection: &Nsid<DefaultStr>,
112 rkey: &Rkey<DefaultStr>,
113 value: Value,
114 ) {
115 let uri = format!(
116 "at://{}/{}/{}",
117 did.as_ref(),
118 collection.as_ref(),
119 rkey.as_ref()
120 );
121 Mock::given(method("GET"))
122 .and(path("/xrpc/com.atproto.repo.getRecord"))
123 .and(query_param("repo", did.as_ref()))
124 .and(query_param("collection", collection.as_ref()))
125 .and(query_param("rkey", rkey.as_ref()))
126 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
127 "uri": uri,
128 "cid": CID,
129 "value": value,
130 })))
131 .mount(&self.server)
132 .await;
133 }
134}
135
136fn list_request(endpoint: &str, subject: &str, extras: &[(&str, &str)]) -> Request<Body> {
137 let mut qs = format!("subject={subject}");
138 extras.iter().for_each(|(k, v)| {
139 qs.push('&');
140 qs.push_str(k);
141 qs.push('=');
142 qs.push_str(&encode(v));
143 });
144 Request::builder()
145 .uri(format!("/xrpc/{endpoint}?{qs}"))
146 .body(Body::empty())
147 .unwrap()
148}
149
150fn list_request_escaped_subject(
151 endpoint: &str,
152 subject: &str,
153 extras: &[(&str, &str)],
154) -> Request<Body> {
155 let mut qs = format!("subject={}", encode(subject));
156 extras.iter().for_each(|(k, v)| {
157 qs.push('&');
158 qs.push_str(k);
159 qs.push('=');
160 qs.push_str(&encode(v));
161 });
162 Request::builder()
163 .uri(format!("/xrpc/{endpoint}?{qs}"))
164 .body(Body::empty())
165 .unwrap()
166}
167
168fn encode(s: &str) -> String {
169 byte_serialize(s.as_bytes()).collect()
170}
171
172async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) {
173 let status = resp.status();
174 let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap();
175 let parsed: Value = serde_json::from_slice(&bytes).expect("JSON body");
176 (status, parsed)
177}
178
179fn label_definition_body(name: &str) -> Value {
180 json!({
181 "$type": "sh.tangled.label.definition",
182 "createdAt": "2026-05-01T00:00:00Z",
183 "name": name,
184 "scope": ["sh.tangled.repo.issue"],
185 "valueType": {"type": "boolean", "format": "any"}
186 })
187}
188
189fn label_op_body(subject: &AtUri<DefaultStr>, def_uri: &AtUri<DefaultStr>, value: &str) -> Value {
190 json!({
191 "$type": "sh.tangled.label.op",
192 "performedAt": "2026-05-01T00:00:00Z",
193 "subject": subject.as_ref(),
194 "add": [{"key": def_uri.as_ref(), "value": value}],
195 "delete": []
196 })
197}
198
199fn pipeline_body(repo_did: &Did<DefaultStr>) -> Value {
200 json!({
201 "$type": "sh.tangled.pipeline",
202 "workflows": [],
203 "triggerMetadata": {
204 "kind": "manual",
205 "repo": {
206 "did": "did:plc:teq",
207 "repoDid": repo_did.as_ref(),
208 "knot": "nel.pet",
209 "defaultBranch": "main"
210 }
211 }
212 })
213}
214
215fn pipeline_body_owner_only(owner_did: &Did<DefaultStr>) -> Value {
216 json!({
217 "$type": "sh.tangled.pipeline",
218 "workflows": [],
219 "triggerMetadata": {
220 "kind": "manual",
221 "repo": {
222 "did": owner_did.as_ref(),
223 "knot": "nel.pet",
224 "defaultBranch": "main"
225 }
226 }
227 })
228}
229
230fn pipeline_status_body(pipeline_uri: &AtUri<DefaultStr>) -> Value {
231 json!({
232 "$type": "sh.tangled.pipeline.status",
233 "createdAt": "2026-05-01T00:00:00Z",
234 "pipeline": pipeline_uri.as_ref(),
235 "workflow": pipeline_uri.as_ref(),
236 "status": "success"
237 })
238}
239
240fn artifact_body(repo_did: &Did<DefaultStr>, name: &str) -> Value {
241 json!({
242 "$type": "sh.tangled.repo.artifact",
243 "createdAt": "2026-05-01T00:00:00Z",
244 "name": name,
245 "repoDid": repo_did.as_ref(),
246 "tag": {"$bytes": TAG_BYTES},
247 "artifact": {
248 "$type": "blob",
249 "ref": {"$link": ARTIFACT_LINK},
250 "mimeType": "application/octet-stream",
251 "size": 12
252 }
253 })
254}
255
256fn knot_member_body(subject_did: &Did<DefaultStr>) -> Value {
257 json!({
258 "$type": "sh.tangled.knot.member",
259 "createdAt": "2026-05-01T00:00:00Z",
260 "subject": subject_did.as_ref(),
261 "domain": "oyster.cafe"
262 })
263}
264
265fn spindle_member_body(subject_did: &Did<DefaultStr>) -> Value {
266 json!({
267 "$type": "sh.tangled.spindle.member",
268 "createdAt": "2026-05-01T00:00:00Z",
269 "subject": subject_did.as_ref(),
270 "instance": "spin.nel.pet"
271 })
272}
273
274fn string_body(filename: &str, contents: &str) -> Value {
275 json!({
276 "$type": "sh.tangled.string",
277 "createdAt": "2026-05-01T00:00:00Z",
278 "filename": filename,
279 "description": "test fixture",
280 "contents": contents
281 })
282}
283
284#[tokio::test]
285async fn list_label_definitions_keys_on_owner_did() {
286 let h = Harness::new().await;
287 let owner = did("did:plc:abalone");
288 let rk = rkey("bug");
289 let source = at(&format!(
290 "at://{}/sh.tangled.label.definition/{}",
291 owner.as_ref(),
292 rk.as_ref()
293 ));
294 h.add_edge(
295 &nsid("sh.tangled.label.definition"),
296 &at(&format!("at://{}", owner.as_ref())),
297 &source,
298 );
299 h.mount(
300 &owner,
301 &nsid("sh.tangled.label.definition"),
302 &rk,
303 label_definition_body("bug"),
304 )
305 .await;
306
307 let app = router(h.state.clone());
308 let (status, body) = json_response(
309 app.oneshot(list_request(
310 "sh.tangled.label.listDefinitions",
311 &format!("at://{}", owner.as_ref()),
312 &[],
313 ))
314 .await
315 .unwrap(),
316 )
317 .await;
318 assert_eq!(status, StatusCode::OK);
319 let items = body["items"].as_array().unwrap();
320 assert_eq!(items.len(), 1);
321 assert_eq!(items[0]["value"]["name"], json!("bug"));
322 assert_eq!(
323 items[0]["value"]["scope"][0],
324 json!("sh.tangled.repo.issue")
325 );
326}
327
328#[tokio::test]
329async fn count_label_definitions_dedupes_per_author() {
330 let h = Harness::new().await;
331 let owner = did("did:plc:abalone");
332 let subject = at(&format!("at://{}", owner.as_ref()));
333 h.add_edge(
334 &nsid("sh.tangled.label.definition"),
335 &subject,
336 &at(&format!(
337 "at://{}/sh.tangled.label.definition/bug",
338 owner.as_ref()
339 )),
340 );
341 h.add_edge(
342 &nsid("sh.tangled.label.definition"),
343 &subject,
344 &at(&format!(
345 "at://{}/sh.tangled.label.definition/wontfix",
346 owner.as_ref()
347 )),
348 );
349
350 let app = router(h.state.clone());
351 let (_, body) = json_response(
352 app.oneshot(list_request(
353 "sh.tangled.label.countDefinitions",
354 subject.as_ref(),
355 &[],
356 ))
357 .await
358 .unwrap(),
359 )
360 .await;
361 assert_eq!(body["count"], json!(2));
362 assert_eq!(body["distinctAuthors"], json!(1));
363}
364
365#[tokio::test]
366async fn list_label_ops_accepts_issue_subject() {
367 let h = Harness::new().await;
368 let issue_uri = at("at://did:plc:abalone/sh.tangled.repo.issue/i1");
369 let author = did("did:plc:nel");
370 let rk = rkey("op1");
371 let def_uri = at("at://did:plc:abalone/sh.tangled.label.definition/bug");
372 h.add_edge(
373 &nsid("sh.tangled.label.op"),
374 &issue_uri,
375 &at(&format!(
376 "at://{}/sh.tangled.label.op/{}",
377 author.as_ref(),
378 rk.as_ref()
379 )),
380 );
381 h.mount(
382 &author,
383 &nsid("sh.tangled.label.op"),
384 &rk,
385 label_op_body(&issue_uri, &def_uri, "true"),
386 )
387 .await;
388
389 let app = router(h.state.clone());
390 let (status, body) = json_response(
391 app.oneshot(list_request(
392 "sh.tangled.label.listOps",
393 issue_uri.as_ref(),
394 &[],
395 ))
396 .await
397 .unwrap(),
398 )
399 .await;
400 assert_eq!(status, StatusCode::OK);
401 let items = body["items"].as_array().unwrap();
402 assert_eq!(items.len(), 1);
403 assert_eq!(items[0]["value"]["subject"], json!(issue_uri.as_ref()));
404 assert_eq!(items[0]["value"]["add"][0]["key"], json!(def_uri.as_ref()));
405}
406
407#[tokio::test]
408async fn list_ops_accepts_percent_escaped_subject() {
409 let h = Harness::new().await;
410 let issue_uri = at("at://did:plc:clam/sh.tangled.repo.issue/i1");
411 let author = did("did:plc:nel");
412 let rk = rkey("op1");
413 let def_uri = at("at://did:plc:clam/sh.tangled.label.definition/bug");
414 h.add_edge(
415 &nsid("sh.tangled.label.op"),
416 &issue_uri,
417 &at(&format!(
418 "at://{}/sh.tangled.label.op/{}",
419 author.as_ref(),
420 rk.as_ref()
421 )),
422 );
423 h.mount(
424 &author,
425 &nsid("sh.tangled.label.op"),
426 &rk,
427 label_op_body(&issue_uri, &def_uri, "true"),
428 )
429 .await;
430
431 let app = router(h.state.clone());
432 let (status, body) = json_response(
433 app.oneshot(list_request_escaped_subject(
434 "sh.tangled.label.listOps",
435 issue_uri.as_ref(),
436 &[],
437 ))
438 .await
439 .unwrap(),
440 )
441 .await;
442 assert_eq!(
443 status,
444 StatusCode::OK,
445 "escaped subject must still be accepted"
446 );
447 assert_eq!(body["items"].as_array().unwrap().len(), 1);
448}
449
450#[tokio::test]
451async fn list_label_ops_pull_subject_round_trip() {
452 let h = Harness::new().await;
453 let pull_uri = at("at://did:plc:abalone/sh.tangled.repo.pull/p1");
454 let author = did("did:plc:bailey");
455 let rk = rkey("op1");
456 let def_uri = at("at://did:plc:abalone/sh.tangled.label.definition/wontfix");
457 let source = at(&format!(
458 "at://{}/sh.tangled.label.op/{}",
459 author.as_ref(),
460 rk.as_ref()
461 ));
462 let body = label_op_body(&pull_uri, &def_uri, "true");
463 let parsed =
464 bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.label.op"), body.clone())
465 .expect("parse label.op record");
466 parsed
467 .extract_edges(&source)
468 .expect("extract")
469 .into_iter()
470 .for_each(|e| h.edges.add(e));
471 h.mount(&author, &nsid("sh.tangled.label.op"), &rk, body)
472 .await;
473
474 let app = router(h.state.clone());
475 let (status, json) = json_response(
476 app.oneshot(list_request(
477 "sh.tangled.label.listOps",
478 pull_uri.as_ref(),
479 &[],
480 ))
481 .await
482 .unwrap(),
483 )
484 .await;
485 assert_eq!(status, StatusCode::OK);
486 let items = json["items"].as_array().unwrap();
487 assert_eq!(items.len(), 1);
488 assert_eq!(items[0]["value"]["subject"], json!(pull_uri.as_ref()));
489 assert_eq!(items[0]["value"]["add"][0]["key"], json!(def_uri.as_ref()));
490}
491
492#[tokio::test]
493async fn list_label_ops_rejects_bare_did_subject() {
494 let h = Harness::new().await;
495 let app = router(h.state.clone());
496 let (status, body) = json_response(
497 app.oneshot(list_request(
498 "sh.tangled.label.listOps",
499 "at://did:plc:abalone",
500 &[],
501 ))
502 .await
503 .unwrap(),
504 )
505 .await;
506 assert_eq!(status, StatusCode::BAD_REQUEST);
507 let msg = body["message"].as_str().unwrap_or_default();
508 assert!(
509 msg.contains("sh.tangled.repo.issue") && msg.contains("sh.tangled.repo.pull"),
510 "message must list allowed collections, got {msg}"
511 );
512}
513
514#[tokio::test]
515async fn list_label_ops_rejects_unrelated_collection() {
516 let h = Harness::new().await;
517 let app = router(h.state.clone());
518 let (status, _) = json_response(
519 app.oneshot(list_request(
520 "sh.tangled.label.listOps",
521 "at://did:plc:abalone/sh.tangled.repo/r1",
522 &[],
523 ))
524 .await
525 .unwrap(),
526 )
527 .await;
528 assert_eq!(status, StatusCode::BAD_REQUEST);
529}
530
531#[tokio::test]
532async fn list_pipelines_keys_on_repo_did() {
533 let h = Harness::new().await;
534 let repo_did = did("did:plc:abalone");
535 let subject = at(&format!("at://{}", repo_did.as_ref()));
536 let spindle_did = did("did:plc:lyna");
537 let rk = rkey("pl1");
538 h.add_edge(
539 &nsid("sh.tangled.pipeline"),
540 &subject,
541 &at(&format!(
542 "at://{}/sh.tangled.pipeline/{}",
543 spindle_did.as_ref(),
544 rk.as_ref()
545 )),
546 );
547 h.mount(
548 &spindle_did,
549 &nsid("sh.tangled.pipeline"),
550 &rk,
551 pipeline_body(&repo_did),
552 )
553 .await;
554
555 let app = router(h.state.clone());
556 let (status, body) = json_response(
557 app.oneshot(list_request(
558 "sh.tangled.pipeline.listPipelines",
559 subject.as_ref(),
560 &[],
561 ))
562 .await
563 .unwrap(),
564 )
565 .await;
566 assert_eq!(status, StatusCode::OK);
567 let items = body["items"].as_array().unwrap();
568 assert_eq!(items.len(), 1);
569 assert_eq!(
570 items[0]["value"]["triggerMetadata"]["repo"]["repoDid"],
571 json!(repo_did.as_ref())
572 );
573}
574
575#[tokio::test]
576async fn count_pipelines_returns_zero_when_no_edges() {
577 let h = Harness::new().await;
578 let app = router(h.state.clone());
579 let (_, body) = json_response(
580 app.oneshot(list_request(
581 "sh.tangled.pipeline.countPipelines",
582 "at://did:plc:abalone",
583 &[],
584 ))
585 .await
586 .unwrap(),
587 )
588 .await;
589 assert_eq!(body["count"], json!(0));
590}
591
592#[tokio::test]
593async fn list_pipeline_statuses_keys_on_pipeline_uri() {
594 let h = Harness::new().await;
595 let pipeline_uri = at("at://did:plc:lyna/sh.tangled.pipeline/pl1");
596 let author = did("did:plc:bailey");
597 let rk = rkey("s1");
598 h.add_edge(
599 &nsid("sh.tangled.pipeline.status"),
600 &pipeline_uri,
601 &at(&format!(
602 "at://{}/sh.tangled.pipeline.status/{}",
603 author.as_ref(),
604 rk.as_ref()
605 )),
606 );
607 h.mount(
608 &author,
609 &nsid("sh.tangled.pipeline.status"),
610 &rk,
611 pipeline_status_body(&pipeline_uri),
612 )
613 .await;
614
615 let app = router(h.state.clone());
616 let (status, body) = json_response(
617 app.oneshot(list_request(
618 "sh.tangled.pipeline.listStatuses",
619 pipeline_uri.as_ref(),
620 &[],
621 ))
622 .await
623 .unwrap(),
624 )
625 .await;
626 assert_eq!(status, StatusCode::OK);
627 let items = body["items"].as_array().unwrap();
628 assert_eq!(items.len(), 1);
629 assert_eq!(items[0]["value"]["pipeline"], json!(pipeline_uri.as_ref()));
630 assert_eq!(items[0]["value"]["status"], json!("success"));
631}
632
633#[tokio::test]
634async fn pipeline_status_endpoint_rejects_bare_did_subject() {
635 let h = Harness::new().await;
636 let app = router(h.state.clone());
637 let (status, body) = json_response(
638 app.oneshot(list_request(
639 "sh.tangled.pipeline.listStatuses",
640 "at://did:plc:lyna",
641 &[],
642 ))
643 .await
644 .unwrap(),
645 )
646 .await;
647 assert_eq!(status, StatusCode::BAD_REQUEST);
648 assert!(
649 body["message"]
650 .as_str()
651 .unwrap_or_default()
652 .contains("sh.tangled.pipeline/<rkey>"),
653 "{body}"
654 );
655}
656
657#[tokio::test]
658async fn list_artifacts_keys_on_repo_did() {
659 let h = Harness::new().await;
660 let repo_did = did("did:plc:abalone");
661 let subject = at(&format!("at://{}", repo_did.as_ref()));
662 let owner = did("did:plc:nel");
663 let rk = rkey("a1");
664 h.add_edge(
665 &nsid("sh.tangled.repo.artifact"),
666 &subject,
667 &at(&format!(
668 "at://{}/sh.tangled.repo.artifact/{}",
669 owner.as_ref(),
670 rk.as_ref()
671 )),
672 );
673 h.mount(
674 &owner,
675 &nsid("sh.tangled.repo.artifact"),
676 &rk,
677 artifact_body(&repo_did, "out.bin"),
678 )
679 .await;
680
681 let app = router(h.state.clone());
682 let (status, body) = json_response(
683 app.oneshot(list_request(
684 "sh.tangled.repo.listArtifacts",
685 subject.as_ref(),
686 &[],
687 ))
688 .await
689 .unwrap(),
690 )
691 .await;
692 assert_eq!(status, StatusCode::OK);
693 let items = body["items"].as_array().unwrap();
694 assert_eq!(items.len(), 1);
695 assert_eq!(items[0]["value"]["name"], json!("out.bin"));
696 assert_eq!(items[0]["value"]["repoDid"], json!(repo_did.as_ref()));
697}
698
699#[tokio::test]
700async fn list_knot_members_keys_on_subject_did() {
701 let h = Harness::new().await;
702 let subject_did = did("did:plc:nel");
703 let subject = at(&format!("at://{}", subject_did.as_ref()));
704 let admin = did("did:plc:teq");
705 let rk = rkey("m1");
706 h.add_edge(
707 &nsid("sh.tangled.knot.member"),
708 &subject,
709 &at(&format!(
710 "at://{}/sh.tangled.knot.member/{}",
711 admin.as_ref(),
712 rk.as_ref()
713 )),
714 );
715 h.mount(
716 &admin,
717 &nsid("sh.tangled.knot.member"),
718 &rk,
719 knot_member_body(&subject_did),
720 )
721 .await;
722
723 let app = router(h.state.clone());
724 let (status, body) = json_response(
725 app.oneshot(list_request(
726 "sh.tangled.knot.listMembers",
727 subject.as_ref(),
728 &[],
729 ))
730 .await
731 .unwrap(),
732 )
733 .await;
734 assert_eq!(status, StatusCode::OK);
735 let items = body["items"].as_array().unwrap();
736 assert_eq!(items.len(), 1);
737 assert_eq!(items[0]["value"]["subject"], json!(subject_did.as_ref()));
738 assert_eq!(items[0]["value"]["domain"], json!("oyster.cafe"));
739}
740
741#[tokio::test]
742async fn list_spindle_members_keys_on_subject_did() {
743 let h = Harness::new().await;
744 let subject_did = did("did:plc:olaren");
745 let subject = at(&format!("at://{}", subject_did.as_ref()));
746 let admin = did("did:plc:teq");
747 let rk = rkey("m1");
748 h.add_edge(
749 &nsid("sh.tangled.spindle.member"),
750 &subject,
751 &at(&format!(
752 "at://{}/sh.tangled.spindle.member/{}",
753 admin.as_ref(),
754 rk.as_ref()
755 )),
756 );
757 h.mount(
758 &admin,
759 &nsid("sh.tangled.spindle.member"),
760 &rk,
761 spindle_member_body(&subject_did),
762 )
763 .await;
764
765 let app = router(h.state.clone());
766 let (status, body) = json_response(
767 app.oneshot(list_request(
768 "sh.tangled.spindle.listMembers",
769 subject.as_ref(),
770 &[],
771 ))
772 .await
773 .unwrap(),
774 )
775 .await;
776 assert_eq!(status, StatusCode::OK);
777 let items = body["items"].as_array().unwrap();
778 assert_eq!(items.len(), 1);
779 assert_eq!(items[0]["value"]["subject"], json!(subject_did.as_ref()));
780 assert_eq!(items[0]["value"]["instance"], json!("spin.nel.pet"));
781}
782
783#[tokio::test]
784async fn list_strings_keys_on_owner_did() {
785 let h = Harness::new().await;
786 let owner = did("did:plc:abalone");
787 let subject = at(&format!("at://{}", owner.as_ref()));
788 let rk = rkey("k1");
789 h.add_edge(
790 &nsid("sh.tangled.string"),
791 &subject,
792 &at(&format!(
793 "at://{}/sh.tangled.string/{}",
794 owner.as_ref(),
795 rk.as_ref()
796 )),
797 );
798 h.mount(
799 &owner,
800 &nsid("sh.tangled.string"),
801 &rk,
802 string_body("snippet.rs", "fn main() {}"),
803 )
804 .await;
805
806 let app = router(h.state.clone());
807 let (status, body) = json_response(
808 app.oneshot(list_request(
809 "sh.tangled.string.listStrings",
810 subject.as_ref(),
811 &[],
812 ))
813 .await
814 .unwrap(),
815 )
816 .await;
817 assert_eq!(status, StatusCode::OK);
818 let items = body["items"].as_array().unwrap();
819 assert_eq!(items.len(), 1);
820 assert_eq!(items[0]["value"]["filename"], json!("snippet.rs"));
821 assert_eq!(items[0]["value"]["contents"], json!("fn main() {}"));
822}
823
824#[tokio::test]
825async fn count_strings_dedupes_per_owner() {
826 let h = Harness::new().await;
827 let owner = did("did:plc:abalone");
828 let subject = at(&format!("at://{}", owner.as_ref()));
829 ["k1", "k2", "k3"].iter().for_each(|r| {
830 h.add_edge(
831 &nsid("sh.tangled.string"),
832 &subject,
833 &at(&format!("at://{}/sh.tangled.string/{}", owner.as_ref(), r)),
834 );
835 });
836
837 let app = router(h.state.clone());
838 let (_, body) = json_response(
839 app.oneshot(list_request(
840 "sh.tangled.string.countStrings",
841 subject.as_ref(),
842 &[],
843 ))
844 .await
845 .unwrap(),
846 )
847 .await;
848 assert_eq!(body["count"], json!(3));
849 assert_eq!(body["distinctAuthors"], json!(1));
850}
851
852#[tokio::test]
853async fn extractor_to_xrpc_round_trip_for_pipeline() {
854 let h = Harness::new().await;
855 let repo_did = did("did:plc:abalone");
856 let spindle_did = did("did:plc:lyna");
857 let rk = rkey("pl1");
858 let source = at(&format!(
859 "at://{}/sh.tangled.pipeline/{}",
860 spindle_did.as_ref(),
861 rk.as_ref()
862 ));
863 let body = pipeline_body(&repo_did);
864 let parsed =
865 bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.pipeline"), body.clone())
866 .expect("parse pipeline record");
867 parsed
868 .extract_edges(&source)
869 .expect("extract")
870 .into_iter()
871 .for_each(|e| h.edges.add(e));
872 h.mount(&spindle_did, &nsid("sh.tangled.pipeline"), &rk, body)
873 .await;
874
875 let app = router(h.state.clone());
876 let (status, json) = json_response(
877 app.oneshot(list_request(
878 "sh.tangled.pipeline.listPipelines",
879 &format!("at://{}", repo_did.as_ref()),
880 &[],
881 ))
882 .await
883 .unwrap(),
884 )
885 .await;
886 assert_eq!(
887 status,
888 StatusCode::OK,
889 "extractor key must match handler subject, body was {json}"
890 );
891 let items = json["items"].as_array().unwrap();
892 assert_eq!(items.len(), 1, "expected exactly one pipeline edge");
893}
894
895#[tokio::test]
896async fn list_pipelines_drops_records_without_resolvable_repo_did() {
897 let h = Harness::new().await;
898 let owner_did = did("did:plc:nel");
899 let spindle_did = did("did:plc:lyna");
900 let rk = rkey("pl1");
901 let source = at(&format!(
902 "at://{}/sh.tangled.pipeline/{}",
903 spindle_did.as_ref(),
904 rk.as_ref()
905 ));
906 let body = pipeline_body_owner_only(&owner_did);
907 let parsed =
908 bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.pipeline"), body.clone())
909 .expect("parse pipeline record");
910 parsed
911 .extract_edges(&source)
912 .expect("extract")
913 .into_iter()
914 .for_each(|e| h.edges.add(e));
915 h.mount(&spindle_did, &nsid("sh.tangled.pipeline"), &rk, body)
916 .await;
917
918 let app = router(h.state.clone());
919 let (status, json) = json_response(
920 app.oneshot(list_request(
921 "sh.tangled.pipeline.listPipelines",
922 &format!("at://{}", owner_did.as_ref()),
923 &[],
924 ))
925 .await
926 .unwrap(),
927 )
928 .await;
929 assert_eq!(status, StatusCode::OK, "body was {json}");
930 let items = json["items"].as_array().unwrap();
931 assert!(
932 items.is_empty(),
933 "pipeline without resolvable repoDid must be dropped, got {json}"
934 );
935}