Skip to main content

ntex/web/types/
mod.rs

1//! Extractor types
2
3pub(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
18/// Reads the complete body, failing with `overflow(size)` once `size > limit`.
19///
20/// A single-chunk body is copied only if it is a part of a larger buffer,
21/// which it would otherwise keep alive.
22async 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        // the body does not keep the larger buffer it is a part of
108        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        // nor the unused part of the buffer it is read into
117        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}