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
62 kB 2061 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={}", encode(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 upstream_failure_during_hydration_drops_only_that_item() { 610 let h = Harness::new().await; 611 let subject = at("at://did:plc:squid"); 612 let kind = nsid("sh.tangled.repo.issue"); 613 let repo = did("did:plc:squid"); 614 h.add_edge( 615 &kind, 616 &subject, 617 &at("at://did:plc:nel/sh.tangled.repo.issue/ok"), 618 ); 619 h.add_edge( 620 &kind, 621 &subject, 622 &at("at://did:plc:teq/sh.tangled.repo.issue/flaky"), 623 ); 624 h.mount( 625 &did("did:plc:nel"), 626 &kind, 627 &rkey("ok"), 628 issue_body(&repo, "kelp survey"), 629 ) 630 .await; 631 Mock::given(method("GET")) 632 .and(path("/xrpc/com.atproto.repo.getRecord")) 633 .and(query_param("repo", "did:plc:teq")) 634 .and(query_param("collection", "sh.tangled.repo.issue")) 635 .and(query_param("rkey", "flaky")) 636 .respond_with(ResponseTemplate::new(503)) 637 .mount(&h.server) 638 .await; 639 640 let app = router(h.state.clone()); 641 let (status, body) = json_response( 642 app.oneshot(list_request( 643 "sh.tangled.repo.listIssues", 644 subject.as_ref(), 645 &[], 646 )) 647 .await 648 .unwrap(), 649 ) 650 .await; 651 assert_eq!(status, StatusCode::OK); 652 let items = body["items"].as_array().expect("items array"); 653 assert_eq!(items.len(), 1, "flaky item dropped, healthy sibling kept"); 654 assert_eq!( 655 items[0]["uri"].as_str().unwrap(), 656 "at://did:plc:nel/sh.tangled.repo.issue/ok", 657 ); 658} 659 660#[tokio::test] 661async fn transient_failure_keeps_edge_so_count_stays_whole() { 662 let h = Harness::new().await; 663 let subject = at("at://did:plc:squid"); 664 let kind = nsid("sh.tangled.repo.issue"); 665 let repo = did("did:plc:squid"); 666 h.add_edge( 667 &kind, 668 &subject, 669 &at("at://did:plc:nel/sh.tangled.repo.issue/ok"), 670 ); 671 h.add_edge( 672 &kind, 673 &subject, 674 &at("at://did:plc:teq/sh.tangled.repo.issue/flaky"), 675 ); 676 h.mount( 677 &did("did:plc:nel"), 678 &kind, 679 &rkey("ok"), 680 issue_body(&repo, "kelp survey"), 681 ) 682 .await; 683 Mock::given(method("GET")) 684 .and(path("/xrpc/com.atproto.repo.getRecord")) 685 .and(query_param("repo", "did:plc:teq")) 686 .and(query_param("collection", "sh.tangled.repo.issue")) 687 .and(query_param("rkey", "flaky")) 688 .respond_with(ResponseTemplate::new(503)) 689 .mount(&h.server) 690 .await; 691 692 let app = router(h.state.clone()); 693 let (status, body) = json_response( 694 app.clone() 695 .oneshot(list_request( 696 "sh.tangled.repo.listIssues", 697 subject.as_ref(), 698 &[], 699 )) 700 .await 701 .unwrap(), 702 ) 703 .await; 704 assert_eq!(status, StatusCode::OK); 705 assert_eq!(body["items"].as_array().unwrap().len(), 1); 706 707 let (cstatus, cbody) = json_response( 708 app.oneshot(list_request( 709 "sh.tangled.repo.countIssues", 710 subject.as_ref(), 711 &[], 712 )) 713 .await 714 .unwrap(), 715 ) 716 .await; 717 assert_eq!(cstatus, StatusCode::OK); 718 assert_eq!( 719 cbody["count"], 720 json!(2), 721 "a transient 503 must not evict the edge, count stays whole", 722 ); 723} 724 725#[tokio::test] 726async fn gone_item_is_evicted_so_count_converges_to_list() { 727 let h = Harness::new().await; 728 let subject = at("at://did:plc:squid"); 729 let kind = nsid("sh.tangled.repo.issue"); 730 let repo = did("did:plc:squid"); 731 h.add_edge( 732 &kind, 733 &subject, 734 &at("at://did:plc:nel/sh.tangled.repo.issue/ok"), 735 ); 736 h.add_edge( 737 &kind, 738 &subject, 739 &at("at://did:plc:teq/sh.tangled.repo.issue/gone"), 740 ); 741 h.mount( 742 &did("did:plc:nel"), 743 &kind, 744 &rkey("ok"), 745 issue_body(&repo, "kelp survey"), 746 ) 747 .await; 748 Mock::given(method("GET")) 749 .and(path("/xrpc/com.atproto.repo.getRecord")) 750 .and(query_param("repo", "did:plc:teq")) 751 .and(query_param("collection", "sh.tangled.repo.issue")) 752 .and(query_param("rkey", "gone")) 753 .respond_with(ResponseTemplate::new(404)) 754 .mount(&h.server) 755 .await; 756 757 let app = router(h.state.clone()); 758 let (status, body) = json_response( 759 app.clone() 760 .oneshot(list_request( 761 "sh.tangled.repo.listIssues", 762 subject.as_ref(), 763 &[], 764 )) 765 .await 766 .unwrap(), 767 ) 768 .await; 769 assert_eq!(status, StatusCode::OK); 770 assert_eq!( 771 body["items"].as_array().unwrap().len(), 772 1, 773 "gone item dropped from the page" 774 ); 775 776 let (cstatus, cbody) = json_response( 777 app.oneshot(list_request( 778 "sh.tangled.repo.countIssues", 779 subject.as_ref(), 780 &[], 781 )) 782 .await 783 .unwrap(), 784 ) 785 .await; 786 assert_eq!(cstatus, StatusCode::OK); 787 assert_eq!( 788 cbody["count"], 789 json!(1), 790 "a definitive 404 must evict the dead edge so count matches the list", 791 ); 792} 793 794#[tokio::test] 795async fn handle_authority_subject_is_400() { 796 let h = Harness::new().await; 797 let app = router(h.state.clone()); 798 let cases = [ 799 "sh.tangled.feed.listStars", 800 "sh.tangled.feed.countStars", 801 "sh.tangled.graph.listFollows", 802 "sh.tangled.graph.countFollows", 803 "sh.tangled.repo.listIssues", 804 "sh.tangled.repo.countIssues", 805 "sh.tangled.repo.listPulls", 806 "sh.tangled.repo.countPulls", 807 "sh.tangled.repo.issue.listComments", 808 "sh.tangled.repo.issue.countComments", 809 ]; 810 stream::iter(cases) 811 .for_each(|endpoint| { 812 let app = app.clone(); 813 async move { 814 let resp = app 815 .oneshot(list_request(endpoint, "at://oyster.cafe", &[])) 816 .await 817 .unwrap(); 818 let (status, body) = json_response(resp).await; 819 assert_eq!(status, StatusCode::BAD_REQUEST, "{endpoint}"); 820 assert_eq!(body["error"], "InvalidRequest", "{endpoint}"); 821 assert!( 822 body["message"] 823 .as_str() 824 .unwrap_or_default() 825 .contains("did, not a handle"), 826 "{endpoint}: {}", 827 body["message"] 828 ); 829 } 830 }) 831 .await; 832} 833 834#[tokio::test] 835async fn empty_subject_is_400() { 836 let h = Harness::new().await; 837 let app = router(h.state.clone()); 838 let resp = app 839 .oneshot(list_request("sh.tangled.repo.listIssues", "", &[])) 840 .await 841 .unwrap(); 842 let (status, body) = json_response(resp).await; 843 assert_eq!(status, StatusCode::BAD_REQUEST); 844 assert_eq!(body["error"], "InvalidRequest"); 845} 846 847#[tokio::test] 848async fn limit_below_min_or_above_max_is_400() { 849 let h = Harness::new().await; 850 let app = router(h.state.clone()); 851 let cases = [("0", "below"), ("1001", "above")]; 852 stream::iter(cases) 853 .for_each(|(limit, label)| { 854 let app = app.clone(); 855 async move { 856 let resp = app 857 .oneshot(list_request( 858 "sh.tangled.repo.listIssues", 859 "at://did:plc:abalone", 860 &[("limit", limit)], 861 )) 862 .await 863 .unwrap(); 864 let (status, body) = json_response(resp).await; 865 assert_eq!(status, StatusCode::BAD_REQUEST, "limit {label}"); 866 assert_eq!(body["error"], "InvalidRequest", "limit {label}"); 867 } 868 }) 869 .await; 870} 871 872#[tokio::test] 873async fn count_after_remove_source_returns_zero() { 874 let h = Harness::new().await; 875 let subject = at("at://did:plc:abalone"); 876 let source = at("at://did:plc:nel/sh.tangled.feed.star/s1"); 877 h.add_edge(&nsid("sh.tangled.feed.star"), &subject, &source); 878 h.edges.remove_source(&source); 879 880 let app = router(h.state.clone()); 881 let (_, body) = json_response( 882 app.oneshot(list_request( 883 "sh.tangled.feed.countStars", 884 subject.as_ref(), 885 &[], 886 )) 887 .await 888 .unwrap(), 889 ) 890 .await; 891 assert_eq!(body["count"], json!(0)); 892 assert_eq!(body["distinctAuthors"], json!(0)); 893} 894 895#[tokio::test] 896async fn list_issue_comments_hydrates_end_to_end() { 897 let h = Harness::new().await; 898 let issue_uri = at("at://did:plc:abalone/sh.tangled.repo.issue/i1"); 899 let nel = did("did:plc:nel"); 900 let rk = rkey("c1"); 901 h.add_edge( 902 &nsid("sh.tangled.repo.issue.comment"), 903 &issue_uri, 904 &at(&format!( 905 "at://{}/sh.tangled.repo.issue.comment/{}", 906 nel.as_ref(), 907 rk.as_ref() 908 )), 909 ); 910 h.mount( 911 &nel, 912 &nsid("sh.tangled.repo.issue.comment"), 913 &rk, 914 json!({ 915 "$type": "sh.tangled.repo.issue.comment", 916 "issue": issue_uri.as_ref(), 917 "body": "thoughts", 918 "createdAt": "2026-05-01T00:00:00Z" 919 }), 920 ) 921 .await; 922 923 let app = router(h.state.clone()); 924 let (status, body) = json_response( 925 app.oneshot(list_request( 926 "sh.tangled.repo.issue.listComments", 927 issue_uri.as_ref(), 928 &[], 929 )) 930 .await 931 .unwrap(), 932 ) 933 .await; 934 assert_eq!(status, StatusCode::OK); 935 let items = body["items"].as_array().unwrap(); 936 assert_eq!(items.len(), 1); 937 assert_eq!(items[0]["value"]["body"], json!("thoughts")); 938 assert_eq!(items[0]["value"]["issue"], json!(issue_uri.as_ref())); 939} 940 941#[tokio::test] 942async fn list_item_cid_is_present() { 943 let h = Harness::new().await; 944 let subject = at("at://did:plc:abalone"); 945 let nel = did("did:plc:nel"); 946 h.add_edge( 947 &nsid("sh.tangled.feed.star"), 948 &subject, 949 &at(&format!("at://{}/sh.tangled.feed.star/s1", nel.as_ref())), 950 ); 951 h.mount( 952 &nel, 953 &nsid("sh.tangled.feed.star"), 954 &rkey("s1"), 955 star_body(&did("did:plc:abalone")), 956 ) 957 .await; 958 959 let app = router(h.state.clone()); 960 let (_, body) = json_response( 961 app.oneshot(list_request( 962 "sh.tangled.feed.listStars", 963 subject.as_ref(), 964 &[], 965 )) 966 .await 967 .unwrap(), 968 ) 969 .await; 970 let item = &body["items"][0]; 971 assert!( 972 item.as_object().unwrap().contains_key("cid"), 973 "list items must mirror getRecord output shape and include cid" 974 ); 975 assert_eq!(item["cid"], json!(CID)); 976} 977 978#[tokio::test] 979async fn count_issue_comments_subjects_on_issue_uri() { 980 let h = Harness::new().await; 981 let issue_uri = at("at://did:plc:abalone/sh.tangled.repo.issue/i1"); 982 h.add_edge( 983 &nsid("sh.tangled.repo.issue.comment"), 984 &issue_uri, 985 &at("at://did:plc:nel/sh.tangled.repo.issue.comment/c1"), 986 ); 987 h.add_edge( 988 &nsid("sh.tangled.repo.issue.comment"), 989 &issue_uri, 990 &at("at://did:plc:olaren/sh.tangled.repo.issue.comment/c2"), 991 ); 992 993 let app = router(h.state.clone()); 994 let (status, body) = json_response( 995 app.oneshot(list_request( 996 "sh.tangled.repo.issue.countComments", 997 issue_uri.as_ref(), 998 &[], 999 )) 1000 .await 1001 .unwrap(), 1002 ) 1003 .await; 1004 assert_eq!(status, StatusCode::OK); 1005 assert_eq!(body["count"], json!(2)); 1006 assert_eq!(body["distinctAuthors"], json!(2)); 1007} 1008 1009#[tokio::test] 1010async fn list_item_404_dropped_not_404_for_subject() { 1011 let h = Harness::new().await; 1012 let subject = at("at://did:plc:squid"); 1013 let kind = nsid("sh.tangled.repo.issue"); 1014 let repo = did("did:plc:squid"); 1015 h.add_edge( 1016 &kind, 1017 &subject, 1018 &at("at://did:plc:nel/sh.tangled.repo.issue/live"), 1019 ); 1020 h.add_edge( 1021 &kind, 1022 &subject, 1023 &at("at://did:plc:teq/sh.tangled.repo.issue/missing"), 1024 ); 1025 h.mount( 1026 &did("did:plc:nel"), 1027 &kind, 1028 &rkey("live"), 1029 issue_body(&repo, "kelp survives"), 1030 ) 1031 .await; 1032 Mock::given(method("GET")) 1033 .and(path("/xrpc/com.atproto.repo.getRecord")) 1034 .and(query_param("repo", "did:plc:teq")) 1035 .and(query_param("collection", "sh.tangled.repo.issue")) 1036 .and(query_param("rkey", "missing")) 1037 .respond_with(ResponseTemplate::new(404).set_body_json(json!({ 1038 "error": "RecordNotFound", 1039 "message": "could not find record" 1040 }))) 1041 .mount(&h.server) 1042 .await; 1043 1044 let app = router(h.state.clone()); 1045 let (status, body) = json_response( 1046 app.oneshot(list_request( 1047 "sh.tangled.repo.listIssues", 1048 subject.as_ref(), 1049 &[], 1050 )) 1051 .await 1052 .unwrap(), 1053 ) 1054 .await; 1055 assert_eq!( 1056 status, 1057 StatusCode::OK, 1058 "a stale-index 404 drops that item, it must not 404 or 502 the subject's list", 1059 ); 1060 let items = body["items"].as_array().expect("items array"); 1061 assert_eq!(items.len(), 1, "stale 404 item dropped, live sibling kept"); 1062 assert_eq!( 1063 items[0]["uri"].as_str().unwrap(), 1064 "at://did:plc:nel/sh.tangled.repo.issue/live", 1065 ); 1066} 1067 1068#[tokio::test] 1069async fn list_item_with_wrong_type_tag_dropped() { 1070 let h = Harness::new().await; 1071 let subject = at("at://did:plc:squid"); 1072 let kind = nsid("sh.tangled.feed.star"); 1073 h.add_edge( 1074 &kind, 1075 &subject, 1076 &at("at://did:plc:nel/sh.tangled.feed.star/good"), 1077 ); 1078 h.add_edge( 1079 &kind, 1080 &subject, 1081 &at("at://did:plc:teq/sh.tangled.feed.star/wrong"), 1082 ); 1083 h.mount( 1084 &did("did:plc:nel"), 1085 &kind, 1086 &rkey("good"), 1087 star_body(&did("did:plc:squid")), 1088 ) 1089 .await; 1090 Mock::given(method("GET")) 1091 .and(path("/xrpc/com.atproto.repo.getRecord")) 1092 .and(query_param("repo", "did:plc:teq")) 1093 .and(query_param("collection", "sh.tangled.feed.star")) 1094 .and(query_param("rkey", "wrong")) 1095 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 1096 "uri": "at://did:plc:teq/sh.tangled.feed.star/wrong", 1097 "cid": CID, 1098 "value": { 1099 "$type": "sh.tangled.feed.reaction", 1100 "createdAt": "2026-05-01T00:00:00Z", 1101 "subject": "at://did:plc:squid" 1102 } 1103 }))) 1104 .mount(&h.server) 1105 .await; 1106 1107 let app = router(h.state.clone()); 1108 let (status, body) = json_response( 1109 app.oneshot(list_request( 1110 "sh.tangled.feed.listStars", 1111 subject.as_ref(), 1112 &[], 1113 )) 1114 .await 1115 .unwrap(), 1116 ) 1117 .await; 1118 assert_eq!(status, StatusCode::OK); 1119 let items = body["items"].as_array().expect("items array"); 1120 assert_eq!(items.len(), 1, "wrong-type item dropped, valid star kept"); 1121 assert_eq!( 1122 items[0]["uri"].as_str().unwrap(), 1123 "at://did:plc:nel/sh.tangled.feed.star/good", 1124 ); 1125} 1126 1127#[tokio::test] 1128async fn list_item_with_mismatched_collection_dropped() { 1129 let h = Harness::new().await; 1130 let subject = at("at://did:plc:squid"); 1131 let kind = nsid("sh.tangled.repo.issue"); 1132 let repo = did("did:plc:squid"); 1133 h.add_edge( 1134 &kind, 1135 &subject, 1136 &at("at://did:plc:nel/sh.tangled.repo.issue/live"), 1137 ); 1138 h.add_edge( 1139 &kind, 1140 &subject, 1141 &at("at://did:plc:teq/sh.tangled.feed.star/whelk"), 1142 ); 1143 h.mount( 1144 &did("did:plc:nel"), 1145 &kind, 1146 &rkey("live"), 1147 issue_body(&repo, "kelp survives"), 1148 ) 1149 .await; 1150 1151 let app = router(h.state.clone()); 1152 let (status, body) = json_response( 1153 app.oneshot(list_request( 1154 "sh.tangled.repo.listIssues", 1155 subject.as_ref(), 1156 &[], 1157 )) 1158 .await 1159 .unwrap(), 1160 ) 1161 .await; 1162 assert_eq!( 1163 status, 1164 StatusCode::OK, 1165 "a mismatched-collection index edge must not 400 the subject's list", 1166 ); 1167 let items = body["items"].as_array().expect("items array"); 1168 assert_eq!( 1169 items.len(), 1170 1, 1171 "mismatched-collection edge dropped, live sibling kept" 1172 ); 1173 assert_eq!( 1174 items[0]["uri"].as_str().unwrap(), 1175 "at://did:plc:nel/sh.tangled.repo.issue/live", 1176 ); 1177} 1178 1179#[tokio::test] 1180async fn bare_did_endpoints_reject_at_uri_subject() { 1181 let h = Harness::new().await; 1182 let app = router(h.state.clone()); 1183 let cases = [ 1184 "sh.tangled.graph.listFollows", 1185 "sh.tangled.graph.countFollows", 1186 ]; 1187 stream::iter(cases) 1188 .for_each(|endpoint| { 1189 let app = app.clone(); 1190 async move { 1191 let resp = app 1192 .oneshot(list_request( 1193 endpoint, 1194 "at://did:plc:abalone/sh.tangled.repo/r1", 1195 &[], 1196 )) 1197 .await 1198 .unwrap(); 1199 let (status, body) = json_response(resp).await; 1200 assert_eq!(status, StatusCode::BAD_REQUEST, "{endpoint}"); 1201 assert_eq!(body["error"], "InvalidRequest", "{endpoint}"); 1202 assert!( 1203 body["message"] 1204 .as_str() 1205 .unwrap_or_default() 1206 .contains("bare did"), 1207 "{endpoint}: {}", 1208 body["message"], 1209 ); 1210 } 1211 }) 1212 .await; 1213} 1214 1215#[tokio::test] 1216async fn repo_pointing_endpoints_reject_at_uri_subject() { 1217 let h = Harness::new().await; 1218 let app = router(h.state.clone()); 1219 let cases = [ 1220 "sh.tangled.repo.listIssues", 1221 "sh.tangled.repo.countIssues", 1222 "sh.tangled.repo.listPulls", 1223 "sh.tangled.repo.countPulls", 1224 "sh.tangled.repo.listArtifacts", 1225 "sh.tangled.repo.countArtifacts", 1226 ]; 1227 stream::iter(cases) 1228 .for_each(|endpoint| { 1229 let app = app.clone(); 1230 async move { 1231 let resp = app 1232 .oneshot(list_request( 1233 endpoint, 1234 "at://did:plc:abalone/sh.tangled.repo/r1", 1235 &[], 1236 )) 1237 .await 1238 .unwrap(); 1239 let (status, body) = json_response(resp).await; 1240 assert_eq!( 1241 status, 1242 StatusCode::BAD_REQUEST, 1243 "{endpoint} must reject rkey-form subjects since rkeys are unstable; clients must send the repoDID", 1244 ); 1245 assert!( 1246 body["message"] 1247 .as_str() 1248 .unwrap_or_default() 1249 .contains("bare did"), 1250 "{endpoint}: {}", 1251 body["message"], 1252 ); 1253 } 1254 }) 1255 .await; 1256} 1257 1258#[tokio::test] 1259async fn repo_pointing_endpoints_accept_bare_did() { 1260 let h = Harness::new().await; 1261 let app = router(h.state.clone()); 1262 let cases = [ 1263 "sh.tangled.repo.listIssues", 1264 "sh.tangled.repo.countIssues", 1265 "sh.tangled.repo.listPulls", 1266 "sh.tangled.repo.countPulls", 1267 "sh.tangled.repo.listArtifacts", 1268 "sh.tangled.repo.countArtifacts", 1269 ]; 1270 stream::iter(cases) 1271 .for_each(|endpoint| { 1272 let app = app.clone(); 1273 async move { 1274 let resp = app 1275 .oneshot(list_request(endpoint, "did:plc:abalone", &[])) 1276 .await 1277 .unwrap(); 1278 let (status, _body) = json_response(resp).await; 1279 assert_eq!(status, StatusCode::OK, "{endpoint} must accept bare did"); 1280 } 1281 }) 1282 .await; 1283} 1284 1285#[tokio::test] 1286async fn issue_collection_endpoints_reject_bare_did_or_wrong_collection() { 1287 let h = Harness::new().await; 1288 let app = router(h.state.clone()); 1289 let endpoints = [ 1290 "sh.tangled.repo.issue.listComments", 1291 "sh.tangled.repo.issue.countComments", 1292 ]; 1293 let inputs = [ 1294 "at://did:plc:abalone", 1295 "at://did:plc:abalone/sh.tangled.repo/r1", 1296 ]; 1297 let cases = endpoints 1298 .iter() 1299 .copied() 1300 .flat_map(|endpoint| inputs.iter().copied().map(move |input| (endpoint, input))); 1301 stream::iter(cases) 1302 .for_each(|(endpoint, input)| { 1303 let app = app.clone(); 1304 async move { 1305 let resp = app 1306 .oneshot(list_request(endpoint, input, &[])) 1307 .await 1308 .unwrap(); 1309 let (status, body) = json_response(resp).await; 1310 assert_eq!(status, StatusCode::BAD_REQUEST, "{endpoint} input={input}"); 1311 assert!( 1312 body["message"] 1313 .as_str() 1314 .unwrap_or_default() 1315 .contains("sh.tangled.repo.issue/<rkey>"), 1316 "{endpoint} input={input}: {}", 1317 body["message"], 1318 ); 1319 } 1320 }) 1321 .await; 1322} 1323 1324#[tokio::test] 1325async fn star_endpoints_reject_unrelated_collection() { 1326 let h = Harness::new().await; 1327 let app = router(h.state.clone()); 1328 let endpoints = ["sh.tangled.feed.listStars", "sh.tangled.feed.countStars"]; 1329 stream::iter(endpoints) 1330 .for_each(|endpoint| { 1331 let app = app.clone(); 1332 async move { 1333 let resp = app 1334 .oneshot(list_request( 1335 endpoint, 1336 "at://did:plc:abalone/sh.tangled.knot/k1", 1337 &[], 1338 )) 1339 .await 1340 .unwrap(); 1341 let (status, body) = json_response(resp).await; 1342 assert_eq!(status, StatusCode::BAD_REQUEST, "{endpoint}"); 1343 let msg = body["message"].as_str().unwrap_or_default(); 1344 assert!(msg.contains("sh.tangled.string"), "{endpoint}: {msg}",); 1345 } 1346 }) 1347 .await; 1348} 1349 1350#[tokio::test] 1351async fn star_endpoints_reject_repo_uri_subject() { 1352 let h = Harness::new().await; 1353 let app = router(h.state.clone()); 1354 let resp = app 1355 .oneshot(list_request( 1356 "sh.tangled.feed.countStars", 1357 "at://did:plc:abalone/sh.tangled.repo/r1", 1358 &[], 1359 )) 1360 .await 1361 .unwrap(); 1362 let (status, body) = json_response(resp).await; 1363 assert_eq!( 1364 status, 1365 StatusCode::BAD_REQUEST, 1366 "rkey-form repo URI must be rejected; clients must send the repoDID directly", 1367 ); 1368 let msg = body["message"].as_str().unwrap_or_default(); 1369 assert!(msg.contains("sh.tangled.string"), "{msg}"); 1370} 1371 1372#[tokio::test] 1373async fn star_endpoints_accept_string_subject_form() { 1374 let h = Harness::new().await; 1375 let app = router(h.state.clone()); 1376 let resp = app 1377 .oneshot(list_request( 1378 "sh.tangled.feed.countStars", 1379 "at://did:plc:abalone/sh.tangled.string/k1", 1380 &[], 1381 )) 1382 .await 1383 .unwrap(); 1384 let (status, body) = json_response(resp).await; 1385 assert_eq!(status, StatusCode::OK); 1386 assert_eq!(body["count"], json!(0)); 1387} 1388 1389#[tokio::test] 1390async fn list_after_remove_source_returns_empty_items() { 1391 let h = Harness::new().await; 1392 let subject = at("at://did:plc:abalone"); 1393 let source = at("at://did:plc:nel/sh.tangled.feed.star/s1"); 1394 h.add_edge(&nsid("sh.tangled.feed.star"), &subject, &source); 1395 h.edges.remove_source(&source); 1396 1397 let app = router(h.state.clone()); 1398 let (status, body) = json_response( 1399 app.oneshot(list_request( 1400 "sh.tangled.feed.listStars", 1401 subject.as_ref(), 1402 &[], 1403 )) 1404 .await 1405 .unwrap(), 1406 ) 1407 .await; 1408 assert_eq!(status, StatusCode::OK); 1409 assert_eq!(body["items"], json!([])); 1410 assert!(body["cursor"].is_null()); 1411} 1412 1413#[tokio::test] 1414async fn list_pulls_hydrates_via_slingshot_when_edges_present() { 1415 let h = Harness::new().await; 1416 let target_did = did("did:plc:abalone"); 1417 let subject = at(&format!("at://{}", target_did.as_ref())); 1418 let source_did = did("did:plc:nel"); 1419 let rk = rkey("p1"); 1420 h.add_edge( 1421 &nsid("sh.tangled.repo.pull"), 1422 &subject, 1423 &at(&format!( 1424 "at://{}/sh.tangled.repo.pull/{}", 1425 source_did.as_ref(), 1426 rk.as_ref() 1427 )), 1428 ); 1429 h.mount( 1430 &source_did, 1431 &nsid("sh.tangled.repo.pull"), 1432 &rk, 1433 json!({ 1434 "$type": "sh.tangled.repo.pull", 1435 "title": "ship it", 1436 "createdAt": "2026-05-01T00:00:00Z", 1437 "rounds": [], 1438 "target": {"repo": target_did.as_ref(), "branch": "main"}, 1439 }), 1440 ) 1441 .await; 1442 let app = router(h.state.clone()); 1443 let (status, body) = json_response( 1444 app.oneshot(list_request( 1445 "sh.tangled.repo.listPulls", 1446 subject.as_ref(), 1447 &[], 1448 )) 1449 .await 1450 .unwrap(), 1451 ) 1452 .await; 1453 assert_eq!(status, StatusCode::OK); 1454 let items = body["items"].as_array().unwrap(); 1455 assert_eq!(items.len(), 1); 1456 assert_eq!(items[0]["value"]["title"], json!("ship it")); 1457 assert_eq!( 1458 items[0]["value"]["target"]["repo"], 1459 json!(target_did.as_ref()) 1460 ); 1461} 1462 1463#[tokio::test] 1464async fn count_pulls_returns_distinct_authors() { 1465 let h = Harness::new().await; 1466 let subject = at("at://did:plc:abalone"); 1467 h.add_edge( 1468 &nsid("sh.tangled.repo.pull"), 1469 &subject, 1470 &at("at://did:plc:nel/sh.tangled.repo.pull/p1"), 1471 ); 1472 h.add_edge( 1473 &nsid("sh.tangled.repo.pull"), 1474 &subject, 1475 &at("at://did:plc:olaren/sh.tangled.repo.pull/p2"), 1476 ); 1477 h.add_edge( 1478 &nsid("sh.tangled.repo.pull"), 1479 &subject, 1480 &at("at://did:plc:nel/sh.tangled.repo.pull/p3"), 1481 ); 1482 let app = router(h.state.clone()); 1483 let (_, body) = json_response( 1484 app.oneshot(list_request( 1485 "sh.tangled.repo.countPulls", 1486 subject.as_ref(), 1487 &[], 1488 )) 1489 .await 1490 .unwrap(), 1491 ) 1492 .await; 1493 assert_eq!(body["count"], json!(3)); 1494 assert_eq!(body["distinctAuthors"], json!(2)); 1495} 1496 1497#[tokio::test] 1498async fn extractor_to_xrpc_round_trip_for_star() { 1499 let h = Harness::new().await; 1500 let subject_did = did("did:plc:abalone"); 1501 let source_did = did("did:plc:nel"); 1502 let rk = rkey("s1"); 1503 let source = at(&format!( 1504 "at://{}/sh.tangled.feed.star/{}", 1505 source_did.as_ref(), 1506 rk.as_ref() 1507 )); 1508 let body = star_body(&subject_did); 1509 let parsed = 1510 bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.feed.star"), body.clone()) 1511 .expect("parse star record"); 1512 parsed 1513 .extract_edges(&source) 1514 .expect("extract") 1515 .into_iter() 1516 .for_each(|e| h.edges.add(e)); 1517 h.mount(&source_did, &nsid("sh.tangled.feed.star"), &rk, body) 1518 .await; 1519 1520 let app = router(h.state.clone()); 1521 let (status, json) = json_response( 1522 app.oneshot(list_request( 1523 "sh.tangled.feed.listStars", 1524 &format!("at://{}", subject_did.as_ref()), 1525 &[], 1526 )) 1527 .await 1528 .unwrap(), 1529 ) 1530 .await; 1531 assert_eq!( 1532 status, 1533 StatusCode::OK, 1534 "extractor key must match handler subject, body was {json}", 1535 ); 1536 let items = json["items"].as_array().unwrap(); 1537 assert_eq!(items.len(), 1, "expected exactly one star edge"); 1538 assert_eq!( 1539 items[0]["value"]["subject"]["did"], 1540 json!(subject_did.as_ref()) 1541 ); 1542} 1543 1544#[tokio::test] 1545async fn list_issues_includes_state_comment_count_and_state_updated_at() { 1546 let h = Harness::new().await; 1547 let repo = did("did:plc:limpet"); 1548 let subject = at(&format!("at://{}", repo.as_ref())); 1549 let issue_uri = at("at://did:plc:nel/sh.tangled.repo.issue/i1"); 1550 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri); 1551 h.mount( 1552 &did("did:plc:nel"), 1553 &nsid("sh.tangled.repo.issue"), 1554 &rkey("i1"), 1555 issue_body(&repo, "hi"), 1556 ) 1557 .await; 1558 h.add_edge( 1559 &nsid("sh.tangled.repo.issue.comment"), 1560 &issue_uri, 1561 &at("at://did:plc:olaren/sh.tangled.repo.issue.comment/c1"), 1562 ); 1563 h.add_edge( 1564 &nsid("sh.tangled.repo.issue.comment"), 1565 &issue_uri, 1566 &at("at://did:plc:teq/sh.tangled.repo.issue.comment/c2"), 1567 ); 1568 1569 h.state.issue_states.upsert( 1570 at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"), 1571 issue_uri.clone(), 1572 1_777_593_600_000_000, 1573 IssueStateKind::Open, 1574 ); 1575 h.state.issue_states.upsert( 1576 at("at://did:plc:nel/sh.tangled.repo.issue.state/s2"), 1577 issue_uri.clone(), 1578 1_777_593_700_000_000, 1579 IssueStateKind::Closed, 1580 ); 1581 1582 let app = router(h.state.clone()); 1583 let resp = app 1584 .oneshot(list_request( 1585 "sh.tangled.repo.listIssues", 1586 subject.as_ref(), 1587 &[], 1588 )) 1589 .await 1590 .unwrap(); 1591 let (status, body) = json_response(resp).await; 1592 assert_eq!(status, StatusCode::OK); 1593 let item = &body["items"][0]; 1594 assert_eq!(item["state"], json!("closed")); 1595 assert_eq!(item["commentCount"], json!(2)); 1596 let updated = item["stateUpdatedAt"] 1597 .as_str() 1598 .expect("stateUpdatedAt must serialize as RFC3339 string"); 1599 assert!( 1600 updated.starts_with("2026-"), 1601 "expected 2026 timestamp, got {updated}" 1602 ); 1603} 1604 1605#[tokio::test] 1606async fn list_issues_defaults_to_open_when_no_state_record() { 1607 let h = Harness::new().await; 1608 let repo = did("did:plc:limpet"); 1609 let subject = at(&format!("at://{}", repo.as_ref())); 1610 let issue_uri = at("at://did:plc:nel/sh.tangled.repo.issue/i1"); 1611 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri); 1612 h.mount( 1613 &did("did:plc:nel"), 1614 &nsid("sh.tangled.repo.issue"), 1615 &rkey("i1"), 1616 issue_body(&repo, "no state yet"), 1617 ) 1618 .await; 1619 1620 let app = router(h.state.clone()); 1621 let (_status, body) = json_response( 1622 app.oneshot(list_request( 1623 "sh.tangled.repo.listIssues", 1624 subject.as_ref(), 1625 &[], 1626 )) 1627 .await 1628 .unwrap(), 1629 ) 1630 .await; 1631 let item = &body["items"][0]; 1632 assert_eq!( 1633 item["state"], 1634 json!("open"), 1635 "absent state record defaults to open" 1636 ); 1637 assert!( 1638 item.get("stateUpdatedAt").is_none(), 1639 "stateUpdatedAt must be absent without a state record", 1640 ); 1641 assert_eq!(item["commentCount"], json!(0)); 1642} 1643 1644#[tokio::test] 1645async fn list_issues_author_filter_restricts_to_matching_did() { 1646 let h = Harness::new().await; 1647 let repo = did("did:plc:limpet"); 1648 let subject = at(&format!("at://{}", repo.as_ref())); 1649 let owners = [ 1650 ("did:plc:nel", "n1"), 1651 ("did:plc:nel", "n2"), 1652 ("did:plc:olaren", "o1"), 1653 ("did:plc:olaren", "o2"), 1654 ]; 1655 stream::iter(owners) 1656 .for_each(|(d, r)| { 1657 let h = &h; 1658 let subject = subject.clone(); 1659 let repo = repo.clone(); 1660 async move { 1661 let d_did = did(d); 1662 let rk = rkey(r); 1663 h.add_edge( 1664 &nsid("sh.tangled.repo.issue"), 1665 &subject, 1666 &at(&format!( 1667 "at://{}/sh.tangled.repo.issue/{}", 1668 d_did.as_ref(), 1669 rk.as_ref() 1670 )), 1671 ); 1672 h.mount( 1673 &d_did, 1674 &nsid("sh.tangled.repo.issue"), 1675 &rk, 1676 issue_body(&repo, &format!("issue-{}", rk.as_ref())), 1677 ) 1678 .await; 1679 } 1680 }) 1681 .await; 1682 1683 let app = router(h.state.clone()); 1684 let (status, body) = json_response( 1685 app.oneshot(list_request( 1686 "sh.tangled.repo.listIssues", 1687 subject.as_ref(), 1688 &[("author", "did:plc:nel")], 1689 )) 1690 .await 1691 .unwrap(), 1692 ) 1693 .await; 1694 assert_eq!(status, StatusCode::OK); 1695 let items = body["items"].as_array().expect("items array"); 1696 assert_eq!(items.len(), 2, "two issues authored by nel"); 1697 let all_nel = items 1698 .iter() 1699 .all(|i| i["uri"].as_str().unwrap().starts_with("at://did:plc:nel/")); 1700 assert!(all_nel, "every returned uri must be authored by nel"); 1701} 1702 1703#[tokio::test] 1704async fn list_issues_invalid_author_returns_400() { 1705 let h = Harness::new().await; 1706 let subject = "at://did:plc:limpet".to_owned(); 1707 let app = router(h.state.clone()); 1708 let resp = app 1709 .oneshot(list_request( 1710 "sh.tangled.repo.listIssues", 1711 &subject, 1712 &[("author", "not-a-did")], 1713 )) 1714 .await 1715 .unwrap(); 1716 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 1717} 1718 1719#[tokio::test] 1720async fn list_pulls_includes_merged_state_and_comment_count() { 1721 let h = Harness::new().await; 1722 let repo = did("did:plc:limpet"); 1723 let subject = at(&format!("at://{}", repo.as_ref())); 1724 let pull_uri = at("at://did:plc:nel/sh.tangled.repo.pull/p1"); 1725 h.add_edge(&nsid("sh.tangled.repo.pull"), &subject, &pull_uri); 1726 h.mount( 1727 &did("did:plc:nel"), 1728 &nsid("sh.tangled.repo.pull"), 1729 &rkey("p1"), 1730 pull_body(&repo, "fix bug"), 1731 ) 1732 .await; 1733 h.add_edge( 1734 &nsid("sh.tangled.repo.pull.comment"), 1735 &pull_uri, 1736 &at("at://did:plc:teq/sh.tangled.repo.pull.comment/c1"), 1737 ); 1738 h.state.pull_statuses.upsert( 1739 at("at://did:plc:nel/sh.tangled.repo.pull.status/s1"), 1740 pull_uri.clone(), 1741 1_777_593_600_000_000, 1742 PullStatusKind::Open, 1743 ); 1744 h.state.pull_statuses.upsert( 1745 at("at://did:plc:nel/sh.tangled.repo.pull.status/s2"), 1746 pull_uri.clone(), 1747 1_777_593_800_000_000, 1748 PullStatusKind::Merged, 1749 ); 1750 1751 let app = router(h.state.clone()); 1752 let (status, body) = json_response( 1753 app.oneshot(list_request( 1754 "sh.tangled.repo.listPulls", 1755 subject.as_ref(), 1756 &[], 1757 )) 1758 .await 1759 .unwrap(), 1760 ) 1761 .await; 1762 assert_eq!(status, StatusCode::OK); 1763 let item = &body["items"][0]; 1764 assert_eq!(item["state"], json!("merged")); 1765 assert_eq!(item["commentCount"], json!(1)); 1766} 1767 1768#[tokio::test] 1769async fn list_issues_state_filter_open_includes_records_without_state() { 1770 let h = Harness::new().await; 1771 let repo = did("did:plc:limpet"); 1772 let subject = at(&format!("at://{}", repo.as_ref())); 1773 let issue_uri = at("at://did:plc:nel/sh.tangled.repo.issue/i1"); 1774 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri); 1775 h.mount( 1776 &did("did:plc:nel"), 1777 &nsid("sh.tangled.repo.issue"), 1778 &rkey("i1"), 1779 issue_body(&repo, "fresh"), 1780 ) 1781 .await; 1782 1783 let app = router(h.state.clone()); 1784 let (status, body) = json_response( 1785 app.oneshot(list_request( 1786 "sh.tangled.repo.listIssues", 1787 subject.as_ref(), 1788 &[("state", "open")], 1789 )) 1790 .await 1791 .unwrap(), 1792 ) 1793 .await; 1794 assert_eq!(status, StatusCode::OK); 1795 let items = body["items"].as_array().expect("items array"); 1796 assert_eq!( 1797 items.len(), 1798 1, 1799 "absent state record still matches state=open" 1800 ); 1801} 1802 1803#[tokio::test] 1804async fn list_issues_state_filter_ignores_third_party_state_source() { 1805 let h = Harness::new().await; 1806 let repo = did("did:plc:limpet"); 1807 let subject = at(&format!("at://{}", repo.as_ref())); 1808 let issue_uri = at("at://did:plc:nel/sh.tangled.repo.issue/i1"); 1809 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri); 1810 h.mount( 1811 &did("did:plc:nel"), 1812 &nsid("sh.tangled.repo.issue"), 1813 &rkey("i1"), 1814 issue_body(&repo, "open issue"), 1815 ) 1816 .await; 1817 h.state.issue_states.upsert( 1818 at("at://did:plc:nautilus/sh.tangled.repo.issue.state/spoof"), 1819 issue_uri.clone(), 1820 1_777_593_800_000_000, 1821 IssueStateKind::Closed, 1822 ); 1823 1824 let app = router(h.state.clone()); 1825 let (status, body) = json_response( 1826 app.oneshot(list_request( 1827 "sh.tangled.repo.listIssues", 1828 subject.as_ref(), 1829 &[("state", "open")], 1830 )) 1831 .await 1832 .unwrap(), 1833 ) 1834 .await; 1835 assert_eq!(status, StatusCode::OK); 1836 let items = body["items"].as_array().expect("items array"); 1837 assert_eq!( 1838 items.len(), 1839 1, 1840 "third-party Closed record must not flip filter result for state=open", 1841 ); 1842 assert_eq!(items[0]["state"], json!("open")); 1843 assert!( 1844 items[0].get("stateUpdatedAt").is_none(), 1845 "third-party state source must not surface stateUpdatedAt", 1846 ); 1847} 1848 1849#[tokio::test] 1850async fn list_pulls_status_filter_ignores_third_party_status_source() { 1851 let h = Harness::new().await; 1852 let repo = did("did:plc:limpet"); 1853 let subject = at(&format!("at://{}", repo.as_ref())); 1854 let pull_uri = at("at://did:plc:nel/sh.tangled.repo.pull/p1"); 1855 h.add_edge(&nsid("sh.tangled.repo.pull"), &subject, &pull_uri); 1856 h.mount( 1857 &did("did:plc:nel"), 1858 &nsid("sh.tangled.repo.pull"), 1859 &rkey("p1"), 1860 pull_body(&repo, "wip"), 1861 ) 1862 .await; 1863 h.state.pull_statuses.upsert( 1864 at("at://did:plc:nautilus/sh.tangled.repo.pull.status/spoof"), 1865 pull_uri.clone(), 1866 1_777_593_800_000_000, 1867 PullStatusKind::Merged, 1868 ); 1869 1870 let app = router(h.state.clone()); 1871 let (status, body) = json_response( 1872 app.oneshot(list_request( 1873 "sh.tangled.repo.listPulls", 1874 subject.as_ref(), 1875 &[("status", "merged")], 1876 )) 1877 .await 1878 .unwrap(), 1879 ) 1880 .await; 1881 assert_eq!(status, StatusCode::OK); 1882 let items = body["items"].as_array().expect("items array"); 1883 assert_eq!( 1884 items.len(), 1885 0, 1886 "third-party Merged record must not satisfy status=merged" 1887 ); 1888} 1889 1890#[tokio::test] 1891async fn list_issues_state_filter_accepts_repo_owner_state_source() { 1892 let h = Harness::new().await; 1893 let repo_owner = did("did:plc:limpet"); 1894 let subject = at(&format!("at://{}", repo_owner.as_ref())); 1895 let issue_uri = at("at://did:plc:nel/sh.tangled.repo.issue/i1"); 1896 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri); 1897 h.mount( 1898 &did("did:plc:nel"), 1899 &nsid("sh.tangled.repo.issue"), 1900 &rkey("i1"), 1901 issue_body(&repo_owner, "owner closed"), 1902 ) 1903 .await; 1904 h.state.issue_states.upsert( 1905 at("at://did:plc:limpet/sh.tangled.repo.issue.state/legit"), 1906 issue_uri.clone(), 1907 1_777_593_800_000_000, 1908 IssueStateKind::Closed, 1909 ); 1910 1911 let app = router(h.state.clone()); 1912 let (status, body) = json_response( 1913 app.oneshot(list_request( 1914 "sh.tangled.repo.listIssues", 1915 subject.as_ref(), 1916 &[("state", "closed")], 1917 )) 1918 .await 1919 .unwrap(), 1920 ) 1921 .await; 1922 assert_eq!(status, StatusCode::OK); 1923 let items = body["items"].as_array().expect("items array"); 1924 assert_eq!( 1925 items.len(), 1926 1, 1927 "repo-owner state record must satisfy state=closed" 1928 ); 1929 assert_eq!(items[0]["state"], json!("closed")); 1930} 1931 1932#[tokio::test] 1933async fn list_issues_order_asc_returns_oldest_first() { 1934 let h = Harness::new().await; 1935 let repo = did("did:plc:limpet"); 1936 let subject = at(&format!("at://{}", repo.as_ref())); 1937 let rkeys = ["a", "b", "c"]; 1938 stream::iter(rkeys) 1939 .for_each(|r| { 1940 let h = &h; 1941 let subject = subject.clone(); 1942 let repo = repo.clone(); 1943 async move { 1944 let rk = rkey(r); 1945 let issue_uri = at(&format!( 1946 "at://did:plc:nel/sh.tangled.repo.issue/{}", 1947 rk.as_ref() 1948 )); 1949 h.add_edge(&nsid("sh.tangled.repo.issue"), &subject, &issue_uri); 1950 h.mount( 1951 &did("did:plc:nel"), 1952 &nsid("sh.tangled.repo.issue"), 1953 &rk, 1954 issue_body(&repo, &format!("issue-{}", rk.as_ref())), 1955 ) 1956 .await; 1957 } 1958 }) 1959 .await; 1960 1961 let app = router(h.state.clone()); 1962 let (_, asc) = json_response( 1963 app.clone() 1964 .oneshot(list_request( 1965 "sh.tangled.repo.listIssues", 1966 subject.as_ref(), 1967 &[("order", "asc")], 1968 )) 1969 .await 1970 .unwrap(), 1971 ) 1972 .await; 1973 let (_, desc) = json_response( 1974 app.oneshot(list_request( 1975 "sh.tangled.repo.listIssues", 1976 subject.as_ref(), 1977 &[("order", "desc")], 1978 )) 1979 .await 1980 .unwrap(), 1981 ) 1982 .await; 1983 let asc_uris: Vec<_> = asc["items"] 1984 .as_array() 1985 .unwrap() 1986 .iter() 1987 .map(|i| i["uri"].as_str().unwrap().to_owned()) 1988 .collect(); 1989 let desc_uris: Vec<_> = desc["items"] 1990 .as_array() 1991 .unwrap() 1992 .iter() 1993 .map(|i| i["uri"].as_str().unwrap().to_owned()) 1994 .collect(); 1995 let mut reversed = asc_uris.clone(); 1996 reversed.reverse(); 1997 assert_eq!(asc_uris.len(), 3); 1998 assert_eq!(desc_uris, reversed, "desc must be exact reverse of asc"); 1999} 2000 2001#[tokio::test] 2002async fn list_issues_by_state_filter_narrows_results() { 2003 let h = Harness::new().await; 2004 let author = did("did:plc:nel"); 2005 let repo = did("did:plc:limpet"); 2006 let open_uri = at("at://did:plc:nel/sh.tangled.repo.issue/open1"); 2007 let closed_uri = at("at://did:plc:nel/sh.tangled.repo.issue/closed1"); 2008 let author_subject = at(&format!("at://{}", author.as_ref())); 2009 h.edges.add(Edge { 2010 kind: nsid("sh.tangled.repo.issue.by"), 2011 subject: SubjectRef::Did(author.clone()), 2012 source: open_uri.clone(), 2013 sort_micros: next_sort_micros(), 2014 }); 2015 h.edges.add(Edge { 2016 kind: nsid("sh.tangled.repo.issue.by"), 2017 subject: SubjectRef::Did(author.clone()), 2018 source: closed_uri.clone(), 2019 sort_micros: next_sort_micros(), 2020 }); 2021 h.mount( 2022 &author, 2023 &nsid("sh.tangled.repo.issue"), 2024 &rkey("open1"), 2025 issue_body(&repo, "still open"), 2026 ) 2027 .await; 2028 h.mount( 2029 &author, 2030 &nsid("sh.tangled.repo.issue"), 2031 &rkey("closed1"), 2032 issue_body(&repo, "shut"), 2033 ) 2034 .await; 2035 h.state.issue_states.upsert( 2036 at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"), 2037 closed_uri.clone(), 2038 1_777_593_800_000_000, 2039 IssueStateKind::Closed, 2040 ); 2041 2042 let app = router(h.state.clone()); 2043 let (status, body) = json_response( 2044 app.oneshot(list_request( 2045 "sh.tangled.repo.listIssuesBy", 2046 author_subject.as_ref(), 2047 &[("state", "closed")], 2048 )) 2049 .await 2050 .unwrap(), 2051 ) 2052 .await; 2053 assert_eq!(status, StatusCode::OK); 2054 let items = body["items"].as_array().expect("items array"); 2055 assert_eq!( 2056 items.len(), 2057 1, 2058 "only the closed issue survives state=closed" 2059 ); 2060 assert_eq!(items[0]["uri"], json!(closed_uri.as_ref())); 2061}