This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / bobbin / crates / xrpc / tests / aggregation.rs
69 kB 2283 lines
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}