This repository has no description
1use std::sync::Arc;
2
3use axum::body::{Body, to_bytes};
4use bobbin_edge_index::{
5 Coverage, CoverageWatch, EdgeStore, HydrantCursor, IssueStateKind, PageToken, PullStatusKind,
6 StateIndex,
7};
8use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig};
9use bobbin_record_lru::{CacheCapacity, LruRecordStore};
10use bobbin_resolver::RepoIdResolver;
11use bobbin_runtime::{RuntimeHasher, SystemClock};
12use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader};
13use bobbin_slingshot_client::SlingshotClient;
14use bobbin_types::edges::Edge;
15use bobbin_types::ids::SubjectRef;
16use bobbin_xrpc::{AppState, router};
17use futures::stream::{self, StreamExt};
18use http::{Request, StatusCode};
19use jacquard_common::DefaultStr;
20use jacquard_common::types::did::Did;
21use jacquard_common::types::nsid::Nsid;
22use jacquard_common::types::recordkey::Rkey;
23use jacquard_common::types::string::AtUri;
24use serde_json::{Value, json};
25use tower::ServiceExt;
26use url::Url;
27use url::form_urlencoded::byte_serialize;
28use wiremock::matchers::{method, path, query_param};
29use wiremock::{Mock, MockServer, ResponseTemplate};
30
31const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i";
32
33fn at(s: &str) -> AtUri<DefaultStr> {
34 AtUri::new_owned(s).unwrap()
35}
36
37fn did(s: &str) -> Did<DefaultStr> {
38 Did::new_owned(s).unwrap()
39}
40
41fn rkey(s: &str) -> Rkey<DefaultStr> {
42 Rkey::new_owned(s).unwrap()
43}
44
45fn nsid(s: &'static str) -> Nsid<DefaultStr> {
46 Nsid::new_static(s).unwrap()
47}
48
49fn subj(s: &str) -> SubjectRef {
50 Did::<DefaultStr>::new_owned(s)
51 .map(SubjectRef::Did)
52 .unwrap_or_else(|_| SubjectRef::Uri(AtUri::new_owned(s).unwrap()))
53}
54
55struct Harness {
56 server: MockServer,
57 edges: Arc<EdgeStore>,
58 coverage: Arc<CoverageWatch>,
59 state: AppState,
60}
61
62static EDGE_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
63
64fn next_sort_micros() -> u64 {
65 EDGE_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
66}
67
68impl Harness {
69 async fn new() -> Self {
70 let server = MockServer::start().await;
71 let edges = Arc::new(EdgeStore::new(RuntimeHasher::default()));
72 let issue_states = Arc::new(StateIndex::new(RuntimeHasher::default()));
73 let pull_statuses = Arc::new(StateIndex::new(RuntimeHasher::default()));
74 let coverage = Arc::new(CoverageWatch::new());
75 let state = AppState::new(
76 Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))),
77 SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(),
78 edges.clone(),
79 issue_states.clone(),
80 pull_statuses.clone(),
81 coverage.clone(),
82 Arc::new(
83 KnotProxy::new(
84 KnotProxyConfig::default(),
85 KnotHttpConfig::default(),
86 Arc::new(SystemClock::new()),
87 RuntimeHasher::default(),
88 )
89 .unwrap(),
90 ),
91 Arc::new(
92 SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES, Arc::new(SystemClock::new())).unwrap(),
93 ) as Arc<dyn SearchReader>,
94 Arc::new(RepoIdResolver::detached(RuntimeHasher::default())),
95 );
96 Self {
97 server,
98 edges,
99 coverage,
100 state,
101 }
102 }
103
104 fn add_edge(
105 &self,
106 kind: &Nsid<DefaultStr>,
107 subject: &AtUri<DefaultStr>,
108 source: &AtUri<DefaultStr>,
109 ) {
110 self.edges.add(Edge {
111 kind: kind.clone(),
112 subject: subj(subject.as_ref()),
113 source: source.clone(),
114 sort_micros: next_sort_micros(),
115 });
116 }
117
118 async fn mount(
119 &self,
120 did: &Did<DefaultStr>,
121 collection: &Nsid<DefaultStr>,
122 rkey: &Rkey<DefaultStr>,
123 value: Value,
124 ) {
125 let uri = format!(
126 "at://{}/{}/{}",
127 did.as_ref(),
128 collection.as_ref(),
129 rkey.as_ref()
130 );
131 let body = json!({ "uri": uri, "cid": CID, "value": value });
132 Mock::given(method("GET"))
133 .and(path("/xrpc/com.atproto.repo.getRecord"))
134 .and(query_param("repo", did.as_ref()))
135 .and(query_param("collection", collection.as_ref()))
136 .and(query_param("rkey", rkey.as_ref()))
137 .respond_with(ResponseTemplate::new(200).set_body_json(body))
138 .mount(&self.server)
139 .await;
140 }
141
142 fn promote_ready(&self, events: u64, cursor: u64) {
143 self.coverage.update(|_| Coverage::Ready {
144 events_processed: events,
145 last_cursor: HydrantCursor::new(cursor),
146 });
147 }
148
149 fn warming(&self, events: u64, cursor: u64) {
150 self.coverage.update(|_| Coverage::Warming {
151 events_processed: events,
152 last_cursor: HydrantCursor::new(cursor),
153 });
154 }
155}
156
157fn list_request(endpoint: &str, subject: &str, extras: &[(&str, &str)]) -> Request<Body> {
158 let mut qs = format!("subject={subject}");
159 extras.iter().for_each(|(k, v)| {
160 qs.push('&');
161 qs.push_str(k);
162 qs.push('=');
163 qs.push_str(&encode(v));
164 });
165 Request::builder()
166 .uri(format!("/xrpc/{endpoint}?{qs}"))
167 .body(Body::empty())
168 .unwrap()
169}
170
171fn encode(s: &str) -> String {
172 byte_serialize(s.as_bytes()).collect()
173}
174
175async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) {
176 let status = resp.status();
177 let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap();
178 let parsed: Value = serde_json::from_slice(&bytes).expect("JSON body");
179 (status, parsed)
180}
181
182fn issue_body(repo_did: &Did<DefaultStr>, title: &str) -> Value {
183 json!({
184 "$type": "sh.tangled.repo.issue",
185 "repo": repo_did.as_ref(),
186 "title": title,
187 "createdAt": "2026-05-01T00:00:00Z"
188 })
189}
190
191fn pull_body(repo_did: &Did<DefaultStr>, title: &str) -> Value {
192 json!({
193 "$type": "sh.tangled.repo.pull",
194 "title": title,
195 "createdAt": "2026-05-01T00:00:00Z",
196 "rounds": [],
197 "target": {
198 "repo": repo_did.as_ref(),
199 "branch": "main"
200 }
201 })
202}
203
204fn star_body(subject_did: &Did<DefaultStr>) -> Value {
205 json!({
206 "$type": "sh.tangled.feed.star",
207 "createdAt": "2026-05-01T00:00:00Z",
208 "subject": {
209 "$type": "sh.tangled.feed.star#repo",
210 "did": subject_did.as_ref()
211 }
212 })
213}
214
215fn follow_body(subject_did: &Did<DefaultStr>) -> Value {
216 json!({
217 "$type": "sh.tangled.graph.follow",
218 "createdAt": "2026-05-01T00:00:00Z",
219 "subject": subject_did.as_ref()
220 })
221}
222
223#[tokio::test]
224async fn list_issues_with_no_edges_returns_empty_items() {
225 let h = Harness::new().await;
226 let app = router(h.state.clone());
227 let resp = app
228 .oneshot(list_request(
229 "sh.tangled.repo.listIssues",
230 "at://did:plc:abalone",
231 &[],
232 ))
233 .await
234 .unwrap();
235 let (status, body) = json_response(resp).await;
236 assert_eq!(status, StatusCode::OK);
237 assert_eq!(body["items"], json!([]));
238 assert!(body["cursor"].is_null());
239}
240
241#[tokio::test]
242async fn count_issues_with_no_edges_returns_zero() {
243 let h = Harness::new().await;
244 let app = router(h.state.clone());
245 let resp = app
246 .oneshot(list_request(
247 "sh.tangled.repo.countIssues",
248 "at://did:plc:abalone",
249 &[],
250 ))
251 .await
252 .unwrap();
253 let (status, body) = json_response(resp).await;
254 assert_eq!(status, StatusCode::OK);
255 assert_eq!(body["count"], json!(0));
256 assert_eq!(body["distinctAuthors"], json!(0));
257}
258
259#[tokio::test]
260async fn list_issues_hydrates_via_slingshot_when_edges_present() {
261 let h = Harness::new().await;
262 let repo = did("did:plc:abalone");
263 let subject = at(&format!("at://{}", repo.as_ref()));
264 let owners = [
265 ("did:plc:nel", "i1", "first"),
266 ("did:plc:olaren", "i2", "second"),
267 ];
268 stream::iter(owners)
269 .for_each(|(d, r, title)| {
270 let h = &h;
271 let subject = subject.clone();
272 let repo = repo.clone();
273 async move {
274 let d_did = did(d);
275 let rk = rkey(r);
276 h.add_edge(
277 &nsid("sh.tangled.repo.issue"),
278 &subject,
279 &at(&format!(
280 "at://{}/sh.tangled.repo.issue/{}",
281 d_did.as_ref(),
282 rk.as_ref()
283 )),
284 );
285 h.mount(
286 &d_did,
287 &nsid("sh.tangled.repo.issue"),
288 &rk,
289 issue_body(&repo, title),
290 )
291 .await;
292 }
293 })
294 .await;
295
296 let app = router(h.state.clone());
297 let resp = app
298 .oneshot(list_request(
299 "sh.tangled.repo.listIssues",
300 subject.as_ref(),
301 &[],
302 ))
303 .await
304 .unwrap();
305 let (status, body) = json_response(resp).await;
306 assert_eq!(status, StatusCode::OK);
307 let items = body["items"].as_array().expect("items array");
308 assert_eq!(items.len(), 2);
309 let titles: Vec<&str> = items
310 .iter()
311 .map(|v| v["value"]["title"].as_str().unwrap())
312 .collect();
313 assert!(titles.contains(&"first"));
314 assert!(titles.contains(&"second"));
315 assert_eq!(items[0]["cid"], CID);
316 assert!(items[0]["uri"].as_str().unwrap().starts_with("at://"));
317}
318
319#[tokio::test]
320async fn count_distinct_authors_dedupes_per_author() {
321 let h = Harness::new().await;
322 let subject = at("at://did:plc:abalone");
323 h.add_edge(
324 &nsid("sh.tangled.feed.star"),
325 &subject,
326 &at("at://did:plc:nel/sh.tangled.feed.star/s1"),
327 );
328 h.add_edge(
329 &nsid("sh.tangled.feed.star"),
330 &subject,
331 &at("at://did:plc:nel/sh.tangled.feed.star/s2"),
332 );
333 h.add_edge(
334 &nsid("sh.tangled.feed.star"),
335 &subject,
336 &at("at://did:plc:olaren/sh.tangled.feed.star/s3"),
337 );
338
339 let app = router(h.state.clone());
340 let resp = app
341 .oneshot(list_request(
342 "sh.tangled.feed.countStars",
343 subject.as_ref(),
344 &[],
345 ))
346 .await
347 .unwrap();
348 let (_, body) = json_response(resp).await;
349 assert_eq!(body["count"], json!(3));
350 assert_eq!(body["distinctAuthors"], json!(2));
351}
352
353#[tokio::test]
354async fn list_items_stable_across_coverage_promotion() {
355 let h = Harness::new().await;
356 let subject = at("at://did:plc:abalone");
357 let nel = did("did:plc:nel");
358 h.add_edge(
359 &nsid("sh.tangled.feed.star"),
360 &subject,
361 &at(&format!("at://{}/sh.tangled.feed.star/s1", nel.as_ref())),
362 );
363 h.mount(
364 &nel,
365 &nsid("sh.tangled.feed.star"),
366 &rkey("s1"),
367 star_body(&did("did:plc:abalone")),
368 )
369 .await;
370
371 let app = router(h.state.clone());
372 h.warming(1, 5);
373 let (_, before) = json_response(
374 app.clone()
375 .oneshot(list_request(
376 "sh.tangled.feed.listStars",
377 subject.as_ref(),
378 &[],
379 ))
380 .await
381 .unwrap(),
382 )
383 .await;
384 assert_eq!(before["items"].as_array().unwrap().len(), 1);
385
386 h.promote_ready(2, 9);
387 let (_, after) = json_response(
388 app.oneshot(list_request(
389 "sh.tangled.feed.listStars",
390 subject.as_ref(),
391 &[],
392 ))
393 .await
394 .unwrap(),
395 )
396 .await;
397 assert_eq!(
398 after["items"].as_array().unwrap().len(),
399 before["items"].as_array().unwrap().len(),
400 );
401 assert_eq!(after["items"], before["items"]);
402}
403
404#[tokio::test]
405async fn list_paginates_via_cursor() {
406 let h = Harness::new().await;
407 let subject = at("at://did:plc:abalone");
408 let repo = did("did:plc:abalone");
409 let owners = [
410 ("did:plc:nel", "i1"),
411 ("did:plc:olaren", "i2"),
412 ("did:plc:teq", "i3"),
413 ("did:plc:lyna", "i4"),
414 ("did:plc:bailey", "i5"),
415 ];
416 stream::iter(owners)
417 .for_each(|(d, r)| {
418 let h = &h;
419 let subject = subject.clone();
420 let repo = repo.clone();
421 async move {
422 let d_did = did(d);
423 let rk = rkey(r);
424 h.add_edge(
425 &nsid("sh.tangled.repo.issue"),
426 &subject,
427 &at(&format!(
428 "at://{}/sh.tangled.repo.issue/{}",
429 d_did.as_ref(),
430 rk.as_ref()
431 )),
432 );
433 h.mount(
434 &d_did,
435 &nsid("sh.tangled.repo.issue"),
436 &rk,
437 issue_body(&repo, &format!("issue-{}", rk.as_ref())),
438 )
439 .await;
440 }
441 })
442 .await;
443
444 let app = router(h.state.clone());
445 let (_, page1) = json_response(
446 app.clone()
447 .oneshot(list_request(
448 "sh.tangled.repo.listIssues",
449 subject.as_ref(),
450 &[("limit", "2")],
451 ))
452 .await
453 .unwrap(),
454 )
455 .await;
456 let page1_items = page1["items"].as_array().unwrap().clone();
457 assert_eq!(page1_items.len(), 2);
458 let cursor = page1["cursor"]
459 .as_str()
460 .expect("first page must yield a cursor")
461 .to_owned();
462 assert!(
463 PageToken::decode_token(&cursor).is_ok(),
464 "cursor must be a TID-shaped token"
465 );
466
467 let (_, page2) = json_response(
468 app.oneshot(list_request(
469 "sh.tangled.repo.listIssues",
470 subject.as_ref(),
471 &[("limit", "10"), ("cursor", &cursor)],
472 ))
473 .await
474 .unwrap(),
475 )
476 .await;
477 let page2_items = page2["items"].as_array().unwrap().clone();
478 assert_eq!(page2_items.len(), 3);
479 assert!(page2["cursor"].is_null(), "tail page must not promise more");
480
481 let union: Vec<&str> = page1_items
482 .iter()
483 .chain(page2_items.iter())
484 .map(|item| item["uri"].as_str().unwrap())
485 .collect();
486 assert_eq!(union.len(), owners.len(), "union covers every owner");
487 let mut sorted = union.clone();
488 sorted.sort();
489 sorted.dedup();
490 assert_eq!(sorted.len(), owners.len(), "no duplicates across pages");
491}
492
493#[tokio::test]
494async fn pagination_unaffected_by_coverage_promotion() {
495 let h = Harness::new().await;
496 let subject = at("at://did:plc:abalone");
497 let repo = did("did:plc:abalone");
498 let owners = [("did:plc:nel", "i1"), ("did:plc:olaren", "i2")];
499 stream::iter(owners)
500 .for_each(|(d, r)| {
501 let h = &h;
502 let subject = subject.clone();
503 let repo = repo.clone();
504 async move {
505 let d_did = did(d);
506 let rk = rkey(r);
507 h.add_edge(
508 &nsid("sh.tangled.repo.issue"),
509 &subject,
510 &at(&format!(
511 "at://{}/sh.tangled.repo.issue/{}",
512 d_did.as_ref(),
513 rk.as_ref()
514 )),
515 );
516 h.mount(
517 &d_did,
518 &nsid("sh.tangled.repo.issue"),
519 &rk,
520 issue_body(&repo, &format!("issue-{}", rk.as_ref())),
521 )
522 .await;
523 }
524 })
525 .await;
526
527 h.warming(1, 5);
528 let app = router(h.state.clone());
529 let (_, page1) = json_response(
530 app.clone()
531 .oneshot(list_request(
532 "sh.tangled.repo.listIssues",
533 subject.as_ref(),
534 &[("limit", "1")],
535 ))
536 .await
537 .unwrap(),
538 )
539 .await;
540 let cursor = page1["cursor"].as_str().unwrap().to_owned();
541
542 h.promote_ready(2, 9);
543 let (_, page2) = json_response(
544 app.oneshot(list_request(
545 "sh.tangled.repo.listIssues",
546 subject.as_ref(),
547 &[("limit", "10"), ("cursor", &cursor)],
548 ))
549 .await
550 .unwrap(),
551 )
552 .await;
553 assert_eq!(page2["items"].as_array().unwrap().len(), 1);
554}
555
556#[tokio::test]
557async fn invalid_cursor_returns_400() {
558 let h = Harness::new().await;
559 let app = router(h.state.clone());
560 let resp = app
561 .oneshot(list_request(
562 "sh.tangled.repo.listIssues",
563 "at://did:plc:abalone",
564 &[("cursor", "not-a-number")],
565 ))
566 .await
567 .unwrap();
568 let (status, body) = json_response(resp).await;
569 assert_eq!(status, StatusCode::BAD_REQUEST);
570 assert_eq!(body["error"], "InvalidRequest");
571}
572
573#[tokio::test]
574async fn list_follows_subject_is_followee_did() {
575 let h = Harness::new().await;
576 let followee = did("did:plc:bailey");
577 let subject = at(&format!("at://{}", followee.as_ref()));
578 h.add_edge(
579 &nsid("sh.tangled.graph.follow"),
580 &subject,
581 &at("at://did:plc:nel/sh.tangled.graph.follow/f1"),
582 );
583 h.mount(
584 &did("did:plc:nel"),
585 &nsid("sh.tangled.graph.follow"),
586 &rkey("f1"),
587 follow_body(&followee),
588 )
589 .await;
590
591 let app = router(h.state.clone());
592 let (status, body) = json_response(
593 app.oneshot(list_request(
594 "sh.tangled.graph.listFollows",
595 subject.as_ref(),
596 &[],
597 ))
598 .await
599 .unwrap(),
600 )
601 .await;
602 assert_eq!(status, StatusCode::OK);
603 let items = body["items"].as_array().unwrap();
604 assert_eq!(items.len(), 1);
605 assert_eq!(items[0]["value"]["subject"], followee.as_ref());
606}
607
608#[tokio::test]
609async fn get_follow_returns_uri_when_present_and_404_otherwise() {
610 let h = Harness::new().await;
611 let followee = did("did:plc:bailey");
612 let subject = at(&format!("at://{}", followee.as_ref()));
613 h.add_edge(
614 &nsid("sh.tangled.graph.follow"),
615 &subject,
616 &at("at://did:plc:nel/sh.tangled.graph.follow/f1"),
617 );
618
619 let app = router(h.state.clone());
620 let (status, body) = json_response(
621 app.clone()
622 .oneshot(list_request(
623 "sh.tangled.graph.getFollow",
624 followee.as_ref(),
625 &[("actor", "did:plc:nel")],
626 ))
627 .await
628 .unwrap(),
629 )
630 .await;
631 assert_eq!(status, StatusCode::OK);
632 assert_eq!(body["uri"], "at://did:plc:nel/sh.tangled.graph.follow/f1");
633
634 // different actor never followed them, so this is 404 not a zero-ish success
635 let (status, _) = json_response(
636 app.oneshot(list_request(
637 "sh.tangled.graph.getFollow",
638 followee.as_ref(),
639 &[("actor", "did:plc:someoneelse")],
640 ))
641 .await
642 .unwrap(),
643 )
644 .await;
645 assert_eq!(status, StatusCode::NOT_FOUND);
646}
647
648#[tokio::test]
649async fn get_star_returns_uri_when_present_and_404_otherwise() {
650 let h = Harness::new().await;
651 let repo_did = did("did:plc:limpet");
652 let subject = at(&format!("at://{}", repo_did.as_ref()));
653 h.add_edge(
654 &nsid("sh.tangled.feed.star"),
655 &subject,
656 &at("at://did:plc:nel/sh.tangled.feed.star/s1"),
657 );
658
659 let app = router(h.state.clone());
660 let (status, body) = json_response(
661 app.clone()
662 .oneshot(list_request(
663 "sh.tangled.feed.getStar",
664 repo_did.as_ref(),
665 &[("actor", "did:plc:nel")],
666 ))
667 .await
668 .unwrap(),
669 )
670 .await;
671 assert_eq!(status, StatusCode::OK);
672 assert_eq!(body["uri"], "at://did:plc:nel/sh.tangled.feed.star/s1");
673
674 let (status, _) = json_response(
675 app.oneshot(list_request(
676 "sh.tangled.feed.getStar",
677 repo_did.as_ref(),
678 &[("actor", "did:plc:someoneelse")],
679 ))
680 .await
681 .unwrap(),
682 )
683 .await;
684 assert_eq!(status, StatusCode::NOT_FOUND);
685}
686
687#[tokio::test]
688async fn upstream_failure_during_hydration_drops_only_that_item() {
689 let h = Harness::new().await;
690 let subject = at("at://did:plc:squid");
691 let kind = nsid("sh.tangled.repo.issue");
692 let repo = did("did:plc:squid");
693 h.add_edge(
694 &kind,
695 &subject,
696 &at("at://did:plc:nel/sh.tangled.repo.issue/ok"),
697 );
698 h.add_edge(
699 &kind,
700 &subject,
701 &at("at://did:plc:teq/sh.tangled.repo.issue/flaky"),
702 );
703 h.mount(
704 &did("did:plc:nel"),
705 &kind,
706 &rkey("ok"),
707 issue_body(&repo, "kelp survey"),
708 )
709 .await;
710 Mock::given(method("GET"))
711 .and(path("/xrpc/com.atproto.repo.getRecord"))
712 .and(query_param("repo", "did:plc:teq"))
713 .and(query_param("collection", "sh.tangled.repo.issue"))
714 .and(query_param("rkey", "flaky"))
715 .respond_with(ResponseTemplate::new(503))
716 .mount(&h.server)
717 .await;
718
719 let app = router(h.state.clone());
720 let (status, body) = json_response(
721 app.oneshot(list_request(
722 "sh.tangled.repo.listIssues",
723 subject.as_ref(),
724 &[],
725 ))
726 .await
727 .unwrap(),
728 )
729 .await;
730 assert_eq!(status, StatusCode::OK);
731 let items = body["items"].as_array().expect("items array");
732 assert_eq!(items.len(), 1, "flaky item dropped, healthy sibling kept");
733 assert_eq!(
734 items[0]["uri"].as_str().unwrap(),
735 "at://did:plc:nel/sh.tangled.repo.issue/ok",
736 );
737}
738
739#[tokio::test]
740async fn transient_failure_keeps_edge_so_count_stays_whole() {
741 let h = Harness::new().await;
742 let subject = at("at://did:plc:squid");
743 let kind = nsid("sh.tangled.repo.issue");
744 let repo = did("did:plc:squid");
745 h.add_edge(
746 &kind,
747 &subject,
748 &at("at://did:plc:nel/sh.tangled.repo.issue/ok"),
749 );
750 h.add_edge(
751 &kind,
752 &subject,
753 &at("at://did:plc:teq/sh.tangled.repo.issue/flaky"),
754 );
755 h.mount(
756 &did("did:plc:nel"),
757 &kind,
758 &rkey("ok"),
759 issue_body(&repo, "kelp survey"),
760 )
761 .await;
762 Mock::given(method("GET"))
763 .and(path("/xrpc/com.atproto.repo.getRecord"))
764 .and(query_param("repo", "did:plc:teq"))
765 .and(query_param("collection", "sh.tangled.repo.issue"))
766 .and(query_param("rkey", "flaky"))
767 .respond_with(ResponseTemplate::new(503))
768 .mount(&h.server)
769 .await;
770
771 let app = router(h.state.clone());
772 let (status, body) = json_response(
773 app.clone()
774 .oneshot(list_request(
775 "sh.tangled.repo.listIssues",
776 subject.as_ref(),
777 &[],
778 ))
779 .await
780 .unwrap(),
781 )
782 .await;
783 assert_eq!(status, StatusCode::OK);
784 assert_eq!(body["items"].as_array().unwrap().len(), 1);
785
786 let (cstatus, cbody) = json_response(
787 app.oneshot(list_request(
788 "sh.tangled.repo.countIssues",
789 subject.as_ref(),
790 &[],
791 ))
792 .await
793 .unwrap(),
794 )
795 .await;
796 assert_eq!(cstatus, StatusCode::OK);
797 assert_eq!(
798 cbody["count"],
799 json!(2),
800 "a transient 503 must not evict the edge, count stays whole",
801 );
802}
803
804#[tokio::test]
805async fn gone_item_is_evicted_so_count_converges_to_list() {
806 let h = Harness::new().await;
807 let subject = at("at://did:plc:squid");
808 let kind = nsid("sh.tangled.repo.issue");
809 let repo = did("did:plc:squid");
810 h.add_edge(
811 &kind,
812 &subject,
813 &at("at://did:plc:nel/sh.tangled.repo.issue/ok"),
814 );
815 h.add_edge(
816 &kind,
817 &subject,
818 &at("at://did:plc:teq/sh.tangled.repo.issue/gone"),
819 );
820 h.mount(
821 &did("did:plc:nel"),
822 &kind,
823 &rkey("ok"),
824 issue_body(&repo, "kelp survey"),
825 )
826 .await;
827 Mock::given(method("GET"))
828 .and(path("/xrpc/com.atproto.repo.getRecord"))
829 .and(query_param("repo", "did:plc:teq"))
830 .and(query_param("collection", "sh.tangled.repo.issue"))
831 .and(query_param("rkey", "gone"))
832 .respond_with(ResponseTemplate::new(404))
833 .mount(&h.server)
834 .await;
835
836 let app = router(h.state.clone());
837 let (status, body) = json_response(
838 app.clone()
839 .oneshot(list_request(
840 "sh.tangled.repo.listIssues",
841 subject.as_ref(),
842 &[],
843 ))
844 .await
845 .unwrap(),
846 )
847 .await;
848 assert_eq!(status, StatusCode::OK);
849 assert_eq!(
850 body["items"].as_array().unwrap().len(),
851 1,
852 "gone item dropped from the page"
853 );
854
855 let (cstatus, cbody) = json_response(
856 app.oneshot(list_request(
857 "sh.tangled.repo.countIssues",
858 subject.as_ref(),
859 &[],
860 ))
861 .await
862 .unwrap(),
863 )
864 .await;
865 assert_eq!(cstatus, StatusCode::OK);
866 assert_eq!(
867 cbody["count"],
868 json!(1),
869 "a definitive 404 must evict the dead edge so count matches the list",
870 );
871}
872
873#[tokio::test]
874async fn handle_authority_subject_is_400() {
875 let h = Harness::new().await;
876 let app = router(h.state.clone());
877 let cases = [
878 "sh.tangled.feed.listStars",
879 "sh.tangled.feed.countStars",
880 "sh.tangled.graph.listFollows",
881 "sh.tangled.graph.countFollows",
882 "sh.tangled.repo.listIssues",
883 "sh.tangled.repo.countIssues",
884 "sh.tangled.repo.listPulls",
885 "sh.tangled.repo.countPulls",
886 "sh.tangled.feed.listComments",
887 "sh.tangled.feed.countComments",
888 ];
889 stream::iter(cases)
890 .for_each(|endpoint| {
891 let app = app.clone();
892 async move {
893 let resp = app
894 .oneshot(list_request(endpoint, "at://oyster.cafe", &[]))
895 .await
896 .unwrap();
897 let (status, body) = json_response(resp).await;
898 assert_eq!(status, StatusCode::BAD_REQUEST, "{endpoint}");
899 assert_eq!(body["error"], "InvalidRequest", "{endpoint}");
900 assert!(
901 body["message"]
902 .as_str()
903 .unwrap_or_default()
904 .contains("did, not a handle"),
905 "{endpoint}: {}",
906 body["message"]
907 );
908 }
909 })
910 .await;
911}
912
913#[tokio::test]
914async fn empty_subject_is_400() {
915 let h = Harness::new().await;
916 let app = router(h.state.clone());
917 let resp = app
918 .oneshot(list_request("sh.tangled.repo.listIssues", "", &[]))
919 .await
920 .unwrap();
921 let (status, body) = json_response(resp).await;
922 assert_eq!(status, StatusCode::BAD_REQUEST);
923 assert_eq!(body["error"], "InvalidRequest");
924}
925
926#[tokio::test]
927async fn limit_below_min_or_above_max_is_400() {
928 let h = Harness::new().await;
929 let app = router(h.state.clone());
930 let cases = [("0", "below"), ("1001", "above")];
931 stream::iter(cases)
932 .for_each(|(limit, label)| {
933 let app = app.clone();
934 async move {
935 let resp = app
936 .oneshot(list_request(
937 "sh.tangled.repo.listIssues",
938 "at://did:plc:abalone",
939 &[("limit", limit)],
940 ))
941 .await
942 .unwrap();
943 let (status, body) = json_response(resp).await;
944 assert_eq!(status, StatusCode::BAD_REQUEST, "limit {label}");
945 assert_eq!(body["error"], "InvalidRequest", "limit {label}");
946 }
947 })
948 .await;
949}
950
951#[tokio::test]
952async fn count_after_remove_source_returns_zero() {
953 let h = Harness::new().await;
954 let subject = at("at://did:plc:abalone");
955 let source = at("at://did:plc:nel/sh.tangled.feed.star/s1");
956 h.add_edge(&nsid("sh.tangled.feed.star"), &subject, &source);
957 h.edges.remove_source(&source);
958
959 let app = router(h.state.clone());
960 let (_, body) = json_response(
961 app.oneshot(list_request(
962 "sh.tangled.feed.countStars",
963 subject.as_ref(),
964 &[],
965 ))
966 .await
967 .unwrap(),
968 )
969 .await;
970 assert_eq!(body["count"], json!(0));
971 assert_eq!(body["distinctAuthors"], json!(0));
972}
973
974#[tokio::test]
975async fn list_feed_comments_hydrates_end_to_end() {
976 let h = Harness::new().await;
977 let issue_uri = at("at://did:plc:abalone/sh.tangled.repo.issue/i1");
978 let nel = did("did:plc:nel");
979 let rk = rkey("c1");
980 h.add_edge(
981 &nsid("sh.tangled.feed.comment"),
982 &issue_uri,
983 &at(&format!(
984 "at://{}/sh.tangled.feed.comment/{}",
985 nel.as_ref(),
986 rk.as_ref()
987 )),
988 );
989 h.mount(
990 &nel,
991 &nsid("sh.tangled.feed.comment"),
992 &rk,
993 json!({
994 "$type": "sh.tangled.feed.comment",
995 "subject": { "uri": issue_uri.as_ref(), "cid": "bafkqaaa" },
996 "body": { "$type": "sh.tangled.markup.markdown", "text": "thoughts" },
997 "createdAt": "2026-05-01T00:00:00Z"
998 }),
999 )
1000 .await;
1001
1002 let app = router(h.state.clone());
1003 let (status, body) = json_response(
1004 app.oneshot(list_request(
1005 "sh.tangled.feed.listComments",
1006 issue_uri.as_ref(),
1007 &[],
1008 ))
1009 .await
1010 .unwrap(),
1011 )
1012 .await;
1013 assert_eq!(status, StatusCode::OK);
1014 let items = body["items"].as_array().unwrap();
1015 assert_eq!(items.len(), 1);
1016 assert_eq!(items[0]["value"]["body"]["text"], json!("thoughts"));
1017 assert_eq!(
1018 items[0]["value"]["subject"]["uri"],
1019 json!(issue_uri.as_ref())
1020 );
1021}
1022
1023#[tokio::test]
1024async fn list_item_cid_is_present() {
1025 let h = Harness::new().await;
1026 let subject = at("at://did:plc:abalone");
1027 let nel = did("did:plc:nel");
1028 h.add_edge(
1029 &nsid("sh.tangled.feed.star"),
1030 &subject,
1031 &at(&format!("at://{}/sh.tangled.feed.star/s1", nel.as_ref())),
1032 );
1033 h.mount(
1034 &nel,
1035 &nsid("sh.tangled.feed.star"),
1036 &rkey("s1"),
1037 star_body(&did("did:plc:abalone")),
1038 )
1039 .await;
1040
1041 let app = router(h.state.clone());
1042 let (_, body) = json_response(
1043 app.oneshot(list_request(
1044 "sh.tangled.feed.listStars",
1045 subject.as_ref(),
1046 &[],
1047 ))
1048 .await
1049 .unwrap(),
1050 )
1051 .await;
1052 let item = &body["items"][0];
1053 assert!(
1054 item.as_object().unwrap().contains_key("cid"),
1055 "list items must mirror getRecord output shape and include cid"
1056 );
1057 assert_eq!(item["cid"], json!(CID));
1058}
1059
1060#[tokio::test]
1061async fn count_feed_comments_subjects_on_issue_uri() {
1062 let h = Harness::new().await;
1063 let issue_uri = at("at://did:plc:abalone/sh.tangled.repo.issue/i1");
1064 h.add_edge(
1065 &nsid("sh.tangled.feed.comment"),
1066 &issue_uri,
1067 &at("at://did:plc:nel/sh.tangled.feed.comment/c1"),
1068 );
1069 h.add_edge(
1070 &nsid("sh.tangled.feed.comment"),
1071 &issue_uri,
1072 &at("at://did:plc:olaren/sh.tangled.feed.comment/c2"),
1073 );
1074
1075 let app = router(h.state.clone());
1076 let (status, body) = json_response(
1077 app.oneshot(list_request(
1078 "sh.tangled.feed.countComments",
1079 issue_uri.as_ref(),
1080 &[],
1081 ))
1082 .await
1083 .unwrap(),
1084 )
1085 .await;
1086 assert_eq!(status, StatusCode::OK);
1087 assert_eq!(body["count"], json!(2));
1088 assert_eq!(body["distinctAuthors"], json!(2));
1089}
1090
1091#[tokio::test]
1092async fn list_item_404_dropped_not_404_for_subject() {
1093 let h = Harness::new().await;
1094 let subject = at("at://did:plc:squid");
1095 let kind = nsid("sh.tangled.repo.issue");
1096 let repo = did("did:plc:squid");
1097 h.add_edge(
1098 &kind,
1099 &subject,
1100 &at("at://did:plc:nel/sh.tangled.repo.issue/live"),
1101 );
1102 h.add_edge(
1103 &kind,
1104 &subject,
1105 &at("at://did:plc:teq/sh.tangled.repo.issue/missing"),
1106 );
1107 h.mount(
1108 &did("did:plc:nel"),
1109 &kind,
1110 &rkey("live"),
1111 issue_body(&repo, "kelp survives"),
1112 )
1113 .await;
1114 Mock::given(method("GET"))
1115 .and(path("/xrpc/com.atproto.repo.getRecord"))
1116 .and(query_param("repo", "did:plc:teq"))
1117 .and(query_param("collection", "sh.tangled.repo.issue"))
1118 .and(query_param("rkey", "missing"))
1119 .respond_with(ResponseTemplate::new(404).set_body_json(json!({
1120 "error": "RecordNotFound",
1121 "message": "could not find record"
1122 })))
1123 .mount(&h.server)
1124 .await;
1125
1126 let app = router(h.state.clone());
1127 let (status, body) = json_response(
1128 app.oneshot(list_request(
1129 "sh.tangled.repo.listIssues",
1130 subject.as_ref(),
1131 &[],
1132 ))
1133 .await
1134 .unwrap(),
1135 )
1136 .await;
1137 assert_eq!(
1138 status,
1139 StatusCode::OK,
1140 "a stale-index 404 drops that item, it must not 404 or 502 the subject's list",
1141 );
1142 let items = body["items"].as_array().expect("items array");
1143 assert_eq!(items.len(), 1, "stale 404 item dropped, live sibling kept");
1144 assert_eq!(
1145 items[0]["uri"].as_str().unwrap(),
1146 "at://did:plc:nel/sh.tangled.repo.issue/live",
1147 );
1148}
1149
1150#[tokio::test]
1151async fn list_item_with_wrong_type_tag_dropped() {
1152 let h = Harness::new().await;
1153 let subject = at("at://did:plc:squid");
1154 let kind = nsid("sh.tangled.feed.star");
1155 h.add_edge(
1156 &kind,
1157 &subject,
1158 &at("at://did:plc:nel/sh.tangled.feed.star/good"),
1159 );
1160 h.add_edge(
1161 &kind,
1162 &subject,
1163 &at("at://did:plc:teq/sh.tangled.feed.star/wrong"),
1164 );
1165 h.mount(
1166 &did("did:plc:nel"),
1167 &kind,
1168 &rkey("good"),
1169 star_body(&did("did:plc:squid")),
1170 )
1171 .await;
1172 Mock::given(method("GET"))
1173 .and(path("/xrpc/com.atproto.repo.getRecord"))
1174 .and(query_param("repo", "did:plc:teq"))
1175 .and(query_param("collection", "sh.tangled.feed.star"))
1176 .and(query_param("rkey", "wrong"))
1177 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
1178 "uri": "at://did:plc:teq/sh.tangled.feed.star/wrong",
1179 "cid": CID,
1180 "value": {
1181 "$type": "sh.tangled.feed.reaction",
1182 "createdAt": "2026-05-01T00:00:00Z",
1183 "subject": "at://did:plc:squid"
1184 }
1185 })))
1186 .mount(&h.server)
1187 .await;
1188
1189 let app = router(h.state.clone());
1190 let (status, body) = json_response(
1191 app.oneshot(list_request(
1192 "sh.tangled.feed.listStars",
1193 subject.as_ref(),
1194 &[],
1195 ))
1196 .await
1197 .unwrap(),
1198 )
1199 .await;
1200 assert_eq!(status, StatusCode::OK);
1201 let items = body["items"].as_array().expect("items array");
1202 assert_eq!(items.len(), 1, "wrong-type item dropped, valid star kept");
1203 assert_eq!(
1204 items[0]["uri"].as_str().unwrap(),
1205 "at://did:plc:nel/sh.tangled.feed.star/good",
1206 );
1207}
1208
1209#[tokio::test]
1210async fn list_item_with_mismatched_collection_dropped() {
1211 let h = Harness::new().await;
1212 let subject = at("at://did:plc:squid");
1213 let kind = nsid("sh.tangled.repo.issue");
1214 let repo = did("did:plc:squid");
1215 h.add_edge(
1216 &kind,
1217 &subject,
1218 &at("at://did:plc:nel/sh.tangled.repo.issue/live"),
1219 );
1220 h.add_edge(
1221 &kind,
1222 &subject,
1223 &at("at://did:plc:teq/sh.tangled.feed.star/whelk"),
1224 );
1225 h.mount(
1226 &did("did:plc:nel"),
1227 &kind,
1228 &rkey("live"),
1229 issue_body(&repo, "kelp survives"),
1230 )
1231 .await;
1232
1233 let app = router(h.state.clone());
1234 let (status, body) = json_response(
1235 app.oneshot(list_request(
1236 "sh.tangled.repo.listIssues",
1237 subject.as_ref(),
1238 &[],
1239 ))
1240 .await
1241 .unwrap(),
1242 )
1243 .await;
1244 assert_eq!(
1245 status,
1246 StatusCode::OK,
1247 "a mismatched-collection index edge must not 400 the subject's list",
1248 );
1249 let items = body["items"].as_array().expect("items array");
1250 assert_eq!(
1251 items.len(),
1252 1,
1253 "mismatched-collection edge dropped, live sibling kept"
1254 );
1255 assert_eq!(
1256 items[0]["uri"].as_str().unwrap(),
1257 "at://did:plc:nel/sh.tangled.repo.issue/live",
1258 );
1259}
1260
1261#[tokio::test]
1262async fn bare_did_endpoints_reject_at_uri_subject() {
1263 let h = Harness::new().await;
1264 let app = router(h.state.clone());
1265 let cases = [
1266 "sh.tangled.graph.listFollows",
1267 "sh.tangled.graph.countFollows",
1268 ];
1269 stream::iter(cases)
1270 .for_each(|endpoint| {
1271 let app = app.clone();
1272 async move {
1273 let resp = app
1274 .oneshot(list_request(
1275 endpoint,
1276 "at://did:plc:abalone/sh.tangled.repo/r1",
1277 &[],
1278 ))
1279 .await
1280 .unwrap();
1281 let (status, body) = json_response(resp).await;
1282 assert_eq!(status, StatusCode::BAD_REQUEST, "{endpoint}");
1283 assert_eq!(body["error"], "InvalidRequest", "{endpoint}");
1284 assert!(
1285 body["message"]
1286 .as_str()
1287 .unwrap_or_default()
1288 .contains("bare did"),
1289 "{endpoint}: {}",
1290 body["message"],
1291 );
1292 }
1293 })
1294 .await;
1295}
1296
1297#[tokio::test]
1298async fn repo_pointing_endpoints_reject_at_uri_subject() {
1299 let h = Harness::new().await;
1300 let app = router(h.state.clone());
1301 let cases = [
1302 "sh.tangled.repo.listIssues",
1303 "sh.tangled.repo.countIssues",
1304 "sh.tangled.repo.listPulls",
1305 "sh.tangled.repo.countPulls",
1306 "sh.tangled.repo.listArtifacts",
1307 "sh.tangled.repo.countArtifacts",
1308 ];
1309 stream::iter(cases)
1310 .for_each(|endpoint| {
1311 let app = app.clone();
1312 async move {
1313 let resp = app
1314 .oneshot(list_request(
1315 endpoint,
1316 "at://did:plc:abalone/sh.tangled.repo/r1",
1317 &[],
1318 ))
1319 .await
1320 .unwrap();
1321 let (status, body) = json_response(resp).await;
1322 assert_eq!(
1323 status,
1324 StatusCode::BAD_REQUEST,
1325 "{endpoint} must reject rkey-form subjects since rkeys are unstable; clients must send the repoDID",
1326 );
1327 assert!(
1328 body["message"]
1329 .as_str()
1330 .unwrap_or_default()
1331 .contains("bare did"),
1332 "{endpoint}: {}",
1333 body["message"],
1334 );
1335 }
1336 })
1337 .await;
1338}
1339
1340#[tokio::test]
1341async fn repo_pointing_endpoints_accept_bare_did() {
1342 let h = Harness::new().await;
1343 let app = router(h.state.clone());
1344 let cases = [
1345 "sh.tangled.repo.listIssues",
1346 "sh.tangled.repo.countIssues",
1347 "sh.tangled.repo.listPulls",
1348 "sh.tangled.repo.countPulls",
1349 "sh.tangled.repo.listArtifacts",
1350 "sh.tangled.repo.countArtifacts",
1351 ];
1352 stream::iter(cases)
1353 .for_each(|endpoint| {
1354 let app = app.clone();
1355 async move {
1356 let resp = app
1357 .oneshot(list_request(endpoint, "did:plc:abalone", &[]))
1358 .await
1359 .unwrap();
1360 let (status, _body) = json_response(resp).await;
1361 assert_eq!(status, StatusCode::OK, "{endpoint} must accept bare did");
1362 }
1363 })
1364 .await;
1365}
1366
1367#[tokio::test]
1368async fn feed_comment_endpoints_reject_bare_did_or_wrong_collection() {
1369 let h = Harness::new().await;
1370 let app = router(h.state.clone());
1371 let endpoints = [
1372 "sh.tangled.feed.listComments",
1373 "sh.tangled.feed.countComments",
1374 ];
1375 let inputs = [
1376 "at://did:plc:abalone",
1377 "at://did:plc:abalone/sh.tangled.repo/r1",
1378 ];
1379 let cases = endpoints
1380 .iter()
1381 .copied()
1382 .flat_map(|endpoint| inputs.iter().copied().map(move |input| (endpoint, input)));
1383 stream::iter(cases)
1384 .for_each(|(endpoint, input)| {
1385 let app = app.clone();
1386 async move {
1387 let resp = app
1388 .oneshot(list_request(endpoint, input, &[]))
1389 .await
1390 .unwrap();
1391 let (status, body) = json_response(resp).await;
1392 assert_eq!(status, StatusCode::BAD_REQUEST, "{endpoint} input={input}");
1393 let msg = body["message"].as_str().unwrap_or_default();
1394 assert!(
1395 msg.contains("sh.tangled.repo.issue")
1396 && msg.contains("sh.tangled.repo.pull")
1397 && msg.contains("sh.tangled.string"),
1398 "{endpoint} input={input}: {msg}",
1399 );
1400 }
1401 })
1402 .await;
1403}
1404
1405#[tokio::test]
1406async fn star_endpoints_reject_unrelated_collection() {
1407 let h = Harness::new().await;
1408 let app = router(h.state.clone());
1409 let endpoints = ["sh.tangled.feed.listStars", "sh.tangled.feed.countStars"];
1410 stream::iter(endpoints)
1411 .for_each(|endpoint| {
1412 let app = app.clone();
1413 async move {
1414 let resp = app
1415 .oneshot(list_request(
1416 endpoint,
1417 "at://did:plc:abalone/sh.tangled.knot/k1",
1418 &[],
1419 ))
1420 .await
1421 .unwrap();
1422 let (status, body) = json_response(resp).await;
1423 assert_eq!(status, StatusCode::BAD_REQUEST, "{endpoint}");
1424 let msg = body["message"].as_str().unwrap_or_default();
1425 assert!(msg.contains("sh.tangled.string"), "{endpoint}: {msg}",);
1426 }
1427 })
1428 .await;
1429}
1430
1431#[tokio::test]
1432async fn star_endpoints_reject_repo_uri_subject() {
1433 let h = Harness::new().await;
1434 let app = router(h.state.clone());
1435 let resp = app
1436 .oneshot(list_request(
1437 "sh.tangled.feed.countStars",
1438 "at://did:plc:abalone/sh.tangled.repo/r1",
1439 &[],
1440 ))
1441 .await
1442 .unwrap();
1443 let (status, body) = json_response(resp).await;
1444 assert_eq!(
1445 status,
1446 StatusCode::BAD_REQUEST,
1447 "rkey-form repo URI must be rejected; clients must send the repoDID directly",
1448 );
1449 let msg = body["message"].as_str().unwrap_or_default();
1450 assert!(msg.contains("sh.tangled.string"), "{msg}");
1451}
1452
1453#[tokio::test]
1454async fn star_endpoints_accept_string_subject_form() {
1455 let h = Harness::new().await;
1456 let app = router(h.state.clone());
1457 let resp = app
1458 .oneshot(list_request(
1459 "sh.tangled.feed.countStars",
1460 "at://did:plc:abalone/sh.tangled.string/k1",
1461 &[],
1462 ))
1463 .await
1464 .unwrap();
1465 let (status, body) = json_response(resp).await;
1466 assert_eq!(status, StatusCode::OK);
1467 assert_eq!(body["count"], json!(0));
1468}
1469
1470#[tokio::test]
1471async fn list_after_remove_source_returns_empty_items() {
1472 let h = Harness::new().await;
1473 let subject = at("at://did:plc:abalone");
1474 let source = at("at://did:plc:nel/sh.tangled.feed.star/s1");
1475 h.add_edge(&nsid("sh.tangled.feed.star"), &subject, &source);
1476 h.edges.remove_source(&source);
1477
1478 let app = router(h.state.clone());
1479 let (status, body) = json_response(
1480 app.oneshot(list_request(
1481 "sh.tangled.feed.listStars",
1482 subject.as_ref(),
1483 &[],
1484 ))
1485 .await
1486 .unwrap(),
1487 )
1488 .await;
1489 assert_eq!(status, StatusCode::OK);
1490 assert_eq!(body["items"], json!([]));
1491 assert!(body["cursor"].is_null());
1492}
1493
1494#[tokio::test]
1495async fn list_pulls_hydrates_via_slingshot_when_edges_present() {
1496 let h = Harness::new().await;
1497 let target_did = did("did:plc:abalone");
1498 let subject = at(&format!("at://{}", target_did.as_ref()));
1499 let source_did = did("did:plc:nel");
1500 let rk = rkey("p1");
1501 h.add_edge(
1502 &nsid("sh.tangled.repo.pull"),
1503 &subject,
1504 &at(&format!(
1505 "at://{}/sh.tangled.repo.pull/{}",
1506 source_did.as_ref(),
1507 rk.as_ref()
1508 )),
1509 );
1510 h.mount(
1511 &source_did,
1512 &nsid("sh.tangled.repo.pull"),
1513 &rk,
1514 json!({
1515 "$type": "sh.tangled.repo.pull",
1516 "title": "ship it",
1517 "createdAt": "2026-05-01T00:00:00Z",
1518 "rounds": [],
1519 "target": {"repo": target_did.as_ref(), "branch": "main"},
1520 }),
1521 )
1522 .await;
1523 let app = router(h.state.clone());
1524 let (status, body) = json_response(
1525 app.oneshot(list_request(
1526 "sh.tangled.repo.listPulls",
1527 subject.as_ref(),
1528 &[],
1529 ))
1530 .await
1531 .unwrap(),
1532 )
1533 .await;
1534 assert_eq!(status, StatusCode::OK);
1535 let items = body["items"].as_array().unwrap();
1536 assert_eq!(items.len(), 1);
1537 assert_eq!(items[0]["value"]["title"], json!("ship it"));
1538 assert_eq!(
1539 items[0]["value"]["target"]["repo"],
1540 json!(target_did.as_ref())
1541 );
1542}
1543
1544#[tokio::test]
1545async fn count_pulls_returns_distinct_authors() {
1546 let h = Harness::new().await;
1547 let subject = at("at://did:plc:abalone");
1548 h.add_edge(
1549 &nsid("sh.tangled.repo.pull"),
1550 &subject,
1551 &at("at://did:plc:nel/sh.tangled.repo.pull/p1"),
1552 );
1553 h.add_edge(
1554 &nsid("sh.tangled.repo.pull"),
1555 &subject,
1556 &at("at://did:plc:olaren/sh.tangled.repo.pull/p2"),
1557 );
1558 h.add_edge(
1559 &nsid("sh.tangled.repo.pull"),
1560 &subject,
1561 &at("at://did:plc:nel/sh.tangled.repo.pull/p3"),
1562 );
1563 let app = router(h.state.clone());
1564 let (_, body) = json_response(
1565 app.oneshot(list_request(
1566 "sh.tangled.repo.countPulls",
1567 subject.as_ref(),
1568 &[],
1569 ))
1570 .await
1571 .unwrap(),
1572 )
1573 .await;
1574 assert_eq!(body["count"], json!(3));
1575 assert_eq!(body["distinctAuthors"], json!(2));
1576}
1577
1578#[tokio::test]
1579async fn extractor_to_xrpc_round_trip_for_star() {
1580 let h = Harness::new().await;
1581 let subject_did = did("did:plc:abalone");
1582 let source_did = did("did:plc:nel");
1583 let rk = rkey("s1");
1584 let source = at(&format!(
1585 "at://{}/sh.tangled.feed.star/{}",
1586 source_did.as_ref(),
1587 rk.as_ref()
1588 ));
1589 let body = star_body(&subject_did);
1590 let parsed =
1591 bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.feed.star"), body.clone())
1592 .expect("parse star record");
1593 parsed
1594 .extract_edges(&source)
1595 .expect("extract")
1596 .into_iter()
1597 .for_each(|e| h.edges.add(e));
1598 h.mount(&source_did, &nsid("sh.tangled.feed.star"), &rk, body)
1599 .await;
1600
1601 let app = router(h.state.clone());
1602 let (status, json) = json_response(
1603 app.oneshot(list_request(
1604 "sh.tangled.feed.listStars",
1605 &format!("at://{}", subject_did.as_ref()),
1606 &[],
1607 ))
1608 .await
1609 .unwrap(),
1610 )
1611 .await;
1612 assert_eq!(
1613 status,
1614 StatusCode::OK,
1615 "extractor key must match handler subject, body was {json}",
1616 );
1617 let items = json["items"].as_array().unwrap();
1618 assert_eq!(items.len(), 1, "expected exactly one star edge");
1619 assert_eq!(
1620 items[0]["value"]["subject"]["did"],
1621 json!(subject_did.as_ref())
1622 );
1623}
1624
1625#[tokio::test]
1626async fn list_issues_includes_state_comment_count_and_state_updated_at() {
1627 let h = Harness::new().await;
1628 let repo = did("did:plc:limpet");
1629 let subject = at(&format!("at://{}", repo.as_ref()));
1630 let issue_uri = at("at://did:plc:nel/sh.tangled.repo.issue/i1");
1631 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri);
1632 h.mount(
1633 &did("did:plc:nel"),
1634 &nsid("sh.tangled.repo.issue"),
1635 &rkey("i1"),
1636 issue_body(&repo, "hi"),
1637 )
1638 .await;
1639 h.add_edge(
1640 &nsid("sh.tangled.feed.comment"),
1641 &issue_uri,
1642 &at("at://did:plc:olaren/sh.tangled.feed.comment/c1"),
1643 );
1644 h.add_edge(
1645 &nsid("sh.tangled.feed.comment"),
1646 &issue_uri,
1647 &at("at://did:plc:teq/sh.tangled.feed.comment/c2"),
1648 );
1649
1650 h.state.issue_states.upsert(
1651 at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"),
1652 issue_uri.clone(),
1653 1_777_593_600_000_000,
1654 IssueStateKind::Open,
1655 );
1656 h.state.issue_states.upsert(
1657 at("at://did:plc:nel/sh.tangled.repo.issue.state/s2"),
1658 issue_uri.clone(),
1659 1_777_593_700_000_000,
1660 IssueStateKind::Closed,
1661 );
1662
1663 let app = router(h.state.clone());
1664 let resp = app
1665 .oneshot(list_request(
1666 "sh.tangled.repo.listIssues",
1667 subject.as_ref(),
1668 &[],
1669 ))
1670 .await
1671 .unwrap();
1672 let (status, body) = json_response(resp).await;
1673 assert_eq!(status, StatusCode::OK);
1674 let item = &body["items"][0];
1675 assert_eq!(item["state"], json!("closed"));
1676 assert_eq!(item["commentCount"], json!(2));
1677 let updated = item["stateUpdatedAt"]
1678 .as_str()
1679 .expect("stateUpdatedAt must serialize as RFC3339 string");
1680 assert!(
1681 updated.starts_with("2026-"),
1682 "expected 2026 timestamp, got {updated}"
1683 );
1684}
1685
1686#[tokio::test]
1687async fn list_issues_defaults_to_open_when_no_state_record() {
1688 let h = Harness::new().await;
1689 let repo = did("did:plc:limpet");
1690 let subject = at(&format!("at://{}", repo.as_ref()));
1691 let issue_uri = at("at://did:plc:nel/sh.tangled.repo.issue/i1");
1692 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri);
1693 h.mount(
1694 &did("did:plc:nel"),
1695 &nsid("sh.tangled.repo.issue"),
1696 &rkey("i1"),
1697 issue_body(&repo, "no state yet"),
1698 )
1699 .await;
1700
1701 let app = router(h.state.clone());
1702 let (_status, body) = json_response(
1703 app.oneshot(list_request(
1704 "sh.tangled.repo.listIssues",
1705 subject.as_ref(),
1706 &[],
1707 ))
1708 .await
1709 .unwrap(),
1710 )
1711 .await;
1712 let item = &body["items"][0];
1713 assert_eq!(
1714 item["state"],
1715 json!("open"),
1716 "absent state record defaults to open"
1717 );
1718 assert!(
1719 item.get("stateUpdatedAt").is_none(),
1720 "stateUpdatedAt must be absent without a state record",
1721 );
1722 assert_eq!(item["commentCount"], json!(0));
1723}
1724
1725#[tokio::test]
1726async fn list_issues_author_filter_restricts_to_matching_did() {
1727 let h = Harness::new().await;
1728 let repo = did("did:plc:limpet");
1729 let subject = at(&format!("at://{}", repo.as_ref()));
1730 let owners = [
1731 ("did:plc:nel", "n1"),
1732 ("did:plc:nel", "n2"),
1733 ("did:plc:olaren", "o1"),
1734 ("did:plc:olaren", "o2"),
1735 ];
1736 stream::iter(owners)
1737 .for_each(|(d, r)| {
1738 let h = &h;
1739 let subject = subject.clone();
1740 let repo = repo.clone();
1741 async move {
1742 let d_did = did(d);
1743 let rk = rkey(r);
1744 h.add_edge(
1745 &nsid("sh.tangled.repo.issue"),
1746 &subject,
1747 &at(&format!(
1748 "at://{}/sh.tangled.repo.issue/{}",
1749 d_did.as_ref(),
1750 rk.as_ref()
1751 )),
1752 );
1753 h.mount(
1754 &d_did,
1755 &nsid("sh.tangled.repo.issue"),
1756 &rk,
1757 issue_body(&repo, &format!("issue-{}", rk.as_ref())),
1758 )
1759 .await;
1760 }
1761 })
1762 .await;
1763
1764 let app = router(h.state.clone());
1765 let (status, body) = json_response(
1766 app.oneshot(list_request(
1767 "sh.tangled.repo.listIssues",
1768 subject.as_ref(),
1769 &[("author", "did:plc:nel")],
1770 ))
1771 .await
1772 .unwrap(),
1773 )
1774 .await;
1775 assert_eq!(status, StatusCode::OK);
1776 let items = body["items"].as_array().expect("items array");
1777 assert_eq!(items.len(), 2, "two issues authored by nel");
1778 let all_nel = items
1779 .iter()
1780 .all(|i| i["uri"].as_str().unwrap().starts_with("at://did:plc:nel/"));
1781 assert!(all_nel, "every returned uri must be authored by nel");
1782}
1783
1784#[tokio::test]
1785async fn list_issues_invalid_author_returns_400() {
1786 let h = Harness::new().await;
1787 let subject = "at://did:plc:limpet".to_owned();
1788 let app = router(h.state.clone());
1789 let resp = app
1790 .oneshot(list_request(
1791 "sh.tangled.repo.listIssues",
1792 &subject,
1793 &[("author", "not-a-did")],
1794 ))
1795 .await
1796 .unwrap();
1797 assert_eq!(resp.status(), StatusCode::BAD_REQUEST);
1798}
1799
1800#[tokio::test]
1801async fn list_pulls_includes_merged_state_and_comment_count() {
1802 let h = Harness::new().await;
1803 let repo = did("did:plc:limpet");
1804 let subject = at(&format!("at://{}", repo.as_ref()));
1805 let pull_uri = at("at://did:plc:nel/sh.tangled.repo.pull/p1");
1806 h.add_edge(&nsid("sh.tangled.repo.pull"), &subject, &pull_uri);
1807 h.mount(
1808 &did("did:plc:nel"),
1809 &nsid("sh.tangled.repo.pull"),
1810 &rkey("p1"),
1811 pull_body(&repo, "fix bug"),
1812 )
1813 .await;
1814 h.add_edge(
1815 &nsid("sh.tangled.feed.comment"),
1816 &pull_uri,
1817 &at("at://did:plc:teq/sh.tangled.feed.comment/c1"),
1818 );
1819 h.state.pull_statuses.upsert(
1820 at("at://did:plc:nel/sh.tangled.repo.pull.status/s1"),
1821 pull_uri.clone(),
1822 1_777_593_600_000_000,
1823 PullStatusKind::Open,
1824 );
1825 h.state.pull_statuses.upsert(
1826 at("at://did:plc:nel/sh.tangled.repo.pull.status/s2"),
1827 pull_uri.clone(),
1828 1_777_593_800_000_000,
1829 PullStatusKind::Merged,
1830 );
1831
1832 let app = router(h.state.clone());
1833 let (status, body) = json_response(
1834 app.oneshot(list_request(
1835 "sh.tangled.repo.listPulls",
1836 subject.as_ref(),
1837 &[],
1838 ))
1839 .await
1840 .unwrap(),
1841 )
1842 .await;
1843 assert_eq!(status, StatusCode::OK);
1844 let item = &body["items"][0];
1845 assert_eq!(item["state"], json!("merged"));
1846 assert_eq!(item["commentCount"], json!(1));
1847}
1848
1849#[tokio::test]
1850async fn list_issues_state_filter_open_includes_records_without_state() {
1851 let h = Harness::new().await;
1852 let repo = did("did:plc:limpet");
1853 let subject = at(&format!("at://{}", repo.as_ref()));
1854 let issue_uri = at("at://did:plc:nel/sh.tangled.repo.issue/i1");
1855 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri);
1856 h.mount(
1857 &did("did:plc:nel"),
1858 &nsid("sh.tangled.repo.issue"),
1859 &rkey("i1"),
1860 issue_body(&repo, "fresh"),
1861 )
1862 .await;
1863
1864 let app = router(h.state.clone());
1865 let (status, body) = json_response(
1866 app.oneshot(list_request(
1867 "sh.tangled.repo.listIssues",
1868 subject.as_ref(),
1869 &[("state", "open")],
1870 ))
1871 .await
1872 .unwrap(),
1873 )
1874 .await;
1875 assert_eq!(status, StatusCode::OK);
1876 let items = body["items"].as_array().expect("items array");
1877 assert_eq!(
1878 items.len(),
1879 1,
1880 "absent state record still matches state=open"
1881 );
1882}
1883
1884#[tokio::test]
1885async fn list_issues_state_filter_ignores_third_party_state_source() {
1886 let h = Harness::new().await;
1887 let repo = did("did:plc:limpet");
1888 let subject = at(&format!("at://{}", repo.as_ref()));
1889 let issue_uri = at("at://did:plc:nel/sh.tangled.repo.issue/i1");
1890 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri);
1891 h.mount(
1892 &did("did:plc:nel"),
1893 &nsid("sh.tangled.repo.issue"),
1894 &rkey("i1"),
1895 issue_body(&repo, "open issue"),
1896 )
1897 .await;
1898 h.state.issue_states.upsert(
1899 at("at://did:plc:nautilus/sh.tangled.repo.issue.state/spoof"),
1900 issue_uri.clone(),
1901 1_777_593_800_000_000,
1902 IssueStateKind::Closed,
1903 );
1904
1905 let app = router(h.state.clone());
1906 let (status, body) = json_response(
1907 app.oneshot(list_request(
1908 "sh.tangled.repo.listIssues",
1909 subject.as_ref(),
1910 &[("state", "open")],
1911 ))
1912 .await
1913 .unwrap(),
1914 )
1915 .await;
1916 assert_eq!(status, StatusCode::OK);
1917 let items = body["items"].as_array().expect("items array");
1918 assert_eq!(
1919 items.len(),
1920 1,
1921 "third-party Closed record must not flip filter result for state=open",
1922 );
1923 assert_eq!(items[0]["state"], json!("open"));
1924 assert!(
1925 items[0].get("stateUpdatedAt").is_none(),
1926 "third-party state source must not surface stateUpdatedAt",
1927 );
1928}
1929
1930#[tokio::test]
1931async fn list_pulls_status_filter_ignores_third_party_status_source() {
1932 let h = Harness::new().await;
1933 let repo = did("did:plc:limpet");
1934 let subject = at(&format!("at://{}", repo.as_ref()));
1935 let pull_uri = at("at://did:plc:nel/sh.tangled.repo.pull/p1");
1936 h.add_edge(&nsid("sh.tangled.repo.pull"), &subject, &pull_uri);
1937 h.mount(
1938 &did("did:plc:nel"),
1939 &nsid("sh.tangled.repo.pull"),
1940 &rkey("p1"),
1941 pull_body(&repo, "wip"),
1942 )
1943 .await;
1944 h.state.pull_statuses.upsert(
1945 at("at://did:plc:nautilus/sh.tangled.repo.pull.status/spoof"),
1946 pull_uri.clone(),
1947 1_777_593_800_000_000,
1948 PullStatusKind::Merged,
1949 );
1950
1951 let app = router(h.state.clone());
1952 let (status, body) = json_response(
1953 app.oneshot(list_request(
1954 "sh.tangled.repo.listPulls",
1955 subject.as_ref(),
1956 &[("status", "merged")],
1957 ))
1958 .await
1959 .unwrap(),
1960 )
1961 .await;
1962 assert_eq!(status, StatusCode::OK);
1963 let items = body["items"].as_array().expect("items array");
1964 assert_eq!(
1965 items.len(),
1966 0,
1967 "third-party Merged record must not satisfy status=merged"
1968 );
1969}
1970
1971#[tokio::test]
1972async fn list_issues_state_filter_accepts_repo_owner_state_source() {
1973 let h = Harness::new().await;
1974 let repo_owner = did("did:plc:limpet");
1975 let subject = at(&format!("at://{}", repo_owner.as_ref()));
1976 let issue_uri = at("at://did:plc:nel/sh.tangled.repo.issue/i1");
1977 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri);
1978 h.mount(
1979 &did("did:plc:nel"),
1980 &nsid("sh.tangled.repo.issue"),
1981 &rkey("i1"),
1982 issue_body(&repo_owner, "owner closed"),
1983 )
1984 .await;
1985 h.state.issue_states.upsert(
1986 at("at://did:plc:limpet/sh.tangled.repo.issue.state/legit"),
1987 issue_uri.clone(),
1988 1_777_593_800_000_000,
1989 IssueStateKind::Closed,
1990 );
1991
1992 let app = router(h.state.clone());
1993 let (status, body) = json_response(
1994 app.oneshot(list_request(
1995 "sh.tangled.repo.listIssues",
1996 subject.as_ref(),
1997 &[("state", "closed")],
1998 ))
1999 .await
2000 .unwrap(),
2001 )
2002 .await;
2003 assert_eq!(status, StatusCode::OK);
2004 let items = body["items"].as_array().expect("items array");
2005 assert_eq!(
2006 items.len(),
2007 1,
2008 "repo-owner state record must satisfy state=closed"
2009 );
2010 assert_eq!(items[0]["state"], json!("closed"));
2011}
2012
2013#[tokio::test]
2014async fn list_issues_order_asc_returns_oldest_first() {
2015 let h = Harness::new().await;
2016 let repo = did("did:plc:limpet");
2017 let subject = at(&format!("at://{}", repo.as_ref()));
2018 let rkeys = ["a", "b", "c"];
2019 stream::iter(rkeys)
2020 .for_each(|r| {
2021 let h = &h;
2022 let subject = subject.clone();
2023 let repo = repo.clone();
2024 async move {
2025 let rk = rkey(r);
2026 let issue_uri = at(&format!(
2027 "at://did:plc:nel/sh.tangled.repo.issue/{}",
2028 rk.as_ref()
2029 ));
2030 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri);
2031 h.mount(
2032 &did("did:plc:nel"),
2033 &nsid("sh.tangled.repo.issue"),
2034 &rk,
2035 issue_body(&repo, &format!("issue-{}", rk.as_ref())),
2036 )
2037 .await;
2038 }
2039 })
2040 .await;
2041
2042 let app = router(h.state.clone());
2043 let (_, asc) = json_response(
2044 app.clone()
2045 .oneshot(list_request(
2046 "sh.tangled.repo.listIssues",
2047 subject.as_ref(),
2048 &[("order", "asc")],
2049 ))
2050 .await
2051 .unwrap(),
2052 )
2053 .await;
2054 let (_, desc) = json_response(
2055 app.oneshot(list_request(
2056 "sh.tangled.repo.listIssues",
2057 subject.as_ref(),
2058 &[("order", "desc")],
2059 ))
2060 .await
2061 .unwrap(),
2062 )
2063 .await;
2064 let asc_uris: Vec<_> = asc["items"]
2065 .as_array()
2066 .unwrap()
2067 .iter()
2068 .map(|i| i["uri"].as_str().unwrap().to_owned())
2069 .collect();
2070 let desc_uris: Vec<_> = desc["items"]
2071 .as_array()
2072 .unwrap()
2073 .iter()
2074 .map(|i| i["uri"].as_str().unwrap().to_owned())
2075 .collect();
2076 let mut reversed = asc_uris.clone();
2077 reversed.reverse();
2078 assert_eq!(asc_uris.len(), 3);
2079 assert_eq!(desc_uris, reversed, "desc must be exact reverse of asc");
2080}
2081
2082#[tokio::test]
2083async fn list_issues_by_state_filter_narrows_results() {
2084 let h = Harness::new().await;
2085 let author = did("did:plc:nel");
2086 let repo = did("did:plc:limpet");
2087 let open_uri = at("at://did:plc:nel/sh.tangled.repo.issue/open1");
2088 let closed_uri = at("at://did:plc:nel/sh.tangled.repo.issue/closed1");
2089 let author_subject = at(&format!("at://{}", author.as_ref()));
2090 h.edges.add(Edge {
2091 kind: nsid("sh.tangled.repo.issue.by"),
2092 subject: SubjectRef::Did(author.clone()),
2093 source: open_uri.clone(),
2094 sort_micros: next_sort_micros(),
2095 });
2096 h.edges.add(Edge {
2097 kind: nsid("sh.tangled.repo.issue.by"),
2098 subject: SubjectRef::Did(author.clone()),
2099 source: closed_uri.clone(),
2100 sort_micros: next_sort_micros(),
2101 });
2102 h.mount(
2103 &author,
2104 &nsid("sh.tangled.repo.issue"),
2105 &rkey("open1"),
2106 issue_body(&repo, "still open"),
2107 )
2108 .await;
2109 h.mount(
2110 &author,
2111 &nsid("sh.tangled.repo.issue"),
2112 &rkey("closed1"),
2113 issue_body(&repo, "shut"),
2114 )
2115 .await;
2116 h.state.issue_states.upsert(
2117 at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"),
2118 closed_uri.clone(),
2119 1_777_593_800_000_000,
2120 IssueStateKind::Closed,
2121 );
2122
2123 let app = router(h.state.clone());
2124 let (status, body) = json_response(
2125 app.oneshot(list_request(
2126 "sh.tangled.repo.listIssuesBy",
2127 author_subject.as_ref(),
2128 &[("state", "closed")],
2129 ))
2130 .await
2131 .unwrap(),
2132 )
2133 .await;
2134 assert_eq!(status, StatusCode::OK);
2135 let items = body["items"].as_array().expect("items array");
2136 assert_eq!(
2137 items.len(),
2138 1,
2139 "only the closed issue survives state=closed"
2140 );
2141 assert_eq!(items[0]["uri"], json!(closed_uri.as_ref()));
2142}
2143
2144#[tokio::test]
2145async fn knot_owned_member_is_synthesized_without_slingshot() {
2146 let harness = Harness::new().await;
2147 let knot = bobbin_types::knot_acl::host_to_knot_did("kt.oyster.cafe").unwrap();
2148 let subject = did("did:plc:boltless");
2149 let created = chrono::DateTime::parse_from_rfc3339("2026-06-01T00:00:00Z").unwrap();
2150 let micros = created.timestamp_micros() as u64;
2151 let (source, edges) = bobbin_types::knot_acl::member_upsert(&knot, &subject, micros).unwrap();
2152 harness.edges.upsert_source(&source, edges);
2153 harness.promote_ready(1, 1);
2154
2155 let (status, body) = json_response(
2156 router(harness.state.clone())
2157 .oneshot(list_request(
2158 "sh.tangled.knot.listMembers",
2159 subject.as_ref(),
2160 &[],
2161 ))
2162 .await
2163 .unwrap(),
2164 )
2165 .await;
2166
2167 assert_eq!(status, StatusCode::OK);
2168 let items = body["items"].as_array().expect("items array");
2169 assert_eq!(
2170 items.len(),
2171 1,
2172 "synthesized member must hydrate with no slingshot mock mounted"
2173 );
2174 assert_eq!(items[0]["uri"], json!(source.as_ref()));
2175 assert!(items[0]["cid"].is_null());
2176 assert_eq!(items[0]["value"]["domain"], json!("kt.oyster.cafe"));
2177 assert_eq!(items[0]["value"]["subject"], json!("did:plc:boltless"));
2178 let got = chrono::DateTime::parse_from_rfc3339(
2179 items[0]["value"]["createdAt"]
2180 .as_str()
2181 .expect("createdAt string"),
2182 )
2183 .unwrap();
2184 assert_eq!(got.timestamp_micros(), micros as i64);
2185}
2186
2187#[tokio::test]
2188async fn knot_owned_member_lists_by_knot_did() {
2189 let harness = Harness::new().await;
2190 let knot = bobbin_types::knot_acl::host_to_knot_did("kt.oyster.cafe").unwrap();
2191 let subject = did("did:plc:boltless");
2192 let created = chrono::DateTime::parse_from_rfc3339("2026-06-01T00:00:00Z").unwrap();
2193 let micros = created.timestamp_micros() as u64;
2194 let (source, edges) = bobbin_types::knot_acl::member_upsert(&knot, &subject, micros).unwrap();
2195 harness.edges.upsert_source(&source, edges);
2196 harness.promote_ready(1, 1);
2197
2198 let (status, body) = json_response(
2199 router(harness.state.clone())
2200 .oneshot(list_request(
2201 "sh.tangled.knot.listMembersBy",
2202 knot.as_ref(),
2203 &[],
2204 ))
2205 .await
2206 .unwrap(),
2207 )
2208 .await;
2209
2210 assert_eq!(status, StatusCode::OK);
2211 let items = body["items"].as_array().expect("items array");
2212 assert_eq!(items.len(), 1);
2213 assert_eq!(items[0]["uri"], json!(source.as_ref()));
2214 assert!(items[0]["cid"].is_null());
2215 assert_eq!(items[0]["value"]["domain"], json!("kt.oyster.cafe"));
2216 assert_eq!(items[0]["value"]["subject"], json!("did:plc:boltless"));
2217}
2218
2219#[tokio::test]
2220async fn knot_owned_collaborator_is_synthesized_without_slingshot() {
2221 let harness = Harness::new().await;
2222 let repo = did("did:plc:scallop");
2223 let subject = did("did:plc:olaren");
2224 let created = chrono::DateTime::parse_from_rfc3339("2026-06-03T12:00:00Z").unwrap();
2225 let micros = created.timestamp_micros() as u64;
2226 let (source, edges) =
2227 bobbin_types::knot_acl::collaborator_upsert(&repo, &subject, micros).unwrap();
2228 harness.edges.upsert_source(&source, edges);
2229 harness.promote_ready(1, 1);
2230
2231 let (status, body) = json_response(
2232 router(harness.state.clone())
2233 .oneshot(list_request(
2234 "sh.tangled.repo.listCollaborators",
2235 repo.as_ref(),
2236 &[],
2237 ))
2238 .await
2239 .unwrap(),
2240 )
2241 .await;
2242
2243 assert_eq!(status, StatusCode::OK);
2244 let items = body["items"].as_array().expect("items array");
2245 assert_eq!(items.len(), 1);
2246 assert_eq!(items[0]["uri"], json!(source.as_ref()));
2247 assert!(items[0]["cid"].is_null());
2248 assert_eq!(items[0]["value"]["repo"], json!("did:plc:scallop"));
2249 assert_eq!(items[0]["value"]["subject"], json!("did:plc:olaren"));
2250}
2251
2252#[tokio::test]
2253async fn knot_owned_collaborator_lists_by_subject_did() {
2254 let harness = Harness::new().await;
2255 let repo = did("did:plc:scallop");
2256 let subject = did("did:plc:olaren");
2257 let created = chrono::DateTime::parse_from_rfc3339("2026-06-03T12:00:00Z").unwrap();
2258 let micros = created.timestamp_micros() as u64;
2259 let (source, edges) =
2260 bobbin_types::knot_acl::collaborator_upsert(&repo, &subject, micros).unwrap();
2261 harness.edges.upsert_source(&source, edges);
2262 harness.promote_ready(1, 1);
2263
2264 let (status, body) = json_response(
2265 router(harness.state.clone())
2266 .oneshot(list_request(
2267 "sh.tangled.repo.listCollaboratorsBy",
2268 subject.as_ref(),
2269 &[],
2270 ))
2271 .await
2272 .unwrap(),
2273 )
2274 .await;
2275
2276 assert_eq!(status, StatusCode::OK);
2277 let items = body["items"].as_array().expect("items array");
2278 assert_eq!(items.len(), 1);
2279 assert_eq!(items[0]["uri"], json!(source.as_ref()));
2280 assert!(items[0]["cid"].is_null());
2281 assert_eq!(items[0]["value"]["repo"], json!("did:plc:scallop"));
2282 assert_eq!(items[0]["value"]["subject"], json!("did:plc:olaren"));
2283}