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 / knot_proxy.rs
34 kB 984 lines
1use std::sync::Arc; 2use std::time::Duration; 3 4use axum::body::{Body, to_bytes}; 5use bobbin_edge_index::{CoverageWatch, EdgeStore, StateIndex}; 6use bobbin_knot_proxy::{FailureThreshold, KnotHttpConfig, KnotProxy, KnotProxyConfig}; 7use bobbin_record_lru::{CacheCapacity, LruRecordStore}; 8use bobbin_resolver::RepoIdResolver; 9use bobbin_runtime::{RuntimeHasher, SystemClock}; 10use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader}; 11use bobbin_slingshot_client::SlingshotClient; 12use bobbin_xrpc::{AppState, router}; 13use http::{Request, StatusCode}; 14use jacquard_common::DefaultStr; 15use jacquard_common::types::did::Did; 16use jacquard_common::types::recordkey::Rkey; 17use serde_json::{Value, json}; 18use tower::ServiceExt; 19use url::Url; 20use url::form_urlencoded::byte_serialize; 21use wiremock::matchers::{header_exists, method, path, query_param}; 22use wiremock::{Mock, MockServer, ResponseTemplate}; 23 24const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; 25 26fn did(s: &str) -> Did<DefaultStr> { 27 Did::new_owned(s).unwrap() 28} 29 30fn rkey(s: &str) -> Rkey<DefaultStr> { 31 Rkey::new_owned(s).unwrap() 32} 33 34fn test_config() -> KnotProxyConfig { 35 KnotProxyConfig { 36 failure_threshold: FailureThreshold::new(2).unwrap(), 37 cooldown: Duration::from_millis(80), 38 allow_private_hosts: true, 39 require_https: false, 40 extra_ca_cert: None, 41 } 42} 43 44fn test_http_config() -> KnotHttpConfig { 45 KnotHttpConfig { 46 connect_timeout: Duration::from_millis(500), 47 read_timeout: Duration::from_secs(2), 48 } 49} 50 51struct Harness { 52 slingshot: MockServer, 53 knot: MockServer, 54 state: AppState, 55} 56 57impl Harness { 58 async fn new() -> Self { 59 Self::with_config(test_config()).await 60 } 61 62 async fn with_config(config: KnotProxyConfig) -> Self { 63 let slingshot_server = MockServer::start().await; 64 let knot_server = MockServer::start().await; 65 let state = AppState::new( 66 Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), 67 SlingshotClient::with_default_http(Url::parse(&slingshot_server.uri()).unwrap()) 68 .unwrap(), 69 Arc::new(EdgeStore::new(RuntimeHasher::default())), 70 Arc::new(StateIndex::new(RuntimeHasher::default())), 71 Arc::new(StateIndex::new(RuntimeHasher::default())), 72 Arc::new(CoverageWatch::new()), 73 Arc::new( 74 KnotProxy::new( 75 config, 76 test_http_config(), 77 Arc::new(SystemClock::new()), 78 RuntimeHasher::default(), 79 ) 80 .unwrap(), 81 ), 82 Arc::new( 83 SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES, Arc::new(SystemClock::new())).unwrap(), 84 ) as Arc<dyn SearchReader>, 85 Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), 86 Arc::new(bobbin_xrpc::default_directory()), 87 ); 88 Self { 89 slingshot: slingshot_server, 90 knot: knot_server, 91 state, 92 } 93 } 94 95 async fn mount_repo_record(&self, did: &Did<DefaultStr>, rkey: &Rkey<DefaultStr>, name: &str) { 96 self.mount_repo_record_inner(did, rkey, Some(name)).await; 97 } 98 99 async fn mount_repo_record_rkey_as_name(&self, did: &Did<DefaultStr>, rkey: &Rkey<DefaultStr>) { 100 self.mount_repo_record_inner(did, rkey, None).await; 101 } 102 103 async fn mount_repo_record_inner( 104 &self, 105 did: &Did<DefaultStr>, 106 rkey: &Rkey<DefaultStr>, 107 name: Option<&str>, 108 ) { 109 let knot_value = self.knot.uri(); 110 let mut record = json!({ 111 "$type": "sh.tangled.repo", 112 "createdAt": "2026-05-01T00:00:00Z", 113 "knot": knot_value, 114 }); 115 if let Some(n) = name { 116 record["name"] = json!(n); 117 } 118 let uri = format!("at://{}/sh.tangled.repo/{}", did.as_ref(), rkey.as_ref()); 119 Mock::given(method("GET")) 120 .and(path("/xrpc/com.atproto.repo.getRecord")) 121 .and(query_param("repo", did.as_ref())) 122 .and(query_param("collection", "sh.tangled.repo")) 123 .and(query_param("rkey", rkey.as_ref())) 124 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 125 "uri": uri, 126 "cid": CID, 127 "value": record, 128 }))) 129 .mount(&self.slingshot) 130 .await; 131 } 132 133 async fn call(&self, path_and_query: &str) -> http::Response<Body> { 134 self.call_with_headers(path_and_query, &[]).await 135 } 136 137 async fn call_with_headers( 138 &self, 139 path_and_query: &str, 140 client_headers: &[(&str, &str)], 141 ) -> http::Response<Body> { 142 let builder = client_headers 143 .iter() 144 .fold(Request::builder().uri(path_and_query), |b, (k, v)| { 145 b.header(*k, *v) 146 }); 147 router(self.state.clone()) 148 .oneshot(builder.body(Body::empty()).unwrap()) 149 .await 150 .expect("router infallible") 151 } 152} 153 154fn enc(s: &str) -> String { 155 byte_serialize(s.as_bytes()).collect() 156} 157 158async fn body_string(resp: http::Response<Body>) -> String { 159 let body = to_bytes(resp.into_body(), 64 * 1024).await.unwrap(); 160 String::from_utf8(body.to_vec()).expect("response body is utf-8") 161} 162 163async fn body_value(resp: http::Response<Body>) -> Value { 164 let s = body_string(resp).await; 165 serde_json::from_str(&s).unwrap_or_else(|e| panic!("body not json: {e}: {s}")) 166} 167 168#[tokio::test] 169async fn proxies_repo_blob_with_did_slash_name_repo_param() { 170 let h = Harness::new().await; 171 let tid = "3jzfcijpj2z2a"; 172 h.mount_repo_record(&did("did:plc:abalone"), &rkey(tid), "barnacle") 173 .await; 174 Mock::given(method("GET")) 175 .and(path("/xrpc/sh.tangled.repo.blob")) 176 .and(query_param("repo", "did:plc:abalone/barnacle")) 177 .and(query_param("ref", "main")) 178 .and(query_param("path", "README.md")) 179 .respond_with( 180 ResponseTemplate::new(200) 181 .set_body_raw(r#"{"path":"README.md","content":"hi"}"#, "application/json"), 182 ) 183 .mount(&h.knot) 184 .await; 185 186 let target = format!( 187 "/xrpc/sh.tangled.repo.blob?repo={}&ref=main&path=README.md", 188 enc(&format!("at://did:plc:abalone/sh.tangled.repo/{tid}")), 189 ); 190 let resp = h.call(&target).await; 191 assert_eq!(resp.status(), StatusCode::OK); 192 assert_eq!( 193 resp.headers().get("content-type").unwrap(), 194 "application/json", 195 ); 196 let v = body_value(resp).await; 197 assert_eq!(v["path"], "README.md"); 198 assert_eq!(v["content"], "hi"); 199} 200 201#[tokio::test] 202async fn modern_rkey_as_name_uses_rkey_even_when_name_field_set() { 203 let h = Harness::new().await; 204 h.mount_repo_record(&did("did:plc:abalone"), &rkey("core"), "Tangled Core") 205 .await; 206 Mock::given(method("GET")) 207 .and(path("/xrpc/sh.tangled.repo.getDefaultBranch")) 208 .and(query_param("repo", "did:plc:abalone/core")) 209 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 210 "hash": "abc", 211 "name": "main", 212 "when": "2026-05-01T00:00:00Z", 213 }))) 214 .mount(&h.knot) 215 .await; 216 217 let target = format!( 218 "/xrpc/sh.tangled.repo.getDefaultBranch?repo={}", 219 enc("at://did:plc:abalone/sh.tangled.repo/core"), 220 ); 221 let resp = h.call(&target).await; 222 assert_eq!(resp.status(), StatusCode::OK); 223 let v = body_value(resp).await; 224 assert_eq!(v["name"], "main"); 225} 226 227#[tokio::test] 228async fn modern_rkey_as_name_works_when_name_field_null() { 229 let h = Harness::new().await; 230 h.mount_repo_record_rkey_as_name(&did("did:plc:abalone"), &rkey("core")) 231 .await; 232 Mock::given(method("GET")) 233 .and(path("/xrpc/sh.tangled.repo.getDefaultBranch")) 234 .and(query_param("repo", "did:plc:abalone/core")) 235 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 236 "hash": "abc", 237 "name": "main", 238 "when": "2026-05-01T00:00:00Z", 239 }))) 240 .mount(&h.knot) 241 .await; 242 243 let target = format!( 244 "/xrpc/sh.tangled.repo.getDefaultBranch?repo={}", 245 enc("at://did:plc:abalone/sh.tangled.repo/core"), 246 ); 247 let resp = h.call(&target).await; 248 assert_eq!(resp.status(), StatusCode::OK); 249} 250 251#[tokio::test] 252async fn legacy_tid_rkey_falls_back_to_name_field() { 253 let h = Harness::new().await; 254 let tid_rkey = "3jzfcijpj2z2a"; 255 h.mount_repo_record(&did("did:plc:abalone"), &rkey(tid_rkey), "dotfiles") 256 .await; 257 Mock::given(method("GET")) 258 .and(path("/xrpc/sh.tangled.repo.getDefaultBranch")) 259 .and(query_param("repo", "did:plc:abalone/dotfiles")) 260 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 261 "hash": "abc", 262 "name": "main", 263 "when": "2026-05-01T00:00:00Z", 264 }))) 265 .mount(&h.knot) 266 .await; 267 268 let target = format!( 269 "/xrpc/sh.tangled.repo.getDefaultBranch?repo={}", 270 enc(&format!("at://did:plc:abalone/sh.tangled.repo/{tid_rkey}")), 271 ); 272 let resp = h.call(&target).await; 273 assert_eq!(resp.status(), StatusCode::OK); 274} 275 276#[tokio::test] 277async fn tid_rkey_without_name_falls_back_to_tid() { 278 let h = Harness::new().await; 279 let tid_rkey = "3jzfcijpj2z2a"; 280 h.mount_repo_record_rkey_as_name(&did("did:plc:abalone"), &rkey(tid_rkey)) 281 .await; 282 Mock::given(method("GET")) 283 .and(path("/xrpc/sh.tangled.repo.getDefaultBranch")) 284 .and(query_param("repo", format!("did:plc:abalone/{tid_rkey}"))) 285 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 286 "hash": "abc", 287 "name": "main", 288 "when": "2026-05-01T00:00:00Z", 289 }))) 290 .mount(&h.knot) 291 .await; 292 let target = format!( 293 "/xrpc/sh.tangled.repo.getDefaultBranch?repo={}", 294 enc(&format!("at://did:plc:abalone/sh.tangled.repo/{tid_rkey}")), 295 ); 296 let resp = h.call(&target).await; 297 assert_eq!(resp.status(), StatusCode::OK); 298} 299 300#[tokio::test] 301async fn streams_binary_archive_through_proxy() { 302 let h = Harness::new().await; 303 let tid = "3jzfcijpj2z2b"; 304 h.mount_repo_record(&did("did:plc:limpet"), &rkey(tid), "kelp") 305 .await; 306 let payload: Vec<u8> = (0u8..=255).collect(); 307 Mock::given(method("GET")) 308 .and(path("/xrpc/sh.tangled.repo.archive")) 309 .and(query_param("repo", "did:plc:limpet/kelp")) 310 .and(query_param("ref", "v1")) 311 .respond_with( 312 ResponseTemplate::new(200) 313 .insert_header("content-type", "application/gzip") 314 .set_body_bytes(payload.clone()), 315 ) 316 .mount(&h.knot) 317 .await; 318 319 let target = format!( 320 "/xrpc/sh.tangled.repo.archive?repo={}&ref=v1", 321 enc(&format!("at://did:plc:limpet/sh.tangled.repo/{tid}")), 322 ); 323 let resp = h.call(&target).await; 324 assert_eq!(resp.status(), StatusCode::OK); 325 assert_eq!( 326 resp.headers().get("content-type").unwrap(), 327 "application/gzip", 328 ); 329 let body = to_bytes(resp.into_body(), 4 * 1024).await.unwrap(); 330 assert_eq!(body.as_ref(), payload.as_slice()); 331} 332 333#[tokio::test] 334async fn missing_repo_param_returns_400() { 335 let h = Harness::new().await; 336 let resp = h.call("/xrpc/sh.tangled.repo.blob?ref=main").await; 337 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 338 let v = body_value(resp).await; 339 assert_eq!(v["error"], "InvalidRequest"); 340} 341 342#[tokio::test] 343async fn unknown_repo_propagates_404_from_slingshot() { 344 let h = Harness::new().await; 345 Mock::given(method("GET")) 346 .and(path("/xrpc/com.atproto.repo.getRecord")) 347 .respond_with(ResponseTemplate::new(404).set_body_string("not found")) 348 .mount(&h.slingshot) 349 .await; 350 let target = format!( 351 "/xrpc/sh.tangled.repo.blob?repo={}&ref=main", 352 enc("at://did:plc:abalone/sh.tangled.repo/missing"), 353 ); 354 let resp = h.call(&target).await; 355 assert_eq!(resp.status(), StatusCode::NOT_FOUND); 356 let v = body_value(resp).await; 357 assert_eq!(v["error"], "RecordNotFound"); 358} 359 360#[tokio::test] 361async fn knot_5xx_routes_to_upstream_failed() { 362 let h = Harness::new().await; 363 h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") 364 .await; 365 Mock::given(method("GET")) 366 .and(path("/xrpc/sh.tangled.repo.blob")) 367 .respond_with(ResponseTemplate::new(503)) 368 .mount(&h.knot) 369 .await; 370 let target = format!( 371 "/xrpc/sh.tangled.repo.blob?repo={}&ref=main&path=x", 372 enc("at://did:plc:abalone/sh.tangled.repo/r1"), 373 ); 374 let resp = h.call(&target).await; 375 assert_eq!(resp.status(), StatusCode::BAD_GATEWAY); 376 let v = body_value(resp).await; 377 assert_eq!(v["error"], "UpstreamFailed"); 378} 379 380#[tokio::test] 381async fn knot_4xx_passes_through_unchanged() { 382 let h = Harness::new().await; 383 h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") 384 .await; 385 Mock::given(method("GET")) 386 .and(path("/xrpc/sh.tangled.repo.blob")) 387 .respond_with(ResponseTemplate::new(404).set_body_raw( 388 r#"{"error":"FileNotFound","message":"nope"}"#, 389 "application/json", 390 )) 391 .mount(&h.knot) 392 .await; 393 let target = format!( 394 "/xrpc/sh.tangled.repo.blob?repo={}&ref=main&path=missing", 395 enc("at://did:plc:abalone/sh.tangled.repo/r1"), 396 ); 397 let resp = h.call(&target).await; 398 assert_eq!(resp.status(), StatusCode::NOT_FOUND); 399 let v = body_value(resp).await; 400 assert_eq!(v["error"], "FileNotFound"); 401} 402 403#[tokio::test] 404async fn breaker_opens_after_threshold_then_short_circuits() { 405 let h = Harness::new().await; 406 h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") 407 .await; 408 Mock::given(method("GET")) 409 .and(path("/xrpc/sh.tangled.repo.blob")) 410 .respond_with(ResponseTemplate::new(503)) 411 .mount(&h.knot) 412 .await; 413 let target = format!( 414 "/xrpc/sh.tangled.repo.blob?repo={}&ref=main&path=x", 415 enc("at://did:plc:abalone/sh.tangled.repo/r1"), 416 ); 417 let r1 = h.call(&target).await; 418 assert_eq!(r1.status(), StatusCode::BAD_GATEWAY); 419 let _ = body_string(r1).await; 420 let r2 = h.call(&target).await; 421 assert_eq!(r2.status(), StatusCode::BAD_GATEWAY); 422 let _ = body_string(r2).await; 423 let r3 = h.call(&target).await; 424 assert_eq!(r3.status(), StatusCode::BAD_GATEWAY); 425 let v = body_value(r3).await; 426 assert!( 427 v["message"] 428 .as_str() 429 .unwrap_or_default() 430 .contains("circuit breaker open"), 431 "third call must be short-circuited by breaker, got {v}", 432 ); 433} 434 435#[tokio::test] 436async fn proxy_owner_uses_knot_query_param() { 437 let h = Harness::new().await; 438 Mock::given(method("GET")) 439 .and(path("/xrpc/sh.tangled.owner")) 440 .respond_with( 441 ResponseTemplate::new(200) 442 .set_body_raw(r#"{"owner":"did:plc:nautilus"}"#, "application/json"), 443 ) 444 .mount(&h.knot) 445 .await; 446 let target = format!("/xrpc/sh.tangled.owner?knot={}", enc(&h.knot.uri())); 447 let resp = h.call(&target).await; 448 assert_eq!(resp.status(), StatusCode::OK); 449 let v = body_value(resp).await; 450 assert_eq!(v["owner"], "did:plc:nautilus"); 451} 452 453#[tokio::test] 454async fn proxy_knot_version_uses_knot_query_param() { 455 let h = Harness::new().await; 456 Mock::given(method("GET")) 457 .and(path("/xrpc/sh.tangled.knot.version")) 458 .respond_with( 459 ResponseTemplate::new(200).set_body_raw(r#"{"version":"0.42"}"#, "application/json"), 460 ) 461 .mount(&h.knot) 462 .await; 463 let target = format!("/xrpc/sh.tangled.knot.version?knot={}", enc(&h.knot.uri())); 464 let resp = h.call(&target).await; 465 assert_eq!(resp.status(), StatusCode::OK); 466 let v = body_value(resp).await; 467 assert_eq!(v["version"], "0.42"); 468} 469 470#[tokio::test] 471async fn proxy_knot_list_keys_forwards_pagination_params() { 472 let h = Harness::new().await; 473 Mock::given(method("GET")) 474 .and(path("/xrpc/sh.tangled.knot.listKeys")) 475 .and(query_param("limit", "5")) 476 .and(query_param("cursor", "abc")) 477 .respond_with(ResponseTemplate::new(200).set_body_raw(r#"{"keys":[]}"#, "application/json")) 478 .mount(&h.knot) 479 .await; 480 let target = format!( 481 "/xrpc/sh.tangled.knot.listKeys?knot={}&limit=5&cursor=abc", 482 enc(&h.knot.uri()), 483 ); 484 let resp = h.call(&target).await; 485 assert_eq!(resp.status(), StatusCode::OK); 486} 487 488#[tokio::test] 489async fn missing_knot_param_on_knot_route_returns_400() { 490 let h = Harness::new().await; 491 let resp = h.call("/xrpc/sh.tangled.knot.version").await; 492 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 493 let v = body_value(resp).await; 494 assert_eq!(v["error"], "InvalidRequest"); 495} 496 497#[tokio::test] 498async fn second_proxy_call_skips_slingshot_via_lru() { 499 let h = Harness::new().await; 500 let tid = "3jzfcijpj2z2c"; 501 h.mount_repo_record(&did("did:plc:abalone"), &rkey(tid), "barnacle") 502 .await; 503 Mock::given(method("GET")) 504 .and(path("/xrpc/sh.tangled.repo.tree")) 505 .and(query_param("repo", "did:plc:abalone/barnacle")) 506 .and(query_param("ref", "main")) 507 .respond_with( 508 ResponseTemplate::new(200) 509 .set_body_raw(r#"{"ref":"main","files":[]}"#, "application/json"), 510 ) 511 .mount(&h.knot) 512 .await; 513 let target = format!( 514 "/xrpc/sh.tangled.repo.tree?repo={}&ref=main", 515 enc(&format!("at://did:plc:abalone/sh.tangled.repo/{tid}")), 516 ); 517 let r1 = h.call(&target).await; 518 assert_eq!(r1.status(), StatusCode::OK); 519 let _ = body_string(r1).await; 520 let r2 = h.call(&target).await; 521 assert_eq!(r2.status(), StatusCode::OK); 522 let _ = body_string(r2).await; 523 let received = h.slingshot.received_requests().await.unwrap(); 524 let getrecord = received 525 .iter() 526 .filter(|r| r.url.path() == "/xrpc/com.atproto.repo.getRecord") 527 .count(); 528 assert_eq!( 529 getrecord, 1, 530 "slingshot must be hit exactly once because the LRU serves the second proxy call", 531 ); 532} 533 534#[tokio::test] 535async fn does_not_inject_auth_or_atproto_proxy_headers() { 536 let h = Harness::new().await; 537 h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") 538 .await; 539 Mock::given(method("GET")) 540 .and(path("/xrpc/sh.tangled.repo.blob")) 541 .and(header_exists("user-agent")) 542 .respond_with( 543 ResponseTemplate::new(200).set_body_raw(r#"{"path":"x"}"#, "application/json"), 544 ) 545 .mount(&h.knot) 546 .await; 547 let target = format!( 548 "/xrpc/sh.tangled.repo.blob?repo={}&ref=main&path=x", 549 enc("at://did:plc:abalone/sh.tangled.repo/r1"), 550 ); 551 let resp = h.call(&target).await; 552 assert_eq!(resp.status(), StatusCode::OK); 553 let received = h.knot.received_requests().await.unwrap(); 554 let knot_call = received 555 .iter() 556 .find(|r| r.url.path() == "/xrpc/sh.tangled.repo.blob") 557 .expect("knot received the proxied call"); 558 assert!( 559 knot_call.headers.get("authorization").is_none(), 560 "bobbin must not inject auth, anonymous read by design", 561 ); 562 assert!( 563 knot_call.headers.get("atproto-proxy").is_none(), 564 "bobbin is not an atproto-proxy chain", 565 ); 566 assert!( 567 knot_call.headers.get("atproto-accept-labelers").is_none(), 568 "bobbin does not negotiate labelers with knots", 569 ); 570} 571 572#[tokio::test] 573async fn forwards_range_and_conditional_request_headers() { 574 let h = Harness::new().await; 575 let tid = "3jzfcijpj2z2d"; 576 h.mount_repo_record(&did("did:plc:limpet"), &rkey(tid), "kelp") 577 .await; 578 Mock::given(method("GET")) 579 .and(path("/xrpc/sh.tangled.repo.archive")) 580 .and(query_param("repo", "did:plc:limpet/kelp")) 581 .respond_with( 582 ResponseTemplate::new(206) 583 .insert_header("content-type", "application/octet-stream") 584 .insert_header("content-range", "bytes 0-99/2048") 585 .insert_header("accept-ranges", "bytes") 586 .insert_header("etag", "\"v1\"") 587 .set_body_bytes(vec![0u8; 100]), 588 ) 589 .mount(&h.knot) 590 .await; 591 592 let target = format!( 593 "/xrpc/sh.tangled.repo.archive?repo={}&ref=v1", 594 enc(&format!("at://did:plc:limpet/sh.tangled.repo/{tid}")), 595 ); 596 let resp = h 597 .call_with_headers( 598 &target, 599 &[ 600 ("range", "bytes=0-99"), 601 ("if-none-match", "\"old\""), 602 ("if-modified-since", "Wed, 01 May 2026 00:00:00 GMT"), 603 ], 604 ) 605 .await; 606 assert_eq!(resp.status(), StatusCode::PARTIAL_CONTENT); 607 assert_eq!( 608 resp.headers().get("content-range").unwrap(), 609 "bytes 0-99/2048" 610 ); 611 assert_eq!(resp.headers().get("accept-ranges").unwrap(), "bytes"); 612 assert_eq!(resp.headers().get("etag").unwrap(), "\"v1\""); 613 614 let received = h.knot.received_requests().await.unwrap(); 615 let knot_call = received 616 .iter() 617 .find(|r| r.url.path() == "/xrpc/sh.tangled.repo.archive") 618 .expect("knot received the proxied call"); 619 assert_eq!(knot_call.headers.get("range").unwrap(), "bytes=0-99"); 620 assert_eq!(knot_call.headers.get("if-none-match").unwrap(), "\"old\""); 621 assert_eq!( 622 knot_call.headers.get("if-modified-since").unwrap(), 623 "Wed, 01 May 2026 00:00:00 GMT", 624 ); 625} 626 627#[tokio::test] 628async fn drops_disallowed_client_headers() { 629 let h = Harness::new().await; 630 h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") 631 .await; 632 Mock::given(method("GET")) 633 .and(path("/xrpc/sh.tangled.repo.blob")) 634 .respond_with( 635 ResponseTemplate::new(200).set_body_raw(r#"{"path":"x"}"#, "application/json"), 636 ) 637 .mount(&h.knot) 638 .await; 639 let target = format!( 640 "/xrpc/sh.tangled.repo.blob?repo={}&path=x", 641 enc("at://did:plc:abalone/sh.tangled.repo/r1"), 642 ); 643 let resp = h 644 .call_with_headers( 645 &target, 646 &[ 647 ("authorization", "Bearer secret"), 648 ("cookie", "sid=evil"), 649 ("x-custom", "should-not-pass"), 650 ], 651 ) 652 .await; 653 assert_eq!(resp.status(), StatusCode::OK); 654 let received = h.knot.received_requests().await.unwrap(); 655 let knot_call = received 656 .iter() 657 .find(|r| r.url.path() == "/xrpc/sh.tangled.repo.blob") 658 .expect("knot received the proxied call"); 659 assert!(knot_call.headers.get("authorization").is_none()); 660 assert!(knot_call.headers.get("cookie").is_none()); 661 assert!(knot_call.headers.get("x-custom").is_none()); 662} 663 664#[tokio::test] 665async fn rejects_client_supplied_loopback_under_strict_config() { 666 let strict = KnotProxyConfig { 667 allow_private_hosts: false, 668 ..test_config() 669 }; 670 let h = Harness::with_config(strict).await; 671 let resp = h 672 .call(&format!( 673 "/xrpc/sh.tangled.knot.version?knot={}", 674 enc("http://127.0.0.1:9"), 675 )) 676 .await; 677 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 678 let v = body_value(resp).await; 679 assert_eq!(v["error"], "InvalidRequest"); 680 let msg = v["message"].as_str().unwrap_or_default().to_owned(); 681 assert!( 682 msg.contains("loopback") || msg.contains("blocked"), 683 "message should explain block reason, got {msg}", 684 ); 685} 686 687#[tokio::test] 688async fn rejects_client_supplied_link_local_metadata_endpoint() { 689 let strict = KnotProxyConfig { 690 allow_private_hosts: false, 691 ..test_config() 692 }; 693 let h = Harness::with_config(strict).await; 694 let resp = h 695 .call(&format!( 696 "/xrpc/sh.tangled.knot.version?knot={}", 697 enc("http://169.254.169.254"), 698 )) 699 .await; 700 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 701 let v = body_value(resp).await; 702 assert_eq!(v["error"], "InvalidRequest"); 703} 704 705#[tokio::test] 706async fn record_with_private_knot_returns_invalid_record() { 707 let strict = KnotProxyConfig { 708 allow_private_hosts: false, 709 ..test_config() 710 }; 711 let h = Harness::with_config(strict).await; 712 let owner = did("did:plc:abalone"); 713 let rk = rkey("r1"); 714 let record = json!({ 715 "$type": "sh.tangled.repo", 716 "createdAt": "2026-05-01T00:00:00Z", 717 "knot": "http://10.0.0.5:3000", 718 "name": "barnacle", 719 }); 720 let uri = format!("at://{}/sh.tangled.repo/{}", owner.as_ref(), rk.as_ref()); 721 Mock::given(method("GET")) 722 .and(path("/xrpc/com.atproto.repo.getRecord")) 723 .and(query_param("repo", owner.as_ref())) 724 .and(query_param("collection", "sh.tangled.repo")) 725 .and(query_param("rkey", rk.as_ref())) 726 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 727 "uri": uri, 728 "cid": CID, 729 "value": record, 730 }))) 731 .mount(&h.slingshot) 732 .await; 733 let resp = h 734 .call(&format!( 735 "/xrpc/sh.tangled.repo.blob?repo={}&path=x", 736 enc(&uri), 737 )) 738 .await; 739 assert_eq!(resp.status(), StatusCode::BAD_GATEWAY); 740 let v = body_value(resp).await; 741 assert_eq!(v["error"], "InvalidRecord"); 742} 743 744#[tokio::test] 745async fn strips_basic_auth_from_credentialed_knot_url() { 746 let h = Harness::new().await; 747 let parsed = Url::parse(&h.knot.uri()).unwrap(); 748 let knot_with_creds = format!( 749 "{}://attacker:secret@{}:{}/", 750 parsed.scheme(), 751 parsed.host_str().unwrap(), 752 parsed.port().unwrap(), 753 ); 754 let owner = did("did:plc:abalone"); 755 let rk = rkey("r1"); 756 let record = json!({ 757 "$type": "sh.tangled.repo", 758 "createdAt": "2026-05-01T00:00:00Z", 759 "knot": knot_with_creds, 760 "name": "barnacle", 761 }); 762 let uri = format!("at://{}/sh.tangled.repo/{}", owner.as_ref(), rk.as_ref()); 763 Mock::given(method("GET")) 764 .and(path("/xrpc/com.atproto.repo.getRecord")) 765 .and(query_param("repo", owner.as_ref())) 766 .and(query_param("collection", "sh.tangled.repo")) 767 .and(query_param("rkey", rk.as_ref())) 768 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 769 "uri": uri, 770 "cid": CID, 771 "value": record, 772 }))) 773 .mount(&h.slingshot) 774 .await; 775 Mock::given(method("GET")) 776 .and(path("/xrpc/sh.tangled.repo.blob")) 777 .respond_with( 778 ResponseTemplate::new(200).set_body_raw(r#"{"path":"x"}"#, "application/json"), 779 ) 780 .mount(&h.knot) 781 .await; 782 let target = format!("/xrpc/sh.tangled.repo.blob?repo={}&path=x", enc(&uri)); 783 let resp = h.call(&target).await; 784 assert_eq!(resp.status(), StatusCode::OK); 785 let received = h.knot.received_requests().await.unwrap(); 786 let knot_call = received 787 .iter() 788 .find(|r| r.url.path() == "/xrpc/sh.tangled.repo.blob") 789 .expect("knot received the proxied call"); 790 assert!( 791 knot_call.headers.get("authorization").is_none(), 792 "userinfo in knot field must not become an Authorization header", 793 ); 794} 795 796#[tokio::test] 797async fn knot_redirect_surfaces_as_upstream_failed() { 798 let h = Harness::new().await; 799 let secondary = MockServer::start().await; 800 h.mount_repo_record(&did("did:plc:abalone"), &rkey("r1"), "barnacle") 801 .await; 802 Mock::given(method("GET")) 803 .and(path("/xrpc/sh.tangled.repo.blob")) 804 .respond_with( 805 ResponseTemplate::new(302) 806 .insert_header("location", &format!("{}/secret", secondary.uri())), 807 ) 808 .mount(&h.knot) 809 .await; 810 Mock::given(method("GET")) 811 .and(path("/secret")) 812 .respond_with(ResponseTemplate::new(200).set_body_string("leaked")) 813 .mount(&secondary) 814 .await; 815 let resp = h 816 .call(&format!( 817 "/xrpc/sh.tangled.repo.blob?repo={}&path=x", 818 enc("at://did:plc:abalone/sh.tangled.repo/r1"), 819 )) 820 .await; 821 assert_eq!(resp.status(), StatusCode::BAD_GATEWAY); 822 let v = body_value(resp).await; 823 assert_eq!(v["error"], "UpstreamFailed"); 824 let received = secondary.received_requests().await.unwrap(); 825 assert!(received.is_empty(), "redirect target must not be dialled"); 826} 827 828#[tokio::test] 829async fn forwards_repeated_query_params() { 830 let h = Harness::new().await; 831 h.mount_repo_record(&did("did:plc:limpet"), &rkey("r4"), "kelp") 832 .await; 833 Mock::given(method("GET")) 834 .and(path("/xrpc/sh.tangled.repo.tags")) 835 .respond_with(ResponseTemplate::new(200).set_body_raw(r#"{"tags":[]}"#, "application/json")) 836 .mount(&h.knot) 837 .await; 838 let target = format!( 839 "/xrpc/sh.tangled.repo.tags?repo={}&filter=alpha&filter=beta", 840 enc("at://did:plc:limpet/sh.tangled.repo/r4"), 841 ); 842 let resp = h.call(&target).await; 843 assert_eq!(resp.status(), StatusCode::OK); 844 let received = h.knot.received_requests().await.unwrap(); 845 let knot_call = received 846 .iter() 847 .find(|r| r.url.path() == "/xrpc/sh.tangled.repo.tags") 848 .expect("knot received the proxied call"); 849 let filters: Vec<String> = knot_call 850 .url 851 .query_pairs() 852 .filter(|(k, _)| k == "filter") 853 .map(|(_, v)| v.into_owned()) 854 .collect(); 855 assert_eq!(filters, vec!["alpha".to_owned(), "beta".to_owned()]); 856} 857 858#[tokio::test] 859async fn duplicate_repo_param_rejected_as_invalid_request() { 860 let h = Harness::new().await; 861 let target = format!( 862 "/xrpc/sh.tangled.repo.blob?repo={}&repo={}", 863 enc("at://did:plc:abalone/sh.tangled.repo/r1"), 864 enc("at://did:plc:limpet/sh.tangled.repo/r2"), 865 ); 866 let resp = h.call(&target).await; 867 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 868 let v = body_value(resp).await; 869 assert_eq!(v["error"], "InvalidRequest"); 870 assert!( 871 v["message"] 872 .as_str() 873 .unwrap_or_default() 874 .contains("repo parameter must appear at most once"), 875 "got {v}", 876 ); 877} 878 879#[tokio::test] 880async fn duplicate_knot_param_rejected_as_invalid_request() { 881 let h = Harness::new().await; 882 let target = format!( 883 "/xrpc/sh.tangled.knot.version?knot={}&knot={}", 884 enc("https://oyster.cafe"), 885 enc("https://nel.pet"), 886 ); 887 let resp = h.call(&target).await; 888 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 889 let v = body_value(resp).await; 890 assert_eq!(v["error"], "InvalidRequest"); 891} 892 893#[tokio::test] 894async fn rejects_client_supplied_plaintext_when_https_required() { 895 let strict = KnotProxyConfig { 896 require_https: true, 897 ..test_config() 898 }; 899 let h = Harness::with_config(strict).await; 900 let resp = h 901 .call(&format!( 902 "/xrpc/sh.tangled.knot.version?knot={}", 903 enc("http://oyster.cafe"), 904 )) 905 .await; 906 assert_eq!(resp.status(), StatusCode::BAD_REQUEST); 907 let v = body_value(resp).await; 908 assert_eq!(v["error"], "InvalidRequest"); 909 assert!( 910 v["message"] 911 .as_str() 912 .unwrap_or_default() 913 .contains("must be https"), 914 "got {v}", 915 ); 916} 917 918#[tokio::test] 919async fn record_with_plaintext_knot_returns_invalid_record_when_https_required() { 920 let strict = KnotProxyConfig { 921 require_https: true, 922 allow_private_hosts: true, 923 ..test_config() 924 }; 925 let h = Harness::with_config(strict).await; 926 let owner = did("did:plc:abalone"); 927 let rk = rkey("r1"); 928 let record = json!({ 929 "$type": "sh.tangled.repo", 930 "createdAt": "2026-05-01T00:00:00Z", 931 "knot": "http://oyster.cafe", 932 "name": "barnacle", 933 }); 934 let uri = format!("at://{}/sh.tangled.repo/{}", owner.as_ref(), rk.as_ref()); 935 Mock::given(method("GET")) 936 .and(path("/xrpc/com.atproto.repo.getRecord")) 937 .and(query_param("repo", owner.as_ref())) 938 .and(query_param("collection", "sh.tangled.repo")) 939 .and(query_param("rkey", rk.as_ref())) 940 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 941 "uri": uri, 942 "cid": CID, 943 "value": record, 944 }))) 945 .mount(&h.slingshot) 946 .await; 947 let resp = h 948 .call(&format!( 949 "/xrpc/sh.tangled.repo.blob?repo={}&path=x", 950 enc(&uri), 951 )) 952 .await; 953 assert_eq!(resp.status(), StatusCode::BAD_GATEWAY); 954 let v = body_value(resp).await; 955 assert_eq!(v["error"], "InvalidRecord"); 956 assert!( 957 v["message"] 958 .as_str() 959 .unwrap_or_default() 960 .contains("requires https"), 961 "got {v}", 962 ); 963} 964 965#[tokio::test] 966async fn knot_not_modified_passes_through() { 967 let h = Harness::new().await; 968 h.mount_repo_record(&did("did:plc:limpet"), &rkey("r5"), "kelp") 969 .await; 970 Mock::given(method("GET")) 971 .and(path("/xrpc/sh.tangled.repo.archive")) 972 .respond_with(ResponseTemplate::new(304).insert_header("etag", "\"v1\"")) 973 .mount(&h.knot) 974 .await; 975 let target = format!( 976 "/xrpc/sh.tangled.repo.archive?repo={}&ref=v1", 977 enc("at://did:plc:limpet/sh.tangled.repo/r5"), 978 ); 979 let resp = h 980 .call_with_headers(&target, &[("if-none-match", "\"v1\"")]) 981 .await; 982 assert_eq!(resp.status(), StatusCode::NOT_MODIFIED); 983 assert_eq!(resp.headers().get("etag").unwrap(), "\"v1\""); 984}