rama_http_types/body/
infinite.rs1use 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
16pub 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 #[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 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 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 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 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}