1pub(in crate::web) mod form;
4pub(in crate::web) mod json;
5mod path;
6pub(in crate::web) mod payload;
7mod query;
8
9pub use self::form::{Form, FormConfig};
10pub use self::json::{Json, JsonConfig};
11pub use self::path::Path;
12pub use self::payload::{Payload, PayloadConfig};
13pub use self::query::Query;
14
15use crate::http::{error::PayloadError, helpers::take_trimmed};
16use crate::util::{Bytes, BytesMut, Stream, stream_recv};
17
18async fn read_body<S, E>(
23 stream: &mut S,
24 limit: usize,
25 length: Option<usize>,
26 overflow: impl FnOnce(usize) -> E,
27) -> Result<Bytes, E>
28where
29 S: Stream<Item = Result<Bytes, PayloadError>> + Unpin,
30 E: From<PayloadError>,
31{
32 let mut first = match stream_recv(stream).await {
33 Some(item) => item?,
34 None => return Ok(Bytes::new()),
35 };
36 if first.len() > limit {
37 return Err(overflow(first.len()));
38 }
39 let Some(chunk) = stream_recv(stream).await else {
40 first.trimdown();
41 return Ok(first);
42 };
43 let chunk = chunk?;
44 let size = first.len() + chunk.len();
45 if size > limit {
46 return Err(overflow(size));
47 }
48
49 let mut body = BytesMut::with_capacity(length.unwrap_or(0).clamp(size, limit));
50 body.extend_from_slice(&first);
51 body.extend_from_slice(&chunk);
52 while let Some(item) = stream_recv(stream).await {
53 let chunk = item?;
54 let size = body.len() + chunk.len();
55 if size > limit {
56 return Err(overflow(size));
57 }
58 body.extend_from_slice(&chunk);
59 }
60 Ok(take_trimmed(&mut body))
61}
62
63#[cfg(test)]
64mod tests {
65 use std::{collections::VecDeque, pin::Pin, task::Context, task::Poll};
66
67 use super::*;
68
69 struct Chunks(VecDeque<Result<Bytes, PayloadError>>);
70
71 impl Stream for Chunks {
72 type Item = Result<Bytes, PayloadError>;
73
74 fn poll_next(mut self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<Option<Self::Item>> {
75 Poll::Ready(self.0.pop_front())
76 }
77 }
78
79 fn chunks(items: &[&'static [u8]]) -> Chunks {
80 Chunks(items.iter().map(|b| Ok(Bytes::from_static(b))).collect())
81 }
82
83 #[crate::rt_test]
84 async fn test_read_body() {
85 let data: &'static [u8] = b"single chunk";
86 let body = read_body(&mut chunks(&[data]), 64, None, |_| PayloadError::Overflow)
87 .await
88 .unwrap();
89 assert_eq!(
90 body.as_ptr(),
91 data.as_ptr(),
92 "single chunk must not be copied"
93 );
94
95 let body = read_body(&mut chunks(&[]), 64, None, |_| PayloadError::Overflow)
96 .await
97 .unwrap();
98 assert!(body.is_empty());
99
100 let body = read_body(&mut chunks(&[b"ab", b"cd", b"ef"]), 6, Some(6), |_| {
101 PayloadError::Overflow
102 })
103 .await
104 .unwrap();
105 assert_eq!(body, Bytes::from_static(b"abcdef"));
106
107 let page = Bytes::from(vec![b'x'; 32 * 1024]);
109 let mut stream = Chunks([Ok(page.slice(..1024))].into());
110 let body = read_body(&mut stream, 64 * 1024, None, |_| PayloadError::Overflow)
111 .await
112 .unwrap();
113 assert_eq!(body, page.slice(..1024));
114 assert_ne!(body.as_ptr(), page.as_ptr());
115
116 let chunk = Bytes::from(vec![b'x'; 1000]);
118 let mut stream = Chunks(std::iter::repeat_n(Ok(chunk), 33).collect());
119 let body = read_body(&mut stream, 64 * 1024, None, |_| PayloadError::Overflow)
120 .await
121 .unwrap();
122 assert_eq!(body.len(), 33 * 1000);
123 let mut copy = body.clone();
124 copy.trimdown();
125 assert_eq!(copy.as_ptr(), body.as_ptr(), "the body has no unused space");
126
127 for items in [
128 &[&b"abcdefg"[..]][..],
129 &[b"abcd", b"efg"],
130 &[b"ab", b"cd", b"efg"],
131 ] {
132 let mut size = 0;
133 let res = read_body(&mut chunks(items), 6, None, |s| {
134 size = s;
135 PayloadError::Overflow
136 })
137 .await;
138 assert!(matches!(res, Err(PayloadError::Overflow)));
139 assert_eq!(size, 7);
140 }
141 }
142}