Skip to main content

rama_http_types/body/
infinite.rs

1use rama_core::stream::io::ReaderStream;
2use rama_core::telemetry::tracing;
3use rama_utils::macros::generate_set_and_with;
4use rand::{Rng, RngExt as _, rng};
5use std::{
6    io,
7    pin::Pin,
8    task::{Context, Poll},
9    time::Duration,
10};
11use tokio::{
12    io::{AsyncRead, ReadBuf},
13    time::Sleep,
14};
15
16/// A(n) (in)finite random byte stream implementing [`AsyncRead`].
17///
18/// Cheap to serve, expensive to download.
19/// Eat while it's hot my dear little bots.
20pub struct InfiniteReader {
21    chunk_size: usize,
22    limit: Option<usize>,
23    byte_count: usize,
24    max_delay: Option<Duration>,
25    sleep: Option<Pin<Box<Sleep>>>,
26}
27
28impl Default for InfiniteReader {
29    #[inline]
30    fn default() -> Self {
31        Self::new()
32    }
33}
34
35impl InfiniteReader {
36    /// Create an new default [`InfiniteReader`].
37    #[must_use]
38    pub fn new() -> Self {
39        Self {
40            chunk_size: 4096,
41            limit: None,
42            byte_count: 0,
43            max_delay: None,
44            sleep: None,
45        }
46    }
47
48    generate_set_and_with! {
49        /// Define the max throttle to be used for the intervals.
50        ///
51        /// Setting it will ensure that we have randomised throttles between
52        /// reads, making it more effective.
53        pub fn throttle(mut self, delay: Option<Duration>) -> Self {
54            self.max_delay = delay.and_then(|d| (!d.is_zero()).then_some(d));
55            self
56        }
57    }
58
59    generate_set_and_with! {
60        /// Set a limit on how much data will be served,
61        /// by default it will be an infinite amount of data.
62        pub fn size_limit(mut self, limit: Option<usize>) -> Self {
63            self.limit = limit.and_then(|n| (n > 0).then_some(n));
64            self
65        }
66    }
67
68    generate_set_and_with! {
69        /// Define the chunk size for downloads.
70        ///
71        /// The default value is used if a value of 0 is given.
72        pub fn chunk_size(mut self, size: usize) -> Self {
73            self.chunk_size = if size == 0 {
74                4096
75            } else {
76                size
77            };
78            self
79        }
80    }
81
82    /// Turn this [`InfiniteReader`] into a [`Body`].
83    ///
84    /// [`Body`]: super::Body
85    pub fn into_body(self) -> super::Body {
86        let stream = ReaderStream::new(self);
87        super::Body::from_stream(stream)
88    }
89}
90
91impl AsyncRead for InfiniteReader {
92    fn poll_read(
93        mut self: Pin<&mut Self>,
94        cx: &mut Context<'_>,
95        buf: &mut ReadBuf<'_>,
96    ) -> Poll<io::Result<()>> {
97        if self.limit.map(|n| n <= self.byte_count).unwrap_or_default() {
98            tracing::trace!(
99                "InfiniteReader finished reading, reached limit ({:?}): {}",
100                self.limit,
101                self.byte_count,
102            );
103            return Poll::Ready(Ok(()));
104        }
105
106        let mut rng = rng();
107
108        if let Some(max_delay) = self.max_delay {
109            if let Some(sleep) = self.sleep.as_mut() {
110                match sleep.as_mut().poll(cx) {
111                    Poll::Ready(_) => {
112                        tracing::trace!(
113                            "InfiniteReader throttle finished (limit: {}ms)",
114                            max_delay.as_millis()
115                        );
116                        self.sleep = None;
117                    }
118                    Poll::Pending => {
119                        tracing::trace!("InfiniteReader still throttling...");
120                        return Poll::Pending;
121                    }
122                }
123            } else {
124                let max_ms = max_delay.as_millis() as u64;
125                let rand_ms = rng.random_range(0..=max_ms);
126                let delay = Duration::from_millis(rand_ms);
127                tracing::trace!("InfiniteReader start throttle: {rand_ms}ms",);
128                let mut sleep = Box::pin(tokio::time::sleep(delay));
129                if sleep.as_mut().poll(cx).is_pending() {
130                    self.sleep = Some(sleep);
131                    return Poll::Pending;
132                };
133            }
134        }
135
136        let len = self.chunk_size.min(buf.remaining());
137        self.byte_count += len;
138        tracing::trace!("InfiniteReader feeding data: {len} random byte(s)");
139        let mut data = vec![0u8; len];
140        rng.fill_bytes(&mut data);
141        buf.put_slice(&data);
142        Poll::Ready(Ok(()))
143    }
144}
145
146impl From<InfiniteReader> for super::Body {
147    #[inline]
148    fn from(reader: InfiniteReader) -> Self {
149        reader.into_body()
150    }
151}