Monorepo for Tangled tangled.org
1

Configure Feed

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

core / bobbin / crates / xrpc / tests / extended.rs
27 kB 935 lines
1use std::sync::Arc; 2 3use axum::body::{Body, to_bytes}; 4use bobbin_edge_index::{CoverageWatch, EdgeStore, StateIndex}; 5use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig}; 6use bobbin_record_lru::{CacheCapacity, LruRecordStore}; 7use bobbin_resolver::RepoIdResolver; 8use bobbin_runtime::{RuntimeHasher, SystemClock}; 9use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader}; 10use bobbin_slingshot_client::SlingshotClient; 11use bobbin_types::edges::Edge; 12use bobbin_types::ids::SubjectRef; 13use bobbin_xrpc::{AppState, router}; 14use http::{Request, StatusCode}; 15use jacquard_common::DefaultStr; 16use jacquard_common::types::did::Did; 17use jacquard_common::types::nsid::Nsid; 18use jacquard_common::types::recordkey::Rkey; 19use jacquard_common::types::string::AtUri; 20use serde_json::{Value, json}; 21use tower::ServiceExt; 22use url::Url; 23use url::form_urlencoded::byte_serialize; 24use wiremock::matchers::{method, path, query_param}; 25use wiremock::{Mock, MockServer, ResponseTemplate}; 26 27const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; 28const TAG_BYTES: &str = "AAAAAAAAAAAAAAAAAAAAAAAAAAA="; 29const ARTIFACT_LINK: &str = "bafkreigh2akiscaildc7gnvtklbsfhdgwz72eolmpckbqr5ej26byp3uli"; 30 31fn at(s: &str) -> AtUri<DefaultStr> { 32 AtUri::new_owned(s).unwrap() 33} 34 35fn did(s: &str) -> Did<DefaultStr> { 36 Did::new_owned(s).unwrap() 37} 38 39fn rkey(s: &str) -> Rkey<DefaultStr> { 40 Rkey::new_owned(s).unwrap() 41} 42 43fn nsid(s: &'static str) -> Nsid<DefaultStr> { 44 Nsid::new_static(s).unwrap() 45} 46 47fn subj(s: &str) -> SubjectRef { 48 Did::<DefaultStr>::new_owned(s) 49 .map(SubjectRef::Did) 50 .unwrap_or_else(|_| SubjectRef::Uri(AtUri::new_owned(s).unwrap())) 51} 52 53struct Harness { 54 server: MockServer, 55 edges: Arc<EdgeStore>, 56 state: AppState, 57} 58 59impl Harness { 60 async fn new() -> Self { 61 let server = MockServer::start().await; 62 let edges = Arc::new(EdgeStore::new(RuntimeHasher::default())); 63 let coverage = Arc::new(CoverageWatch::new()); 64 let state = AppState::new( 65 Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), 66 SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(), 67 edges.clone(), 68 Arc::new(StateIndex::new(RuntimeHasher::default())), 69 Arc::new(StateIndex::new(RuntimeHasher::default())), 70 coverage, 71 Arc::new( 72 KnotProxy::new( 73 KnotProxyConfig::default(), 74 KnotHttpConfig::default(), 75 Arc::new(SystemClock::new()), 76 RuntimeHasher::default(), 77 ) 78 .unwrap(), 79 ), 80 Arc::new( 81 SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES, Arc::new(SystemClock::new())).unwrap(), 82 ) as Arc<dyn SearchReader>, 83 Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), 84 Arc::new(jacquard_identity::JacquardResolver::default()), 85 ); 86 Self { 87 server, 88 edges, 89 state, 90 } 91 } 92 93 fn add_edge( 94 &self, 95 kind: &Nsid<DefaultStr>, 96 subject: &AtUri<DefaultStr>, 97 source: &AtUri<DefaultStr>, 98 ) { 99 static EDGE_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1); 100 self.edges.add(Edge { 101 kind: kind.clone(), 102 subject: subj(subject.as_ref()), 103 source: source.clone(), 104 sort_micros: EDGE_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed), 105 }); 106 } 107 108 async fn mount( 109 &self, 110 did: &Did<DefaultStr>, 111 collection: &Nsid<DefaultStr>, 112 rkey: &Rkey<DefaultStr>, 113 value: Value, 114 ) { 115 let uri = format!( 116 "at://{}/{}/{}", 117 did.as_ref(), 118 collection.as_ref(), 119 rkey.as_ref() 120 ); 121 Mock::given(method("GET")) 122 .and(path("/xrpc/com.atproto.repo.getRecord")) 123 .and(query_param("repo", did.as_ref())) 124 .and(query_param("collection", collection.as_ref())) 125 .and(query_param("rkey", rkey.as_ref())) 126 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 127 "uri": uri, 128 "cid": CID, 129 "value": value, 130 }))) 131 .mount(&self.server) 132 .await; 133 } 134} 135 136fn list_request(endpoint: &str, subject: &str, extras: &[(&str, &str)]) -> Request<Body> { 137 let mut qs = format!("subject={subject}"); 138 extras.iter().for_each(|(k, v)| { 139 qs.push('&'); 140 qs.push_str(k); 141 qs.push('='); 142 qs.push_str(&encode(v)); 143 }); 144 Request::builder() 145 .uri(format!("/xrpc/{endpoint}?{qs}")) 146 .body(Body::empty()) 147 .unwrap() 148} 149 150fn list_request_escaped_subject( 151 endpoint: &str, 152 subject: &str, 153 extras: &[(&str, &str)], 154) -> Request<Body> { 155 let mut qs = format!("subject={}", encode(subject)); 156 extras.iter().for_each(|(k, v)| { 157 qs.push('&'); 158 qs.push_str(k); 159 qs.push('='); 160 qs.push_str(&encode(v)); 161 }); 162 Request::builder() 163 .uri(format!("/xrpc/{endpoint}?{qs}")) 164 .body(Body::empty()) 165 .unwrap() 166} 167 168fn encode(s: &str) -> String { 169 byte_serialize(s.as_bytes()).collect() 170} 171 172async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) { 173 let status = resp.status(); 174 let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap(); 175 let parsed: Value = serde_json::from_slice(&bytes).expect("JSON body"); 176 (status, parsed) 177} 178 179fn label_definition_body(name: &str) -> Value { 180 json!({ 181 "$type": "sh.tangled.label.definition", 182 "createdAt": "2026-05-01T00:00:00Z", 183 "name": name, 184 "scope": ["sh.tangled.repo.issue"], 185 "valueType": {"type": "boolean", "format": "any"} 186 }) 187} 188 189fn label_op_body(subject: &AtUri<DefaultStr>, def_uri: &AtUri<DefaultStr>, value: &str) -> Value { 190 json!({ 191 "$type": "sh.tangled.label.op", 192 "performedAt": "2026-05-01T00:00:00Z", 193 "subject": subject.as_ref(), 194 "add": [{"key": def_uri.as_ref(), "value": value}], 195 "delete": [] 196 }) 197} 198 199fn pipeline_body(repo_did: &Did<DefaultStr>) -> Value { 200 json!({ 201 "$type": "sh.tangled.pipeline", 202 "workflows": [], 203 "triggerMetadata": { 204 "kind": "manual", 205 "repo": { 206 "did": "did:plc:teq", 207 "repoDid": repo_did.as_ref(), 208 "knot": "nel.pet", 209 "defaultBranch": "main" 210 } 211 } 212 }) 213} 214 215fn pipeline_body_owner_only(owner_did: &Did<DefaultStr>) -> Value { 216 json!({ 217 "$type": "sh.tangled.pipeline", 218 "workflows": [], 219 "triggerMetadata": { 220 "kind": "manual", 221 "repo": { 222 "did": owner_did.as_ref(), 223 "knot": "nel.pet", 224 "defaultBranch": "main" 225 } 226 } 227 }) 228} 229 230fn pipeline_status_body(pipeline_uri: &AtUri<DefaultStr>) -> Value { 231 json!({ 232 "$type": "sh.tangled.pipeline.status", 233 "createdAt": "2026-05-01T00:00:00Z", 234 "pipeline": pipeline_uri.as_ref(), 235 "workflow": pipeline_uri.as_ref(), 236 "status": "success" 237 }) 238} 239 240fn artifact_body(repo_did: &Did<DefaultStr>, name: &str) -> Value { 241 json!({ 242 "$type": "sh.tangled.repo.artifact", 243 "createdAt": "2026-05-01T00:00:00Z", 244 "name": name, 245 "repoDid": repo_did.as_ref(), 246 "tag": {"$bytes": TAG_BYTES}, 247 "artifact": { 248 "$type": "blob", 249 "ref": {"$link": ARTIFACT_LINK}, 250 "mimeType": "application/octet-stream", 251 "size": 12 252 } 253 }) 254} 255 256fn knot_member_body(subject_did: &Did<DefaultStr>) -> Value { 257 json!({ 258 "$type": "sh.tangled.knot.member", 259 "createdAt": "2026-05-01T00:00:00Z", 260 "subject": subject_did.as_ref(), 261 "domain": "oyster.cafe" 262 }) 263} 264 265fn spindle_member_body(subject_did: &Did<DefaultStr>) -> Value { 266 json!({ 267 "$type": "sh.tangled.spindle.member", 268 "createdAt": "2026-05-01T00:00:00Z", 269 "subject": subject_did.as_ref(), 270 "instance": "spin.nel.pet" 271 }) 272} 273 274fn string_body(filename: &str, contents: &str) -> Value { 275 json!({ 276 "$type": "sh.tangled.string", 277 "createdAt": "2026-05-01T00:00:00Z", 278 "filename": filename, 279 "description": "test fixture", 280 "contents": contents 281 }) 282} 283 284#[tokio::test] 285async fn list_label_definitions_keys_on_owner_did() { 286 let h = Harness::new().await; 287 let owner = did("did:plc:abalone"); 288 let rk = rkey("bug"); 289 let source = at(&format!( 290 "at://{}/sh.tangled.label.definition/{}", 291 owner.as_ref(), 292 rk.as_ref() 293 )); 294 h.add_edge( 295 &nsid("sh.tangled.label.definition"), 296 &at(&format!("at://{}", owner.as_ref())), 297 &source, 298 ); 299 h.mount( 300 &owner, 301 &nsid("sh.tangled.label.definition"), 302 &rk, 303 label_definition_body("bug"), 304 ) 305 .await; 306 307 let app = router(h.state.clone()); 308 let (status, body) = json_response( 309 app.oneshot(list_request( 310 "sh.tangled.label.listDefinitions", 311 &format!("at://{}", owner.as_ref()), 312 &[], 313 )) 314 .await 315 .unwrap(), 316 ) 317 .await; 318 assert_eq!(status, StatusCode::OK); 319 let items = body["items"].as_array().unwrap(); 320 assert_eq!(items.len(), 1); 321 assert_eq!(items[0]["value"]["name"], json!("bug")); 322 assert_eq!( 323 items[0]["value"]["scope"][0], 324 json!("sh.tangled.repo.issue") 325 ); 326} 327 328#[tokio::test] 329async fn count_label_definitions_dedupes_per_author() { 330 let h = Harness::new().await; 331 let owner = did("did:plc:abalone"); 332 let subject = at(&format!("at://{}", owner.as_ref())); 333 h.add_edge( 334 &nsid("sh.tangled.label.definition"), 335 &subject, 336 &at(&format!( 337 "at://{}/sh.tangled.label.definition/bug", 338 owner.as_ref() 339 )), 340 ); 341 h.add_edge( 342 &nsid("sh.tangled.label.definition"), 343 &subject, 344 &at(&format!( 345 "at://{}/sh.tangled.label.definition/wontfix", 346 owner.as_ref() 347 )), 348 ); 349 350 let app = router(h.state.clone()); 351 let (_, body) = json_response( 352 app.oneshot(list_request( 353 "sh.tangled.label.countDefinitions", 354 subject.as_ref(), 355 &[], 356 )) 357 .await 358 .unwrap(), 359 ) 360 .await; 361 assert_eq!(body["count"], json!(2)); 362 assert_eq!(body["distinctAuthors"], json!(1)); 363} 364 365#[tokio::test] 366async fn list_label_ops_accepts_issue_subject() { 367 let h = Harness::new().await; 368 let issue_uri = at("at://did:plc:abalone/sh.tangled.repo.issue/i1"); 369 let author = did("did:plc:nel"); 370 let rk = rkey("op1"); 371 let def_uri = at("at://did:plc:abalone/sh.tangled.label.definition/bug"); 372 h.add_edge( 373 &nsid("sh.tangled.label.op"), 374 &issue_uri, 375 &at(&format!( 376 "at://{}/sh.tangled.label.op/{}", 377 author.as_ref(), 378 rk.as_ref() 379 )), 380 ); 381 h.mount( 382 &author, 383 &nsid("sh.tangled.label.op"), 384 &rk, 385 label_op_body(&issue_uri, &def_uri, "true"), 386 ) 387 .await; 388 389 let app = router(h.state.clone()); 390 let (status, body) = json_response( 391 app.oneshot(list_request( 392 "sh.tangled.label.listOps", 393 issue_uri.as_ref(), 394 &[], 395 )) 396 .await 397 .unwrap(), 398 ) 399 .await; 400 assert_eq!(status, StatusCode::OK); 401 let items = body["items"].as_array().unwrap(); 402 assert_eq!(items.len(), 1); 403 assert_eq!(items[0]["value"]["subject"], json!(issue_uri.as_ref())); 404 assert_eq!(items[0]["value"]["add"][0]["key"], json!(def_uri.as_ref())); 405} 406 407#[tokio::test] 408async fn list_ops_accepts_percent_escaped_subject() { 409 let h = Harness::new().await; 410 let issue_uri = at("at://did:plc:clam/sh.tangled.repo.issue/i1"); 411 let author = did("did:plc:nel"); 412 let rk = rkey("op1"); 413 let def_uri = at("at://did:plc:clam/sh.tangled.label.definition/bug"); 414 h.add_edge( 415 &nsid("sh.tangled.label.op"), 416 &issue_uri, 417 &at(&format!( 418 "at://{}/sh.tangled.label.op/{}", 419 author.as_ref(), 420 rk.as_ref() 421 )), 422 ); 423 h.mount( 424 &author, 425 &nsid("sh.tangled.label.op"), 426 &rk, 427 label_op_body(&issue_uri, &def_uri, "true"), 428 ) 429 .await; 430 431 let app = router(h.state.clone()); 432 let (status, body) = json_response( 433 app.oneshot(list_request_escaped_subject( 434 "sh.tangled.label.listOps", 435 issue_uri.as_ref(), 436 &[], 437 )) 438 .await 439 .unwrap(), 440 ) 441 .await; 442 assert_eq!( 443 status, 444 StatusCode::OK, 445 "escaped subject must still be accepted" 446 ); 447 assert_eq!(body["items"].as_array().unwrap().len(), 1); 448} 449 450#[tokio::test] 451async fn list_label_ops_pull_subject_round_trip() { 452 let h = Harness::new().await; 453 let pull_uri = at("at://did:plc:abalone/sh.tangled.repo.pull/p1"); 454 let author = did("did:plc:bailey"); 455 let rk = rkey("op1"); 456 let def_uri = at("at://did:plc:abalone/sh.tangled.label.definition/wontfix"); 457 let source = at(&format!( 458 "at://{}/sh.tangled.label.op/{}", 459 author.as_ref(), 460 rk.as_ref() 461 )); 462 let body = label_op_body(&pull_uri, &def_uri, "true"); 463 let parsed = 464 bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.label.op"), body.clone()) 465 .expect("parse label.op record"); 466 parsed 467 .extract_edges(&source) 468 .expect("extract") 469 .into_iter() 470 .for_each(|e| h.edges.add(e)); 471 h.mount(&author, &nsid("sh.tangled.label.op"), &rk, body) 472 .await; 473 474 let app = router(h.state.clone()); 475 let (status, json) = json_response( 476 app.oneshot(list_request( 477 "sh.tangled.label.listOps", 478 pull_uri.as_ref(), 479 &[], 480 )) 481 .await 482 .unwrap(), 483 ) 484 .await; 485 assert_eq!(status, StatusCode::OK); 486 let items = json["items"].as_array().unwrap(); 487 assert_eq!(items.len(), 1); 488 assert_eq!(items[0]["value"]["subject"], json!(pull_uri.as_ref())); 489 assert_eq!(items[0]["value"]["add"][0]["key"], json!(def_uri.as_ref())); 490} 491 492#[tokio::test] 493async fn list_label_ops_rejects_bare_did_subject() { 494 let h = Harness::new().await; 495 let app = router(h.state.clone()); 496 let (status, body) = json_response( 497 app.oneshot(list_request( 498 "sh.tangled.label.listOps", 499 "at://did:plc:abalone", 500 &[], 501 )) 502 .await 503 .unwrap(), 504 ) 505 .await; 506 assert_eq!(status, StatusCode::BAD_REQUEST); 507 let msg = body["message"].as_str().unwrap_or_default(); 508 assert!( 509 msg.contains("sh.tangled.repo.issue") && msg.contains("sh.tangled.repo.pull"), 510 "message must list allowed collections, got {msg}" 511 ); 512} 513 514#[tokio::test] 515async fn list_label_ops_rejects_unrelated_collection() { 516 let h = Harness::new().await; 517 let app = router(h.state.clone()); 518 let (status, _) = json_response( 519 app.oneshot(list_request( 520 "sh.tangled.label.listOps", 521 "at://did:plc:abalone/sh.tangled.repo/r1", 522 &[], 523 )) 524 .await 525 .unwrap(), 526 ) 527 .await; 528 assert_eq!(status, StatusCode::BAD_REQUEST); 529} 530 531#[tokio::test] 532async fn list_pipelines_keys_on_repo_did() { 533 let h = Harness::new().await; 534 let repo_did = did("did:plc:abalone"); 535 let subject = at(&format!("at://{}", repo_did.as_ref())); 536 let spindle_did = did("did:plc:lyna"); 537 let rk = rkey("pl1"); 538 h.add_edge( 539 &nsid("sh.tangled.pipeline"), 540 &subject, 541 &at(&format!( 542 "at://{}/sh.tangled.pipeline/{}", 543 spindle_did.as_ref(), 544 rk.as_ref() 545 )), 546 ); 547 h.mount( 548 &spindle_did, 549 &nsid("sh.tangled.pipeline"), 550 &rk, 551 pipeline_body(&repo_did), 552 ) 553 .await; 554 555 let app = router(h.state.clone()); 556 let (status, body) = json_response( 557 app.oneshot(list_request( 558 "sh.tangled.pipeline.listPipelines", 559 subject.as_ref(), 560 &[], 561 )) 562 .await 563 .unwrap(), 564 ) 565 .await; 566 assert_eq!(status, StatusCode::OK); 567 let items = body["items"].as_array().unwrap(); 568 assert_eq!(items.len(), 1); 569 assert_eq!( 570 items[0]["value"]["triggerMetadata"]["repo"]["repoDid"], 571 json!(repo_did.as_ref()) 572 ); 573} 574 575#[tokio::test] 576async fn count_pipelines_returns_zero_when_no_edges() { 577 let h = Harness::new().await; 578 let app = router(h.state.clone()); 579 let (_, body) = json_response( 580 app.oneshot(list_request( 581 "sh.tangled.pipeline.countPipelines", 582 "at://did:plc:abalone", 583 &[], 584 )) 585 .await 586 .unwrap(), 587 ) 588 .await; 589 assert_eq!(body["count"], json!(0)); 590} 591 592#[tokio::test] 593async fn list_pipeline_statuses_keys_on_pipeline_uri() { 594 let h = Harness::new().await; 595 let pipeline_uri = at("at://did:plc:lyna/sh.tangled.pipeline/pl1"); 596 let author = did("did:plc:bailey"); 597 let rk = rkey("s1"); 598 h.add_edge( 599 &nsid("sh.tangled.pipeline.status"), 600 &pipeline_uri, 601 &at(&format!( 602 "at://{}/sh.tangled.pipeline.status/{}", 603 author.as_ref(), 604 rk.as_ref() 605 )), 606 ); 607 h.mount( 608 &author, 609 &nsid("sh.tangled.pipeline.status"), 610 &rk, 611 pipeline_status_body(&pipeline_uri), 612 ) 613 .await; 614 615 let app = router(h.state.clone()); 616 let (status, body) = json_response( 617 app.oneshot(list_request( 618 "sh.tangled.pipeline.listStatuses", 619 pipeline_uri.as_ref(), 620 &[], 621 )) 622 .await 623 .unwrap(), 624 ) 625 .await; 626 assert_eq!(status, StatusCode::OK); 627 let items = body["items"].as_array().unwrap(); 628 assert_eq!(items.len(), 1); 629 assert_eq!(items[0]["value"]["pipeline"], json!(pipeline_uri.as_ref())); 630 assert_eq!(items[0]["value"]["status"], json!("success")); 631} 632 633#[tokio::test] 634async fn pipeline_status_endpoint_rejects_bare_did_subject() { 635 let h = Harness::new().await; 636 let app = router(h.state.clone()); 637 let (status, body) = json_response( 638 app.oneshot(list_request( 639 "sh.tangled.pipeline.listStatuses", 640 "at://did:plc:lyna", 641 &[], 642 )) 643 .await 644 .unwrap(), 645 ) 646 .await; 647 assert_eq!(status, StatusCode::BAD_REQUEST); 648 assert!( 649 body["message"] 650 .as_str() 651 .unwrap_or_default() 652 .contains("sh.tangled.pipeline/<rkey>"), 653 "{body}" 654 ); 655} 656 657#[tokio::test] 658async fn list_artifacts_keys_on_repo_did() { 659 let h = Harness::new().await; 660 let repo_did = did("did:plc:abalone"); 661 let subject = at(&format!("at://{}", repo_did.as_ref())); 662 let owner = did("did:plc:nel"); 663 let rk = rkey("a1"); 664 h.add_edge( 665 &nsid("sh.tangled.repo.artifact"), 666 &subject, 667 &at(&format!( 668 "at://{}/sh.tangled.repo.artifact/{}", 669 owner.as_ref(), 670 rk.as_ref() 671 )), 672 ); 673 h.mount( 674 &owner, 675 &nsid("sh.tangled.repo.artifact"), 676 &rk, 677 artifact_body(&repo_did, "out.bin"), 678 ) 679 .await; 680 681 let app = router(h.state.clone()); 682 let (status, body) = json_response( 683 app.oneshot(list_request( 684 "sh.tangled.repo.listArtifacts", 685 subject.as_ref(), 686 &[], 687 )) 688 .await 689 .unwrap(), 690 ) 691 .await; 692 assert_eq!(status, StatusCode::OK); 693 let items = body["items"].as_array().unwrap(); 694 assert_eq!(items.len(), 1); 695 assert_eq!(items[0]["value"]["name"], json!("out.bin")); 696 assert_eq!(items[0]["value"]["repoDid"], json!(repo_did.as_ref())); 697} 698 699#[tokio::test] 700async fn list_knot_members_keys_on_subject_did() { 701 let h = Harness::new().await; 702 let subject_did = did("did:plc:nel"); 703 let subject = at(&format!("at://{}", subject_did.as_ref())); 704 let admin = did("did:plc:teq"); 705 let rk = rkey("m1"); 706 h.add_edge( 707 &nsid("sh.tangled.knot.member"), 708 &subject, 709 &at(&format!( 710 "at://{}/sh.tangled.knot.member/{}", 711 admin.as_ref(), 712 rk.as_ref() 713 )), 714 ); 715 h.mount( 716 &admin, 717 &nsid("sh.tangled.knot.member"), 718 &rk, 719 knot_member_body(&subject_did), 720 ) 721 .await; 722 723 let app = router(h.state.clone()); 724 let (status, body) = json_response( 725 app.oneshot(list_request( 726 "sh.tangled.knot.listMembers", 727 subject.as_ref(), 728 &[], 729 )) 730 .await 731 .unwrap(), 732 ) 733 .await; 734 assert_eq!(status, StatusCode::OK); 735 let items = body["items"].as_array().unwrap(); 736 assert_eq!(items.len(), 1); 737 assert_eq!(items[0]["value"]["subject"], json!(subject_did.as_ref())); 738 assert_eq!(items[0]["value"]["domain"], json!("oyster.cafe")); 739} 740 741#[tokio::test] 742async fn list_spindle_members_keys_on_subject_did() { 743 let h = Harness::new().await; 744 let subject_did = did("did:plc:olaren"); 745 let subject = at(&format!("at://{}", subject_did.as_ref())); 746 let admin = did("did:plc:teq"); 747 let rk = rkey("m1"); 748 h.add_edge( 749 &nsid("sh.tangled.spindle.member"), 750 &subject, 751 &at(&format!( 752 "at://{}/sh.tangled.spindle.member/{}", 753 admin.as_ref(), 754 rk.as_ref() 755 )), 756 ); 757 h.mount( 758 &admin, 759 &nsid("sh.tangled.spindle.member"), 760 &rk, 761 spindle_member_body(&subject_did), 762 ) 763 .await; 764 765 let app = router(h.state.clone()); 766 let (status, body) = json_response( 767 app.oneshot(list_request( 768 "sh.tangled.spindle.listMembers", 769 subject.as_ref(), 770 &[], 771 )) 772 .await 773 .unwrap(), 774 ) 775 .await; 776 assert_eq!(status, StatusCode::OK); 777 let items = body["items"].as_array().unwrap(); 778 assert_eq!(items.len(), 1); 779 assert_eq!(items[0]["value"]["subject"], json!(subject_did.as_ref())); 780 assert_eq!(items[0]["value"]["instance"], json!("spin.nel.pet")); 781} 782 783#[tokio::test] 784async fn list_strings_keys_on_owner_did() { 785 let h = Harness::new().await; 786 let owner = did("did:plc:abalone"); 787 let subject = at(&format!("at://{}", owner.as_ref())); 788 let rk = rkey("k1"); 789 h.add_edge( 790 &nsid("sh.tangled.string"), 791 &subject, 792 &at(&format!( 793 "at://{}/sh.tangled.string/{}", 794 owner.as_ref(), 795 rk.as_ref() 796 )), 797 ); 798 h.mount( 799 &owner, 800 &nsid("sh.tangled.string"), 801 &rk, 802 string_body("snippet.rs", "fn main() {}"), 803 ) 804 .await; 805 806 let app = router(h.state.clone()); 807 let (status, body) = json_response( 808 app.oneshot(list_request( 809 "sh.tangled.string.listStrings", 810 subject.as_ref(), 811 &[], 812 )) 813 .await 814 .unwrap(), 815 ) 816 .await; 817 assert_eq!(status, StatusCode::OK); 818 let items = body["items"].as_array().unwrap(); 819 assert_eq!(items.len(), 1); 820 assert_eq!(items[0]["value"]["filename"], json!("snippet.rs")); 821 assert_eq!(items[0]["value"]["contents"], json!("fn main() {}")); 822} 823 824#[tokio::test] 825async fn count_strings_dedupes_per_owner() { 826 let h = Harness::new().await; 827 let owner = did("did:plc:abalone"); 828 let subject = at(&format!("at://{}", owner.as_ref())); 829 ["k1", "k2", "k3"].iter().for_each(|r| { 830 h.add_edge( 831 &nsid("sh.tangled.string"), 832 &subject, 833 &at(&format!("at://{}/sh.tangled.string/{}", owner.as_ref(), r)), 834 ); 835 }); 836 837 let app = router(h.state.clone()); 838 let (_, body) = json_response( 839 app.oneshot(list_request( 840 "sh.tangled.string.countStrings", 841 subject.as_ref(), 842 &[], 843 )) 844 .await 845 .unwrap(), 846 ) 847 .await; 848 assert_eq!(body["count"], json!(3)); 849 assert_eq!(body["distinctAuthors"], json!(1)); 850} 851 852#[tokio::test] 853async fn extractor_to_xrpc_round_trip_for_pipeline() { 854 let h = Harness::new().await; 855 let repo_did = did("did:plc:abalone"); 856 let spindle_did = did("did:plc:lyna"); 857 let rk = rkey("pl1"); 858 let source = at(&format!( 859 "at://{}/sh.tangled.pipeline/{}", 860 spindle_did.as_ref(), 861 rk.as_ref() 862 )); 863 let body = pipeline_body(&repo_did); 864 let parsed = 865 bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.pipeline"), body.clone()) 866 .expect("parse pipeline record"); 867 parsed 868 .extract_edges(&source) 869 .expect("extract") 870 .into_iter() 871 .for_each(|e| h.edges.add(e)); 872 h.mount(&spindle_did, &nsid("sh.tangled.pipeline"), &rk, body) 873 .await; 874 875 let app = router(h.state.clone()); 876 let (status, json) = json_response( 877 app.oneshot(list_request( 878 "sh.tangled.pipeline.listPipelines", 879 &format!("at://{}", repo_did.as_ref()), 880 &[], 881 )) 882 .await 883 .unwrap(), 884 ) 885 .await; 886 assert_eq!( 887 status, 888 StatusCode::OK, 889 "extractor key must match handler subject, body was {json}" 890 ); 891 let items = json["items"].as_array().unwrap(); 892 assert_eq!(items.len(), 1, "expected exactly one pipeline edge"); 893} 894 895#[tokio::test] 896async fn list_pipelines_drops_records_without_resolvable_repo_did() { 897 let h = Harness::new().await; 898 let owner_did = did("did:plc:nel"); 899 let spindle_did = did("did:plc:lyna"); 900 let rk = rkey("pl1"); 901 let source = at(&format!( 902 "at://{}/sh.tangled.pipeline/{}", 903 spindle_did.as_ref(), 904 rk.as_ref() 905 )); 906 let body = pipeline_body_owner_only(&owner_did); 907 let parsed = 908 bobbin_types::edges::Record::from_json_value(&nsid("sh.tangled.pipeline"), body.clone()) 909 .expect("parse pipeline record"); 910 parsed 911 .extract_edges(&source) 912 .expect("extract") 913 .into_iter() 914 .for_each(|e| h.edges.add(e)); 915 h.mount(&spindle_did, &nsid("sh.tangled.pipeline"), &rk, body) 916 .await; 917 918 let app = router(h.state.clone()); 919 let (status, json) = json_response( 920 app.oneshot(list_request( 921 "sh.tangled.pipeline.listPipelines", 922 &format!("at://{}", owner_did.as_ref()), 923 &[], 924 )) 925 .await 926 .unwrap(), 927 ) 928 .await; 929 assert_eq!(status, StatusCode::OK, "body was {json}"); 930 let items = json["items"].as_array().unwrap(); 931 assert!( 932 items.is_empty(), 933 "pipeline without resolvable repoDid must be dropped, got {json}" 934 ); 935}