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 / bulk.rs
19 kB 684 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_xrpc::{AppState, router}; 12use http::{Request, StatusCode}; 13use jacquard_common::DefaultStr; 14use jacquard_common::types::did::Did; 15use jacquard_common::types::handle::Handle; 16use jacquard_common::types::nsid::Nsid; 17use jacquard_common::types::recordkey::Rkey; 18use serde_json::{Value, json}; 19use tower::ServiceExt; 20use url::Url; 21use url::form_urlencoded::byte_serialize; 22use wiremock::matchers::{method, path, query_param}; 23use wiremock::{Mock, MockServer, ResponseTemplate}; 24 25const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; 26 27fn did(s: &str) -> Did<DefaultStr> { 28 Did::new_owned(s).unwrap() 29} 30 31fn rkey(s: &str) -> Rkey<DefaultStr> { 32 Rkey::new_owned(s).unwrap() 33} 34 35fn nsid(s: &'static str) -> Nsid<DefaultStr> { 36 Nsid::new_static(s).unwrap() 37} 38 39fn handle(s: &str) -> Handle<DefaultStr> { 40 Handle::new_owned(s).unwrap() 41} 42 43struct Harness { 44 server: MockServer, 45 state: AppState, 46} 47 48impl Harness { 49 async fn new() -> Self { 50 let server = MockServer::start().await; 51 let coverage = Arc::new(CoverageWatch::new()); 52 let state = AppState::new( 53 Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), 54 SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(), 55 Arc::new(EdgeStore::new(RuntimeHasher::default())), 56 Arc::new(StateIndex::new(RuntimeHasher::default())), 57 Arc::new(StateIndex::new(RuntimeHasher::default())), 58 coverage, 59 Arc::new( 60 KnotProxy::new( 61 KnotProxyConfig::default(), 62 KnotHttpConfig::default(), 63 Arc::new(SystemClock::new()), 64 RuntimeHasher::default(), 65 ) 66 .unwrap(), 67 ), 68 Arc::new( 69 SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES, Arc::new(SystemClock::new())).unwrap(), 70 ) as Arc<dyn SearchReader>, 71 Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), 72 Arc::new(bobbin_xrpc::default_directory()), 73 ); 74 Self { server, state } 75 } 76 77 async fn mount( 78 &self, 79 did: &Did<DefaultStr>, 80 collection: &Nsid<DefaultStr>, 81 rkey: &Rkey<DefaultStr>, 82 value: Value, 83 ) { 84 let uri = format!( 85 "at://{}/{}/{}", 86 did.as_ref(), 87 collection.as_ref(), 88 rkey.as_ref() 89 ); 90 Mock::given(method("GET")) 91 .and(path("/xrpc/com.atproto.repo.getRecord")) 92 .and(query_param("repo", did.as_ref())) 93 .and(query_param("collection", collection.as_ref())) 94 .and(query_param("rkey", rkey.as_ref())) 95 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 96 "uri": uri, 97 "cid": CID, 98 "value": value, 99 }))) 100 .mount(&self.server) 101 .await; 102 } 103 104 async fn mount_404( 105 &self, 106 did: &Did<DefaultStr>, 107 collection: &Nsid<DefaultStr>, 108 rkey: &Rkey<DefaultStr>, 109 ) { 110 Mock::given(method("GET")) 111 .and(path("/xrpc/com.atproto.repo.getRecord")) 112 .and(query_param("repo", did.as_ref())) 113 .and(query_param("collection", collection.as_ref())) 114 .and(query_param("rkey", rkey.as_ref())) 115 .respond_with( 116 ResponseTemplate::new(404) 117 .set_body_json(json!({"error": "RecordNotFound", "message": "missing"})), 118 ) 119 .mount(&self.server) 120 .await; 121 } 122} 123 124fn enc(s: &str) -> String { 125 byte_serialize(s.as_bytes()).collect() 126} 127 128fn bulk_request(endpoint: &str, key: &str, values: &[&str]) -> Request<Body> { 129 let qs = values 130 .iter() 131 .map(|v| format!("{key}={v}")) 132 .collect::<Vec<_>>() 133 .join("&"); 134 Request::builder() 135 .uri(format!("/xrpc/{endpoint}?{qs}")) 136 .body(Body::empty()) 137 .unwrap() 138} 139 140fn bulk_request_escaped(endpoint: &str, key: &str, values: &[&str]) -> Request<Body> { 141 let qs = values 142 .iter() 143 .map(|v| format!("{key}={}", enc(v))) 144 .collect::<Vec<_>>() 145 .join("&"); 146 Request::builder() 147 .uri(format!("/xrpc/{endpoint}?{qs}")) 148 .body(Body::empty()) 149 .unwrap() 150} 151 152async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) { 153 let status = resp.status(); 154 let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap(); 155 let parsed: Value = serde_json::from_slice(&bytes).expect("JSON body"); 156 (status, parsed) 157} 158 159fn issue_body(repo_did: &Did<DefaultStr>, title: &str) -> Value { 160 json!({ 161 "$type": "sh.tangled.repo.issue", 162 "repo": repo_did.as_ref(), 163 "title": title, 164 "createdAt": "2026-05-01T00:00:00Z" 165 }) 166} 167 168fn pull_body(target_repo: &Did<DefaultStr>, title: &str) -> Value { 169 json!({ 170 "$type": "sh.tangled.repo.pull", 171 "title": title, 172 "createdAt": "2026-05-01T00:00:00Z", 173 "rounds": [], 174 "target": { 175 "branch": "main", 176 "repo": target_repo.as_ref() 177 } 178 }) 179} 180 181fn repo_body(name: &str) -> Value { 182 json!({ 183 "$type": "sh.tangled.repo", 184 "name": name, 185 "knot": "oyster.cafe", 186 "createdAt": "2026-05-01T00:00:00Z" 187 }) 188} 189 190fn profile_body(handle: &Handle<DefaultStr>) -> Value { 191 json!({ 192 "$type": "sh.tangled.actor.profile", 193 "bluesky": false, 194 "preferredHandle": handle.as_ref() 195 }) 196} 197 198#[tokio::test] 199async fn get_repos_returns_all_resolved_records() { 200 let h = Harness::new().await; 201 h.mount( 202 &did("did:plc:nel"), 203 &nsid("sh.tangled.repo"), 204 &rkey("abalone"), 205 repo_body("abalone"), 206 ) 207 .await; 208 h.mount( 209 &did("did:plc:teq"), 210 &nsid("sh.tangled.repo"), 211 &rkey("limpet"), 212 repo_body("limpet"), 213 ) 214 .await; 215 let app = router(h.state.clone()); 216 let (status, body) = json_response( 217 app.oneshot(bulk_request( 218 "sh.tangled.repo.getRepos", 219 "repos", 220 &[ 221 "at://did:plc:nel/sh.tangled.repo/abalone", 222 "at://did:plc:teq/sh.tangled.repo/limpet", 223 ], 224 )) 225 .await 226 .unwrap(), 227 ) 228 .await; 229 assert_eq!(status, StatusCode::OK); 230 let items = body["items"].as_array().unwrap(); 231 assert_eq!(items.len(), 2); 232 let names: Vec<&str> = items 233 .iter() 234 .map(|v| v["value"]["name"].as_str().unwrap()) 235 .collect(); 236 assert!(names.contains(&"abalone")); 237 assert!(names.contains(&"limpet")); 238} 239 240#[tokio::test] 241async fn get_profiles_returns_all_resolved_profiles() { 242 let h = Harness::new().await; 243 h.mount( 244 &did("did:plc:nel"), 245 &nsid("sh.tangled.actor.profile"), 246 &rkey("self"), 247 profile_body(&handle("witchcraft.systems")), 248 ) 249 .await; 250 h.mount( 251 &did("did:plc:teq"), 252 &nsid("sh.tangled.actor.profile"), 253 &rkey("self"), 254 profile_body(&handle("olaren.dev")), 255 ) 256 .await; 257 let app = router(h.state.clone()); 258 let (status, body) = json_response( 259 app.oneshot(bulk_request( 260 "sh.tangled.actor.getProfiles", 261 "actors", 262 &[ 263 "at://did:plc:nel/sh.tangled.actor.profile/self", 264 "at://did:plc:teq/sh.tangled.actor.profile/self", 265 ], 266 )) 267 .await 268 .unwrap(), 269 ) 270 .await; 271 assert_eq!(status, StatusCode::OK); 272 let items = body["items"].as_array().unwrap(); 273 assert_eq!(items.len(), 2); 274} 275 276#[tokio::test] 277async fn get_profiles_accepts_percent_escaped_at_uris() { 278 let h = Harness::new().await; 279 h.mount( 280 &did("did:plc:nel"), 281 &nsid("sh.tangled.actor.profile"), 282 &rkey("self"), 283 profile_body(&handle("witchcraft.systems")), 284 ) 285 .await; 286 let app = router(h.state.clone()); 287 let (status, body) = json_response( 288 app.oneshot(bulk_request_escaped( 289 "sh.tangled.actor.getProfiles", 290 "actors", 291 &["at://did:plc:nel/sh.tangled.actor.profile/self"], 292 )) 293 .await 294 .unwrap(), 295 ) 296 .await; 297 assert_eq!(status, StatusCode::OK); 298 assert_eq!(body["items"].as_array().unwrap().len(), 1); 299} 300 301#[tokio::test] 302async fn get_issues_returns_all_resolved_issues() { 303 let h = Harness::new().await; 304 let repo = did("did:plc:abalone"); 305 h.mount( 306 &did("did:plc:nel"), 307 &nsid("sh.tangled.repo.issue"), 308 &rkey("i1"), 309 issue_body(&repo, "first"), 310 ) 311 .await; 312 h.mount( 313 &did("did:plc:olaren"), 314 &nsid("sh.tangled.repo.issue"), 315 &rkey("i2"), 316 issue_body(&repo, "second"), 317 ) 318 .await; 319 let app = router(h.state.clone()); 320 let (status, body) = json_response( 321 app.oneshot(bulk_request( 322 "sh.tangled.repo.getIssues", 323 "issues", 324 &[ 325 "at://did:plc:nel/sh.tangled.repo.issue/i1", 326 "at://did:plc:olaren/sh.tangled.repo.issue/i2", 327 ], 328 )) 329 .await 330 .unwrap(), 331 ) 332 .await; 333 assert_eq!(status, StatusCode::OK); 334 let items = body["items"].as_array().unwrap(); 335 assert_eq!(items.len(), 2); 336} 337 338#[tokio::test] 339async fn get_pulls_returns_all_resolved_pulls() { 340 let h = Harness::new().await; 341 let target = did("did:plc:abalone"); 342 h.mount( 343 &did("did:plc:nel"), 344 &nsid("sh.tangled.repo.pull"), 345 &rkey("p1"), 346 pull_body(&target, "patch one"), 347 ) 348 .await; 349 let app = router(h.state.clone()); 350 let (status, body) = json_response( 351 app.oneshot(bulk_request( 352 "sh.tangled.repo.getPulls", 353 "pulls", 354 &["at://did:plc:nel/sh.tangled.repo.pull/p1"], 355 )) 356 .await 357 .unwrap(), 358 ) 359 .await; 360 assert_eq!(status, StatusCode::OK); 361 let items = body["items"].as_array().unwrap(); 362 assert_eq!(items.len(), 1); 363 assert_eq!(items[0]["value"]["title"], json!("patch one")); 364} 365 366#[tokio::test] 367async fn missing_records_are_dropped_silently() { 368 let h = Harness::new().await; 369 h.mount( 370 &did("did:plc:nel"), 371 &nsid("sh.tangled.repo"), 372 &rkey("abalone"), 373 repo_body("abalone"), 374 ) 375 .await; 376 h.mount_404( 377 &did("did:plc:teq"), 378 &nsid("sh.tangled.repo"), 379 &rkey("ghost"), 380 ) 381 .await; 382 let app = router(h.state.clone()); 383 let (status, body) = json_response( 384 app.oneshot(bulk_request( 385 "sh.tangled.repo.getRepos", 386 "repos", 387 &[ 388 "at://did:plc:nel/sh.tangled.repo/abalone", 389 "at://did:plc:teq/sh.tangled.repo/ghost", 390 ], 391 )) 392 .await 393 .unwrap(), 394 ) 395 .await; 396 assert_eq!(status, StatusCode::OK); 397 let items = body["items"].as_array().unwrap(); 398 assert_eq!( 399 items.len(), 400 1, 401 "missing records must be dropped not fail the bulk call" 402 ); 403 assert_eq!(items[0]["value"]["name"], json!("abalone")); 404} 405 406#[tokio::test] 407async fn transient_failure_drops_only_that_record() { 408 let h = Harness::new().await; 409 h.mount( 410 &did("did:plc:nel"), 411 &nsid("sh.tangled.repo"), 412 &rkey("conch"), 413 repo_body("conch"), 414 ) 415 .await; 416 Mock::given(method("GET")) 417 .and(path("/xrpc/com.atproto.repo.getRecord")) 418 .and(query_param("repo", "did:plc:teq")) 419 .and(query_param("collection", "sh.tangled.repo")) 420 .and(query_param("rkey", "flaky")) 421 .respond_with(ResponseTemplate::new(503)) 422 .mount(&h.server) 423 .await; 424 let app = router(h.state.clone()); 425 let (status, body) = json_response( 426 app.oneshot(bulk_request( 427 "sh.tangled.repo.getRepos", 428 "repos", 429 &[ 430 "at://did:plc:nel/sh.tangled.repo/conch", 431 "at://did:plc:teq/sh.tangled.repo/flaky", 432 ], 433 )) 434 .await 435 .unwrap(), 436 ) 437 .await; 438 assert_eq!(status, StatusCode::OK); 439 let items = body["items"].as_array().unwrap(); 440 assert_eq!( 441 items.len(), 442 1, 443 "transient upstream failure drops that record, it must not fail the bulk call" 444 ); 445 assert_eq!(items[0]["value"]["name"], json!("conch")); 446} 447 448#[tokio::test] 449async fn wrong_collection_uri_fails_bulk_request() { 450 let h = Harness::new().await; 451 h.mount( 452 &did("did:plc:nel"), 453 &nsid("sh.tangled.repo"), 454 &rkey("conch"), 455 repo_body("conch"), 456 ) 457 .await; 458 let app = router(h.state.clone()); 459 let (status, body) = json_response( 460 app.oneshot(bulk_request( 461 "sh.tangled.repo.getRepos", 462 "repos", 463 &[ 464 "at://did:plc:nel/sh.tangled.repo/conch", 465 "at://did:plc:teq/sh.tangled.repo.issue/whelk", 466 ], 467 )) 468 .await 469 .unwrap(), 470 ) 471 .await; 472 assert_eq!( 473 status, 474 StatusCode::BAD_REQUEST, 475 "a wrong-collection uri must fail the whole bulk request" 476 ); 477 assert_eq!(body["error"], "InvalidRequest"); 478} 479 480#[tokio::test] 481async fn empty_uri_list_is_rejected() { 482 let h = Harness::new().await; 483 let app = router(h.state.clone()); 484 let resp = app 485 .oneshot( 486 Request::builder() 487 .uri("/xrpc/sh.tangled.repo.getRepos") 488 .body(Body::empty()) 489 .unwrap(), 490 ) 491 .await 492 .unwrap(); 493 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 494} 495 496#[tokio::test] 497async fn over_limit_uri_list_is_rejected() { 498 let h = Harness::new().await; 499 let uris: Vec<String> = (0..51) 500 .map(|i| format!("at://did:plc:nel/sh.tangled.repo/r{i}")) 501 .collect(); 502 let refs: Vec<&str> = uris.iter().map(|s| s.as_str()).collect(); 503 let app = router(h.state.clone()); 504 let resp = app 505 .oneshot(bulk_request("sh.tangled.repo.getRepos", "repos", &refs)) 506 .await 507 .unwrap(); 508 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 509} 510 511#[tokio::test] 512async fn malformed_uri_in_list_returns_400() { 513 let h = Harness::new().await; 514 let app = router(h.state.clone()); 515 let resp = app 516 .oneshot(bulk_request( 517 "sh.tangled.repo.getRepos", 518 "repos", 519 &["not-a-uri"], 520 )) 521 .await 522 .unwrap(); 523 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 524} 525 526#[tokio::test] 527async fn bulk_items_serialize_a_single_type_key() { 528 let h = Harness::new().await; 529 h.mount( 530 &did("did:plc:teq"), 531 &nsid("sh.tangled.repo"), 532 &rkey("abalone"), 533 repo_body("abalone"), 534 ) 535 .await; 536 let app = router(h.state.clone()); 537 let resp = app 538 .oneshot(bulk_request( 539 "sh.tangled.repo.getRepos", 540 "repos", 541 &["at://did:plc:teq/sh.tangled.repo/abalone"], 542 )) 543 .await 544 .unwrap(); 545 assert_eq!(resp.status(), StatusCode::OK); 546 let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap(); 547 let raw = String::from_utf8(bytes.to_vec()).unwrap(); 548 assert_eq!(raw.matches("\"$type\"").count(), 1, "body: {raw}"); 549} 550 551#[tokio::test] 552async fn get_repos_by_repo_dids_returns_resolved_repos() { 553 let h = Harness::new().await; 554 h.state 555 .resolver 556 .observe( 557 did("did:plc:nel"), 558 rkey("abalone"), 559 Some(did("did:plc:limpet")), 560 None, 561 ) 562 .await; 563 h.state 564 .resolver 565 .observe( 566 did("did:plc:teq"), 567 rkey("coral"), 568 Some(did("did:plc:coral")), 569 None, 570 ) 571 .await; 572 h.mount( 573 &did("did:plc:nel"), 574 &nsid("sh.tangled.repo"), 575 &rkey("abalone"), 576 repo_body("abalone"), 577 ) 578 .await; 579 h.mount( 580 &did("did:plc:teq"), 581 &nsid("sh.tangled.repo"), 582 &rkey("coral"), 583 repo_body("coral"), 584 ) 585 .await; 586 let app = router(h.state.clone()); 587 let (status, body) = json_response( 588 app.oneshot(bulk_request( 589 "sh.tangled.repo.getReposByRepoDids", 590 "dids", 591 &["did:plc:limpet", "did:plc:coral"], 592 )) 593 .await 594 .unwrap(), 595 ) 596 .await; 597 assert_eq!(status, StatusCode::OK, "{body}"); 598 let items = body["items"].as_array().unwrap(); 599 assert_eq!(items.len(), 2, "{body}"); 600 let names: Vec<&str> = items 601 .iter() 602 .map(|v| v["value"]["name"].as_str().unwrap()) 603 .collect(); 604 assert!(names.contains(&"abalone")); 605 assert!(names.contains(&"coral")); 606} 607 608#[tokio::test] 609async fn get_repos_by_repo_dids_skips_unobserved_dids() { 610 let h = Harness::new().await; 611 h.state 612 .resolver 613 .observe( 614 did("did:plc:nel"), 615 rkey("abalone"), 616 Some(did("did:plc:limpet")), 617 None, 618 ) 619 .await; 620 h.mount( 621 &did("did:plc:nel"), 622 &nsid("sh.tangled.repo"), 623 &rkey("abalone"), 624 repo_body("abalone"), 625 ) 626 .await; 627 let app = router(h.state.clone()); 628 let (status, body) = json_response( 629 app.oneshot(bulk_request( 630 "sh.tangled.repo.getReposByRepoDids", 631 "dids", 632 &["did:plc:limpet", "did:plc:ghost"], 633 )) 634 .await 635 .unwrap(), 636 ) 637 .await; 638 // unknown dids are skipped, same as missing bulk records 639 assert_eq!(status, StatusCode::OK, "{body}"); 640 let items = body["items"].as_array().unwrap(); 641 assert_eq!(items.len(), 1, "{body}"); 642 assert_eq!(items[0]["value"]["name"], "abalone"); 643} 644 645#[tokio::test] 646async fn get_repos_by_repo_dids_rejects_bad_requests() { 647 let h = Harness::new().await; 648 let app = router(h.state.clone()); 649 // no dids at all 650 let resp = app 651 .clone() 652 .oneshot( 653 Request::builder() 654 .uri("/xrpc/sh.tangled.repo.getReposByRepoDids") 655 .body(Body::empty()) 656 .unwrap(), 657 ) 658 .await 659 .unwrap(); 660 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 661 // not a did 662 let resp = app 663 .clone() 664 .oneshot(bulk_request( 665 "sh.tangled.repo.getReposByRepoDids", 666 "dids", 667 &["not-a-did"], 668 )) 669 .await 670 .unwrap(); 671 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 672 // over the bulk limit 673 let dids: Vec<String> = (0..51).map(|i| format!("did:plc:d{i}")).collect(); 674 let refs: Vec<&str> = dids.iter().map(|s| s.as_str()).collect(); 675 let resp = app 676 .oneshot(bulk_request( 677 "sh.tangled.repo.getReposByRepoDids", 678 "dids", 679 &refs, 680 )) 681 .await 682 .unwrap(); 683 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 684}