This repository has no description
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}