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