1use std::task::{Context, Poll};
4
5use bytes::{Buf, Bytes, BytesMut};
6use http_body::Body;
7use sha2::{Digest, Sha256};
8use tower_service::Service;
9
10use crate::evaluator::{ChioEvaluator, EvaluationInput};
11use crate::request_metadata::RequestMetadata;
12
13pub const DEFAULT_MAX_BODY_BYTES: usize = 8 * 1024 * 1024;
21
22#[derive(Clone)]
33pub struct ChioService<S> {
34 inner: S,
35 evaluator: ChioEvaluator,
36 max_body_bytes: usize,
37}
38
39impl<S> ChioService<S> {
40 pub fn new(inner: S, evaluator: ChioEvaluator) -> Self {
42 crate::metrics::seed_fail_open_series();
43 Self {
44 inner,
45 evaluator,
46 max_body_bytes: DEFAULT_MAX_BODY_BYTES,
47 }
48 }
49
50 #[must_use]
56 pub fn with_max_body_bytes(mut self, max_body_bytes: usize) -> Self {
57 self.max_body_bytes = max_body_bytes;
58 self
59 }
60}
61
62impl<S, ReqBody, ResBody> Service<http::Request<ReqBody>> for ChioService<S>
63where
64 S: Service<http::Request<ReqBody>, Response = http::Response<ResBody>> + Clone + Send + 'static,
65 S::Future: Send,
66 S::Error: Into<Box<dyn std::error::Error + Send + Sync>>,
67 ReqBody: Body + From<Bytes> + Send + 'static,
68 ReqBody::Data: Send,
69 ReqBody::Error: Into<Box<dyn std::error::Error + Send + Sync>>,
70 ResBody: Default + From<Bytes> + Send + 'static,
71{
72 type Response = http::Response<ResBody>;
73 type Error = Box<dyn std::error::Error + Send + Sync>;
74 type Future = std::pin::Pin<
75 Box<dyn std::future::Future<Output = Result<Self::Response, Self::Error>> + Send>,
76 >;
77
78 fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
79 self.inner.poll_ready(cx).map_err(Into::into)
80 }
81
82 fn call(&mut self, req: http::Request<ReqBody>) -> Self::Future {
83 let evaluator = self.evaluator.clone();
84 let mut inner = self.inner.clone();
85 let max_body_bytes = self.max_body_bytes;
86
87 Box::pin(async move {
88 let method = req.method().as_str().to_string();
89 let path = req.uri().path().to_string();
90 let request_metadata = RequestMetadata::parse(req.uri().query());
91 let headers = req.headers().clone();
92
93 let identity_fn = evaluator.identity_extractor();
98 let caller = identity_fn(&headers);
99
100 let (req, body_hash, body_length) = match buffer_request_body(req, max_body_bytes).await
101 {
102 Ok(parts) => parts,
103 Err(BufferBodyError::TooLarge) => {
104 return Ok(build_payload_too_large_response::<ResBody>(
105 &evaluator,
106 &method,
107 &path,
108 &caller,
109 max_body_bytes,
110 ));
111 }
112 Err(BufferBodyError::Inner(error)) => return Err(error),
113 };
114
115 let evaluation_input = EvaluationInput {
117 method: &method,
118 path: &path,
119 query: request_metadata.query(),
120 caller,
121 headers: &headers,
122 body_hash,
123 body_length,
124 };
125 let prepared = match request_metadata.presented_capability_override() {
126 Some(presented_capability) => evaluator.prepare_with_presented_capability(
127 evaluation_input,
128 Some(presented_capability),
129 ),
130 None => evaluator.prepare(evaluation_input),
131 };
132
133 let prepared = match prepared {
134 Ok(r) => r,
135 Err(e) => {
136 if evaluator.is_fail_open() {
137 crate::metrics::record_fail_open_suspected("tower");
138 tracing::warn!(
141 error = %e,
142 "Chio evaluation failed; fail-open enabled, forwarding request WITHOUT enforcement"
143 );
144 return inner.call(req).await.map_err(Into::into);
145 }
146 tracing::error!("Chio evaluation failed: {e}");
147 let mut response = http::Response::new(ResBody::default());
148 *response.status_mut() = http::StatusCode::BAD_GATEWAY;
149 return Ok(response);
150 }
151 };
152
153 if prepared.verdict.is_denied() {
154 let status = denied_status(&prepared.verdict);
155 let receipt = evaluator
156 .finalize_receipt(&prepared, status.as_u16())
157 .map_err(|error| Box::new(error) as Box<dyn std::error::Error + Send + Sync>)?;
158 evaluator
163 .persist_http_receipt(&receipt)
164 .map_err(|error| Box::new(error) as Box<dyn std::error::Error + Send + Sync>)?;
165
166 let mut response = http::Response::new(ResBody::default());
167 *response.status_mut() = status;
168 response.headers_mut().insert(
169 "x-chio-receipt-id",
170 http::HeaderValue::from_str(&receipt.id)
171 .unwrap_or_else(|_| http::HeaderValue::from_static("unknown")),
172 );
173 response.extensions_mut().insert(receipt);
174 return Ok(response);
175 }
176
177 if evaluator.has_durable_receipt_sink() {
186 let decision_receipt = evaluator
187 .sign_decision_receipt(&prepared)
188 .map_err(|error| Box::new(error) as Box<dyn std::error::Error + Send + Sync>)?;
189 evaluator
190 .persist_http_receipt(&decision_receipt)
191 .map_err(|error| Box::new(error) as Box<dyn std::error::Error + Send + Sync>)?;
192 }
193
194 let mut response = inner.call(req).await.map_err(Into::into)?;
196 let receipt = evaluator
197 .finalize_receipt(&prepared, response.status().as_u16())
198 .map_err(|error| Box::new(error) as Box<dyn std::error::Error + Send + Sync>)?;
199 evaluator
200 .persist_http_receipt(&receipt)
201 .map_err(|error| Box::new(error) as Box<dyn std::error::Error + Send + Sync>)?;
202
203 if let Ok(val) = http::HeaderValue::from_str(&receipt.id) {
205 response.headers_mut().insert("x-chio-receipt-id", val);
206 }
207 response.extensions_mut().insert(receipt);
208
209 Ok(response)
210 })
211 }
212}
213
214fn denied_status(verdict: &chio_http_core::Verdict) -> http::StatusCode {
215 if let chio_http_core::Verdict::Deny { http_status, .. } = verdict {
216 http::StatusCode::from_u16(*http_status).unwrap_or(http::StatusCode::FORBIDDEN)
217 } else {
218 http::StatusCode::FORBIDDEN
219 }
220}
221
222const TRANSPORT_BODY_SIZE_GUARD: &str = "chio_tower_request_body_limit_guard";
225
226fn build_payload_too_large_response<ResBody>(
238 evaluator: &ChioEvaluator,
239 method: &str,
240 path: &str,
241 caller: &chio_http_core::CallerIdentity,
242 max_body_bytes: usize,
243) -> http::Response<ResBody>
244where
245 ResBody: Default + From<Bytes>,
246{
247 use chio_http_core::TransportDenyInput;
248
249 if !evaluator.receipts_are_audited() {
256 tracing::error!(
257 target: "chio::tower",
258 "refusing to sign a transport-deny receipt without durable receipt storage; failing closed"
259 );
260 let mut response = http::Response::new(ResBody::default());
261 *response.status_mut() = http::StatusCode::BAD_GATEWAY;
262 return response;
263 }
264
265 let status = http::StatusCode::PAYLOAD_TOO_LARGE;
266 let caller_identity_hash = caller.identity_hash().ok();
267
268 let receipt = match (
269 crate::evaluator::parse_method(method),
270 caller_identity_hash.as_deref(),
271 ) {
272 (Ok(http_method), Some(caller_hash)) => {
273 let route_pattern = (evaluator.route_resolver())(method, path);
274 let verdict = chio_http_core::Verdict::deny_with_status(
275 format!("request body exceeds {max_body_bytes}-byte limit for chio_tower"),
276 TRANSPORT_BODY_SIZE_GUARD,
277 status.as_u16(),
278 );
279 let request_id = uuid::Uuid::now_v7().to_string();
280 evaluator
281 .sign_transport_deny_receipt(TransportDenyInput {
282 request_id: &request_id,
283 route_pattern: &route_pattern,
284 method: http_method,
285 caller_identity_hash: caller_hash,
286 content_hash: None,
287 verdict,
288 })
289 .ok()
290 }
291 _ => None,
292 };
293
294 let mut response = http::Response::new(ResBody::default());
295 *response.status_mut() = status;
296
297 if let Some(receipt) = receipt {
298 if let Err(error) = evaluator.persist_http_receipt(&receipt) {
302 tracing::error!(
303 target: "chio::tower",
304 %error,
305 "failed to persist transport-deny receipt to durable store; failing closed"
306 );
307 let mut response = http::Response::new(ResBody::default());
308 *response.status_mut() = http::StatusCode::BAD_GATEWAY;
309 return response;
310 }
311 if let Ok(val) = http::HeaderValue::from_str(&receipt.id) {
312 response.headers_mut().insert("x-chio-receipt-id", val);
313 }
314 if let Ok(body_bytes) = serde_json::to_vec(&receipt) {
315 *response.body_mut() = ResBody::from(Bytes::from(body_bytes));
316 response.headers_mut().insert(
317 http::header::CONTENT_TYPE,
318 http::HeaderValue::from_static("application/json"),
319 );
320 }
321 response.extensions_mut().insert(receipt);
322 }
323
324 response
325}
326
327enum BufferBodyError {
329 TooLarge,
333 Inner(Box<dyn std::error::Error + Send + Sync>),
335}
336
337async fn buffer_request_body<ReqBody>(
338 req: http::Request<ReqBody>,
339 max_body_bytes: usize,
340) -> Result<(http::Request<ReqBody>, Option<String>, u64), BufferBodyError>
341where
342 ReqBody: Body + From<Bytes>,
343 ReqBody::Data: Send,
344 ReqBody::Error: Into<Box<dyn std::error::Error + Send + Sync>>,
345{
346 if let Some(upper) = req.body().size_hint().upper() {
349 if upper > max_body_bytes as u64 {
350 return Err(BufferBodyError::TooLarge);
351 }
352 }
353
354 let (parts, body) = req.into_parts();
361 let mut body = std::pin::pin!(body);
362 let mut buffer = BytesMut::new();
363 let limit = max_body_bytes;
364
365 loop {
366 let frame = std::future::poll_fn(|cx| body.as_mut().poll_frame(cx)).await;
367 let frame = match frame {
368 Some(Ok(frame)) => frame,
369 Some(Err(error)) => return Err(BufferBodyError::Inner(error.into())),
370 None => break,
371 };
372 let mut data = match frame.into_data() {
373 Ok(data) => data,
374 Err(_non_data) => continue,
378 };
379 let chunk_len = data.remaining();
380 if buffer.len().saturating_add(chunk_len) > limit {
381 return Err(BufferBodyError::TooLarge);
382 }
383 while data.has_remaining() {
387 let slice = data.chunk();
388 let slice_len = slice.len();
389 buffer.extend_from_slice(slice);
390 data.advance(slice_len);
391 }
392 }
393
394 let collected = buffer.freeze();
395 let body_length = collected.len() as u64;
396 let body_hash = if collected.is_empty() {
397 None
398 } else {
399 let mut hasher = Sha256::new();
400 hasher.update(collected.as_ref());
401 Some(hex::encode(hasher.finalize()))
402 };
403 let replay = http::Request::from_parts(parts, ReqBody::from(collected));
404 Ok((replay, body_hash, body_length))
405}
406
407#[cfg(test)]
408mod tests {
409 use super::*;
410 use crate::evaluator::ChioEvaluator;
411 use bytes::Bytes;
412 use chio_core_types::capability::{
413 scope::ChioScope,
414 token::{CapabilityToken, CapabilityTokenBody},
415 };
416 use chio_core_types::crypto::Keypair;
417 use chio_http_core::{
418 http_authority_tool_grant, http_status_scope, HttpReceipt, CHIO_HTTP_STATUS_SCOPE_FINAL,
419 };
420 use http_body_util::{BodyExt, Full};
421 use tower::ServiceExt;
422
423 type TestBody = Full<Bytes>;
424
425 fn valid_capability_token_json(id: &str, issuer: &Keypair) -> String {
426 let now = chrono::Utc::now().timestamp() as u64;
427 let token = CapabilityToken::sign(
428 CapabilityTokenBody {
429 id: id.to_string(),
430 issuer: issuer.public_key(),
431 subject: issuer.public_key(),
432 scope: ChioScope {
433 grants: vec![http_authority_tool_grant()],
434 ..ChioScope::default()
435 },
436 issued_at: now.saturating_sub(60),
437 expires_at: now + 3600,
438 delegation_chain: Vec::new(),
439 aggregate_invocation_budget: None,
440 },
441 issuer,
442 )
443 .unwrap_or_else(|e| panic!("token sign failed: {e}"));
444 serde_json::to_string(&token).unwrap_or_else(|e| panic!("token serialize failed: {e}"))
445 }
446
447 fn make_service() -> (Keypair, ChioEvaluator) {
448 let keypair = Keypair::generate();
449 let evaluator = ChioEvaluator::new_ephemeral(keypair.clone(), "test-policy".to_string());
450 (keypair, evaluator)
451 }
452
453 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
454 async fn fail_open_branch_increments_suspected_counter() {
455 let (_kp, evaluator) = make_service();
456 let evaluator = evaluator.with_fail_open(true);
460 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
461 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
462 http::Response::new(Full::new(Bytes::new())),
463 )
464 });
465 let mut service = ChioService::new(inner, evaluator);
466
467 let before = {
468 let mut body = String::new();
469 chio_metrics_spec::runtime::families::FAIL_OPEN_SUSPECTED.render(&mut body);
470 body
471 };
472
473 let request = http::Request::builder()
474 .method("FOOBAR")
475 .uri("/anything")
476 .body(Full::new(Bytes::new()))
477 .unwrap_or_else(|e| panic!("request build failed: {e}"));
478 let response = service
479 .ready()
480 .await
481 .unwrap_or_else(|e| panic!("ready failed: {e}"))
482 .call(request)
483 .await;
484 assert!(response.is_ok(), "fail-open forwards to inner");
485
486 let mut after = String::new();
487 chio_metrics_spec::runtime::families::FAIL_OPEN_SUSPECTED.render(&mut after);
488 assert!(
491 after.contains("chio_fail_open_suspected_total{surface=\"tower\"}"),
492 "series must exist: {after}"
493 );
494 assert_ne!(before, after, "the fail-open counter must advance");
495 }
496
497 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
501 async fn service_allows_get() {
502 let (_kp, evaluator) = make_service();
503
504 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
505 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
506 http::Response::new(Full::new(Bytes::new())),
507 )
508 });
509
510 let mut service = ChioService::new(inner, evaluator);
511
512 let req = http::Request::builder()
513 .method("GET")
514 .uri("/pets")
515 .body(Full::new(Bytes::new()))
516 .unwrap_or_else(|e| panic!("build failed: {e}"));
517
518 let resp: http::Response<TestBody> = service
519 .ready()
520 .await
521 .unwrap_or_else(|e| panic!("ready failed: {e}"))
522 .call(req)
523 .await
524 .unwrap_or_else(|e| panic!("call failed: {e}"));
525
526 assert_eq!(resp.status(), http::StatusCode::OK);
527 assert!(resp.headers().contains_key("x-chio-receipt-id"));
528 let receipt = resp
529 .extensions()
530 .get::<HttpReceipt>()
531 .unwrap_or_else(|| panic!("missing receipt extension"));
532 assert_eq!(receipt.response_status, 200);
533 assert_eq!(
534 http_status_scope(receipt.metadata.as_ref()),
535 Some(CHIO_HTTP_STATUS_SCOPE_FINAL)
536 );
537 }
538
539 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
540 async fn service_denies_duplicate_query_capability_even_when_one_value_is_valid() {
541 let (kp, evaluator) = make_service();
542
543 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
544 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
545 http::Response::new(Full::new(Bytes::new())),
546 )
547 });
548
549 let mut service = ChioService::new(inner, evaluator);
550 let valid_capability = url::form_urlencoded::byte_serialize(
551 valid_capability_token_json("cap-query", &kp).as_bytes(),
552 )
553 .collect::<String>();
554 let req = http::Request::builder()
555 .method("GET")
556 .uri(format!(
557 "/pets?chio_capability=not-json&chio_capability={valid_capability}"
558 ))
559 .body(Full::new(Bytes::new()))
560 .unwrap_or_else(|e| panic!("build failed: {e}"));
561
562 let resp: http::Response<TestBody> = service
563 .ready()
564 .await
565 .unwrap_or_else(|e| panic!("ready failed: {e}"))
566 .call(req)
567 .await
568 .unwrap_or_else(|e| panic!("call failed: {e}"));
569
570 assert_eq!(resp.status(), http::StatusCode::FORBIDDEN);
571 assert!(resp.headers().contains_key("x-chio-receipt-id"));
572 let receipt = resp
573 .extensions()
574 .get::<HttpReceipt>()
575 .unwrap_or_else(|| panic!("missing receipt extension"));
576 assert!(receipt.is_denied());
577 assert_eq!(receipt.response_status, 403);
578 assert_eq!(
579 http_status_scope(receipt.metadata.as_ref()),
580 Some(CHIO_HTTP_STATUS_SCOPE_FINAL)
581 );
582 }
583
584 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
585 async fn service_denies_post_without_capability() {
586 let (_kp, evaluator) = make_service();
587
588 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
589 panic!("inner should not be called for denied requests");
590 #[allow(unreachable_code)]
591 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
592 http::Response::new(Full::new(Bytes::new())),
593 )
594 });
595
596 let mut service = ChioService::new(inner, evaluator);
597
598 let req = http::Request::builder()
599 .method("POST")
600 .uri("/pets")
601 .body(Full::new(Bytes::new()))
602 .unwrap_or_else(|e| panic!("build failed: {e}"));
603
604 let resp: http::Response<TestBody> = service
605 .ready()
606 .await
607 .unwrap_or_else(|e| panic!("ready failed: {e}"))
608 .call(req)
609 .await
610 .unwrap_or_else(|e| panic!("call failed: {e}"));
611
612 assert_eq!(resp.status(), http::StatusCode::FORBIDDEN);
613 assert!(resp.headers().contains_key("x-chio-receipt-id"));
614 let receipt = resp
615 .extensions()
616 .get::<HttpReceipt>()
617 .unwrap_or_else(|| panic!("missing receipt extension"));
618 assert_eq!(receipt.response_status, 403);
619 assert_eq!(
620 http_status_scope(receipt.metadata.as_ref()),
621 Some(CHIO_HTTP_STATUS_SCOPE_FINAL)
622 );
623 }
624
625 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
626 async fn service_allows_post_with_capability() {
627 let (kp, evaluator) = make_service();
628
629 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
630 let mut response = http::Response::new(Full::new(Bytes::new()));
631 *response.status_mut() = http::StatusCode::CREATED;
632 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(response)
633 });
634
635 let mut service = ChioService::new(inner, evaluator);
636
637 let req = http::Request::builder()
638 .method("POST")
639 .uri("/pets")
640 .header(
641 "x-chio-capability",
642 valid_capability_token_json("cap-service", &kp),
643 )
644 .body(Full::new(Bytes::from_static(br#"{"name":"Rex"}"#)))
645 .unwrap_or_else(|e| panic!("build failed: {e}"));
646
647 let resp: http::Response<TestBody> = service
648 .ready()
649 .await
650 .unwrap_or_else(|e| panic!("ready failed: {e}"))
651 .call(req)
652 .await
653 .unwrap_or_else(|e| panic!("call failed: {e}"));
654
655 assert_eq!(resp.status(), http::StatusCode::CREATED);
656 assert!(resp.headers().contains_key("x-chio-receipt-id"));
657 let receipt = resp
658 .extensions()
659 .get::<HttpReceipt>()
660 .unwrap_or_else(|| panic!("missing receipt extension"));
661 assert_eq!(receipt.response_status, 201);
662 assert_eq!(
663 http_status_scope(receipt.metadata.as_ref()),
664 Some(CHIO_HTTP_STATUS_SCOPE_FINAL)
665 );
666 }
667
668 #[tokio::test]
669 async fn buffer_request_body_hashes_and_replays_raw_bytes() {
670 let payload = Bytes::from_static(br#"{"hello":"world","count":2}"#);
671 let req = http::Request::builder()
672 .method("POST")
673 .uri("/echo")
674 .body(Full::new(payload.clone()))
675 .unwrap_or_else(|e| panic!("build failed: {e}"));
676
677 let (req, body_hash, body_length) =
678 match buffer_request_body(req, DEFAULT_MAX_BODY_BYTES).await {
679 Ok(parts) => parts,
680 Err(BufferBodyError::TooLarge) => panic!("body unexpectedly too large"),
681 Err(BufferBodyError::Inner(error)) => panic!("buffer failed: {error}"),
682 };
683
684 let replayed = req
685 .into_body()
686 .collect()
687 .await
688 .unwrap_or_else(|e| panic!("collect failed: {e}"))
689 .to_bytes();
690
691 let mut expected = Sha256::new();
692 expected.update(payload.as_ref());
693
694 assert_eq!(body_length, payload.len() as u64);
695 assert_eq!(body_hash, Some(hex::encode(expected.finalize())));
696 assert_eq!(replayed, payload);
697 }
698
699 #[tokio::test]
700 async fn buffer_request_body_rejects_oversized_payload() {
701 let payload = Bytes::from(vec![b'a'; 32]);
702 let req = http::Request::builder()
703 .method("POST")
704 .uri("/oversized")
705 .body(Full::new(payload))
706 .unwrap_or_else(|e| panic!("build failed: {e}"));
707
708 let result = buffer_request_body(req, 16).await;
709 assert!(matches!(result, Err(BufferBodyError::TooLarge)));
710 }
711
712 #[tokio::test]
717 async fn buffer_request_body_aborts_streaming_body_during_collection() {
718 use std::pin::Pin;
719 use std::task::{Context, Poll};
720
721 struct StreamingBody {
722 chunks: Vec<Bytes>,
723 }
724
725 impl Body for StreamingBody {
726 type Data = Bytes;
727 type Error = std::convert::Infallible;
728
729 fn poll_frame(
730 mut self: Pin<&mut Self>,
731 _cx: &mut Context<'_>,
732 ) -> Poll<Option<Result<http_body::Frame<Self::Data>, Self::Error>>> {
733 if self.chunks.is_empty() {
734 Poll::Ready(None)
735 } else {
736 let chunk = self.chunks.remove(0);
737 Poll::Ready(Some(Ok(http_body::Frame::data(chunk))))
738 }
739 }
740
741 fn size_hint(&self) -> http_body::SizeHint {
744 http_body::SizeHint::default()
745 }
746 }
747
748 struct AdaptedBody(StreamingBody);
750 impl Body for AdaptedBody {
751 type Data = Bytes;
752 type Error = std::convert::Infallible;
753 fn poll_frame(
754 mut self: Pin<&mut Self>,
755 cx: &mut Context<'_>,
756 ) -> Poll<Option<Result<http_body::Frame<Self::Data>, Self::Error>>> {
757 Pin::new(&mut self.0).poll_frame(cx)
758 }
759 fn size_hint(&self) -> http_body::SizeHint {
760 self.0.size_hint()
761 }
762 }
763 impl From<Bytes> for AdaptedBody {
764 fn from(_value: Bytes) -> Self {
765 AdaptedBody(StreamingBody { chunks: Vec::new() })
766 }
767 }
768
769 let chunks = vec![
770 Bytes::from(vec![b'x'; 8]),
771 Bytes::from(vec![b'x'; 8]),
772 Bytes::from(vec![b'x'; 8]),
773 Bytes::from(vec![b'x'; 8]),
774 ];
775 let req = http::Request::builder()
776 .method("POST")
777 .uri("/streaming-oversized")
778 .body(AdaptedBody(StreamingBody {
779 chunks: chunks.clone(),
780 }))
781 .unwrap_or_else(|e| panic!("build failed: {e}"));
782
783 let result = buffer_request_body(req, 16).await;
787 assert!(
788 matches!(result, Err(BufferBodyError::TooLarge)),
789 "streaming body exceeding the cap during collection must abort"
790 );
791 }
792
793 #[tokio::test]
794 async fn service_returns_413_when_body_exceeds_limit() {
795 let (_kp, evaluator) = make_service();
796
797 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
798 panic!("inner should not be called for oversized requests");
799 #[allow(unreachable_code)]
800 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
801 http::Response::new(Full::new(Bytes::new())),
802 )
803 });
804
805 let mut service = ChioService::new(inner, evaluator).with_max_body_bytes(16);
806
807 let oversized = Bytes::from(vec![b'x'; 64]);
808 let req = http::Request::builder()
809 .method("POST")
810 .uri("/pets")
811 .body(Full::new(oversized))
812 .unwrap_or_else(|e| panic!("build failed: {e}"));
813
814 let resp: http::Response<TestBody> = service
815 .ready()
816 .await
817 .unwrap_or_else(|e| panic!("call failed: {e}"))
818 .call(req)
819 .await
820 .unwrap_or_else(|e| panic!("call failed: {e}"));
821
822 assert_eq!(resp.status(), http::StatusCode::PAYLOAD_TOO_LARGE);
823 }
824
825 #[tokio::test]
830 async fn service_413_response_has_signed_receipt_body() {
831 let (_kp, evaluator) = make_service();
832
833 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
834 panic!("inner should not be called for oversized requests");
835 #[allow(unreachable_code)]
836 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
837 http::Response::new(Full::new(Bytes::new())),
838 )
839 });
840
841 let mut service = ChioService::new(inner, evaluator).with_max_body_bytes(16);
842
843 let oversized = Bytes::from(vec![b'x'; 64]);
844 let req = http::Request::builder()
845 .method("POST")
846 .uri("/pets")
847 .body(Full::new(oversized))
848 .unwrap_or_else(|e| panic!("build failed: {e}"));
849
850 let resp: http::Response<TestBody> = service
851 .ready()
852 .await
853 .unwrap_or_else(|e| panic!("ready failed: {e}"))
854 .call(req)
855 .await
856 .unwrap_or_else(|e| panic!("call failed: {e}"));
857
858 assert_eq!(resp.status(), http::StatusCode::PAYLOAD_TOO_LARGE);
859
860 let receipt = resp
863 .extensions()
864 .get::<HttpReceipt>()
865 .unwrap_or_else(|| panic!("missing receipt extension on 413"))
866 .clone();
867 assert_eq!(receipt.response_status, 413);
868 assert!(
869 receipt.is_denied(),
870 "413 receipt must record a Deny verdict"
871 );
872 assert!(
873 receipt
874 .verify_signature()
875 .unwrap_or_else(|e| panic!("verify failed: {e}")),
876 "413 receipt signature must verify under embedded kernel key"
877 );
878
879 let header_id = resp
882 .headers()
883 .get("x-chio-receipt-id")
884 .unwrap_or_else(|| panic!("missing x-chio-receipt-id header on 413"))
885 .to_str()
886 .unwrap_or_else(|e| panic!("header was not utf-8: {e}"))
887 .to_string();
888 assert_eq!(header_id, receipt.id);
889
890 assert_eq!(
892 resp.headers()
893 .get(http::header::CONTENT_TYPE)
894 .map(|v| v.to_str().unwrap_or_else(|e| panic!("ctype utf8: {e}"))),
895 Some("application/json"),
896 );
897 let body_bytes = resp
898 .into_body()
899 .collect()
900 .await
901 .unwrap_or_else(|e| panic!("collect failed: {e}"))
902 .to_bytes();
903 assert!(
904 !body_bytes.is_empty(),
905 "413 response body must not be empty"
906 );
907 let parsed: HttpReceipt = serde_json::from_slice(&body_bytes)
908 .unwrap_or_else(|e| panic!("response body must be a serialised HttpReceipt: {e}"));
909 assert_eq!(parsed.id, receipt.id);
910 assert_eq!(parsed.response_status, 413);
911 assert!(
912 parsed
913 .verify_signature()
914 .unwrap_or_else(|e| panic!("verify failed: {e}")),
915 "deserialised receipt body must verify under the embedded kernel key"
916 );
917 }
918
919 #[tokio::test]
925 async fn service_413_fails_closed_without_durable_receipts() {
926 let keypair = Keypair::generate();
927 let evaluator = ChioEvaluator::new(keypair, "fail-closed-policy".to_string());
930
931 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
932 panic!("inner should not be called for oversized requests");
933 #[allow(unreachable_code)]
934 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
935 http::Response::new(Full::new(Bytes::new())),
936 )
937 });
938
939 let mut service = ChioService::new(inner, evaluator).with_max_body_bytes(16);
940
941 let oversized = Bytes::from(vec![b'x'; 64]);
942 let req = http::Request::builder()
943 .method("POST")
944 .uri("/pets")
945 .body(Full::new(oversized))
946 .unwrap_or_else(|e| panic!("build failed: {e}"));
947
948 let resp: http::Response<TestBody> = service
949 .ready()
950 .await
951 .unwrap_or_else(|e| panic!("ready failed: {e}"))
952 .call(req)
953 .await
954 .unwrap_or_else(|e| panic!("call failed: {e}"));
955
956 assert_eq!(resp.status(), http::StatusCode::BAD_GATEWAY);
959 assert!(resp.extensions().get::<HttpReceipt>().is_none());
960 assert!(resp.headers().get("x-chio-receipt-id").is_none());
961 let body_bytes = resp
962 .into_body()
963 .collect()
964 .await
965 .unwrap_or_else(|e| panic!("collect failed: {e}"))
966 .to_bytes();
967 assert!(body_bytes.is_empty());
968 }
969
970 #[tokio::test]
971 async fn service_413_without_receipt_body_omits_json_content_type() {
972 let (_kp, evaluator) = make_service();
973
974 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
975 panic!("inner should not be called for oversized requests");
976 #[allow(unreachable_code)]
977 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
978 http::Response::new(Full::new(Bytes::new())),
979 )
980 });
981
982 let mut service = ChioService::new(inner, evaluator).with_max_body_bytes(16);
983
984 let oversized = Bytes::from(vec![b'x'; 64]);
985 let req = http::Request::builder()
986 .method("BREW")
987 .uri("/pets")
988 .body(Full::new(oversized))
989 .unwrap_or_else(|e| panic!("build failed: {e}"));
990
991 let resp: http::Response<TestBody> = service
992 .ready()
993 .await
994 .unwrap_or_else(|e| panic!("ready failed: {e}"))
995 .call(req)
996 .await
997 .unwrap_or_else(|e| panic!("call failed: {e}"));
998
999 assert_eq!(resp.status(), http::StatusCode::PAYLOAD_TOO_LARGE);
1000 assert!(resp.extensions().get::<HttpReceipt>().is_none());
1001 assert!(resp.headers().get(http::header::CONTENT_TYPE).is_none());
1002 assert!(resp.headers().get("x-chio-receipt-id").is_none());
1003
1004 let body_bytes = resp
1005 .into_body()
1006 .collect()
1007 .await
1008 .unwrap_or_else(|e| panic!("collect failed: {e}"))
1009 .to_bytes();
1010 assert!(body_bytes.is_empty());
1011 }
1012
1013 #[tokio::test]
1014 async fn service_413_unsupported_method_rejects_before_route_resolver() {
1015 fn route_resolver(_method: &str, _path: &str) -> String {
1016 panic!("route resolver must not run for unsupported HTTP methods");
1017 }
1018
1019 let (_kp, evaluator) = make_service();
1020 let evaluator = evaluator.with_route_resolver(route_resolver);
1021
1022 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
1023 panic!("inner should not be called for oversized requests");
1024 #[allow(unreachable_code)]
1025 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
1026 http::Response::new(Full::new(Bytes::new())),
1027 )
1028 });
1029
1030 let mut service = ChioService::new(inner, evaluator).with_max_body_bytes(16);
1031
1032 let oversized = Bytes::from(vec![b'x'; 64]);
1033 let req = http::Request::builder()
1034 .method("BREW")
1035 .uri("/pets")
1036 .body(Full::new(oversized))
1037 .unwrap_or_else(|e| panic!("build failed: {e}"));
1038
1039 let resp: http::Response<TestBody> = service
1040 .ready()
1041 .await
1042 .unwrap_or_else(|e| panic!("ready failed: {e}"))
1043 .call(req)
1044 .await
1045 .unwrap_or_else(|e| panic!("call failed: {e}"));
1046
1047 assert_eq!(resp.status(), http::StatusCode::PAYLOAD_TOO_LARGE);
1048 assert!(resp.extensions().get::<HttpReceipt>().is_none());
1049 assert!(resp.headers().get(http::header::CONTENT_TYPE).is_none());
1050 assert!(resp.headers().get("x-chio-receipt-id").is_none());
1051 }
1052
1053 #[tokio::test]
1054 async fn service_413_receipt_uses_route_resolver_pattern() {
1055 fn route_resolver(_method: &str, path: &str) -> String {
1056 if path.starts_with("/pets/") {
1057 "/pets/{petId}".to_string()
1058 } else {
1059 path.to_string()
1060 }
1061 }
1062
1063 let (_kp, evaluator) = make_service();
1064 let evaluator = evaluator.with_route_resolver(route_resolver);
1065
1066 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
1067 panic!("inner should not be called for oversized requests");
1068 #[allow(unreachable_code)]
1069 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
1070 http::Response::new(Full::new(Bytes::new())),
1071 )
1072 });
1073
1074 let mut service = ChioService::new(inner, evaluator).with_max_body_bytes(16);
1075
1076 let oversized = Bytes::from(vec![b'x'; 64]);
1077 let req = http::Request::builder()
1078 .method("POST")
1079 .uri("/pets/secret-id")
1080 .body(Full::new(oversized))
1081 .unwrap_or_else(|e| panic!("build failed: {e}"));
1082
1083 let resp: http::Response<TestBody> = service
1084 .ready()
1085 .await
1086 .unwrap_or_else(|e| panic!("ready failed: {e}"))
1087 .call(req)
1088 .await
1089 .unwrap_or_else(|e| panic!("call failed: {e}"));
1090
1091 let receipt = resp
1092 .extensions()
1093 .get::<HttpReceipt>()
1094 .unwrap_or_else(|| panic!("missing receipt extension on 413"));
1095 assert_eq!(receipt.route_pattern, "/pets/{petId}");
1096 }
1097
1098 #[derive(Default)]
1102 struct RecordingReceiptStore {
1103 receipts: std::sync::Mutex<Vec<chio_core_types::receipt::body::ChioReceipt>>,
1104 }
1105
1106 impl RecordingReceiptStore {
1107 fn stored_ids(&self) -> Vec<String> {
1108 let receipts = match self.receipts.lock() {
1109 Ok(receipts) => receipts,
1110 Err(poisoned) => poisoned.into_inner(),
1111 };
1112 receipts.iter().map(|receipt| receipt.id.clone()).collect()
1113 }
1114 }
1115
1116 impl chio_kernel::ReceiptStore for RecordingReceiptStore {
1117 fn append_chio_receipt(
1118 &self,
1119 receipt: &chio_core_types::receipt::body::ChioReceipt,
1120 ) -> Result<(), chio_kernel::ReceiptStoreError> {
1121 match self.receipts.lock() {
1122 Ok(mut receipts) => receipts.push(receipt.clone()),
1123 Err(poisoned) => poisoned.into_inner().push(receipt.clone()),
1124 }
1125 Ok(())
1126 }
1127
1128 fn append_child_receipt(
1129 &self,
1130 _receipt: &chio_core_types::receipt::lineage::ChildRequestReceipt,
1131 ) -> Result<(), chio_kernel::ReceiptStoreError> {
1132 Ok(())
1133 }
1134 }
1135
1136 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1141 async fn service_persists_http_receipt_to_durable_store() {
1142 let keypair = Keypair::generate();
1143 let store = std::sync::Arc::new(RecordingReceiptStore::default());
1144 let evaluator = ChioEvaluator::builder(keypair.clone(), "durable-policy".to_string())
1145 .receipt_store(store.clone())
1146 .allow_ephemeral(true)
1147 .build()
1148 .unwrap_or_else(|e| panic!("build failed: {e}"));
1149
1150 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
1151 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
1152 http::Response::new(Full::new(Bytes::new())),
1153 )
1154 });
1155 let mut service = ChioService::new(inner, evaluator);
1156
1157 let req = http::Request::builder()
1158 .method("GET")
1159 .uri("/pets")
1160 .body(Full::new(Bytes::new()))
1161 .unwrap_or_else(|e| panic!("build failed: {e}"));
1162
1163 let resp: http::Response<TestBody> = service
1164 .ready()
1165 .await
1166 .unwrap_or_else(|e| panic!("ready failed: {e}"))
1167 .call(req)
1168 .await
1169 .unwrap_or_else(|e| panic!("call failed: {e}"));
1170
1171 assert_eq!(resp.status(), http::StatusCode::OK);
1172 let http_receipt = resp
1173 .extensions()
1174 .get::<HttpReceipt>()
1175 .unwrap_or_else(|| panic!("missing receipt extension"))
1176 .clone();
1177
1178 let expected = http_receipt
1182 .to_chio_receipt_with_keypair(&keypair)
1183 .unwrap_or_else(|e| panic!("convert failed: {e}"));
1184 let stored_ids = store.stored_ids();
1185 assert!(
1186 stored_ids.contains(&expected.id),
1187 "durable store must contain the HTTP decision receipt {}; stored: {stored_ids:?}",
1188 expected.id
1189 );
1190 }
1191
1192 #[derive(Default)]
1198 struct FailAfterFirstAppend {
1199 appended: std::sync::atomic::AtomicUsize,
1200 }
1201
1202 impl chio_kernel::ReceiptStore for FailAfterFirstAppend {
1203 fn append_chio_receipt(
1204 &self,
1205 _receipt: &chio_core_types::receipt::body::ChioReceipt,
1206 ) -> Result<(), chio_kernel::ReceiptStoreError> {
1207 if self
1208 .appended
1209 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
1210 == 0
1211 {
1212 Ok(())
1213 } else {
1214 Err(chio_kernel::ReceiptStoreError::Conflict(
1215 "durable receipt append failed (volume full)".to_string(),
1216 ))
1217 }
1218 }
1219
1220 fn append_child_receipt(
1221 &self,
1222 _receipt: &chio_core_types::receipt::lineage::ChildRequestReceipt,
1223 ) -> Result<(), chio_kernel::ReceiptStoreError> {
1224 Ok(())
1225 }
1226 }
1227
1228 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1234 async fn service_fails_closed_when_durable_http_receipt_append_fails() {
1235 let keypair = Keypair::generate();
1236 let store = std::sync::Arc::new(FailAfterFirstAppend::default());
1237 let evaluator = ChioEvaluator::builder(keypair, "durable-policy".to_string())
1238 .receipt_store(store)
1239 .allow_ephemeral(true)
1240 .build()
1241 .unwrap_or_else(|e| panic!("build failed: {e}"));
1242
1243 let inner = tower::service_fn(|_req: http::Request<TestBody>| async {
1244 panic!("inner service must not run when the durable HTTP receipt append fails");
1245 #[allow(unreachable_code)]
1246 Ok::<http::Response<TestBody>, Box<dyn std::error::Error + Send + Sync>>(
1247 http::Response::new(Full::new(Bytes::new())),
1248 )
1249 });
1250 let mut service = ChioService::new(inner, evaluator);
1251
1252 let req = http::Request::builder()
1253 .method("GET")
1254 .uri("/pets")
1255 .body(Full::new(Bytes::new()))
1256 .unwrap_or_else(|e| panic!("build failed: {e}"));
1257
1258 let result = service
1259 .ready()
1260 .await
1261 .unwrap_or_else(|e| panic!("ready failed: {e}"))
1262 .call(req)
1263 .await;
1264
1265 if let Ok(response) = result {
1269 assert!(
1270 response.status().is_server_error(),
1271 "a durable HTTP receipt append failure must fail closed, got {}",
1272 response.status()
1273 );
1274 }
1275 }
1276}