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
71 kB 2353 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 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}