praxis_core/subrequest/
body.rs1use std::time::Duration;
5
6use bytes::Bytes;
7use metrics::{counter, histogram};
8use pingora_core::{connectors::http::Connector, protocols::http::client::HttpSession, upstreams::peer::HttpPeer};
9use tracing::{debug, warn};
10
11use super::{
12 internals::{
13 SUBREQUEST_STREAM_BYTES_TOTAL, SUBREQUEST_STREAM_DURATION_SECONDS, SUBREQUEST_STREAMS_TOTAL,
14 check_clean_completion,
15 },
16 types::{SubRequestError, SubResponseBody},
17};
18
19pub(super) async fn dispose_session_abnormal(
25 session: HttpSession<()>,
26 peer: Option<&HttpPeer>,
27 connector: Option<&Connector>,
28) {
29 let mut session = session;
30 session.shutdown().await;
31 if matches!(session, HttpSession::H2(_))
32 && let (Some(peer), Some(connector)) = (peer, connector)
33 {
34 connector.release_http_session(session, peer, None).await;
35 }
36}
37
38impl SubResponseBody {
43 pub(super) fn new_done() -> Self {
45 Self {
46 session: None,
47 peer: None,
48 connector: None,
49 permit: None,
50 read_timeout: None,
51 idle_timeout: Duration::from_secs(30),
52 stream_deadline: None,
53 max_total_bytes: None,
54 received_bytes: 0,
55 chunk_count: 0,
56 stream_started_at: tokio::time::Instant::now(),
57 done: true,
58 }
59 }
60
61 pub fn is_done(&self) -> bool {
63 self.done
64 }
65
66 pub fn received_bytes(&self) -> usize {
68 self.received_bytes
69 }
70
71 pub fn chunk_count(&self) -> u64 {
73 self.chunk_count
74 }
75
76 #[expect(clippy::large_stack_frames, reason = "Pingora session types are large")]
94 #[expect(clippy::too_many_lines, reason = "inline chunk read and deadline enforcement")]
95 #[expect(
96 clippy::expect_used,
97 reason = "session is an invariant: None only when done=true; caller checked !done"
98 )]
99 pub async fn next_chunk(&mut self) -> Result<Option<Bytes>, SubRequestError> {
100 if self.done {
101 return Ok(None);
102 }
103
104 if let Some(deadline) = self.stream_deadline
106 && tokio::time::Instant::now() >= deadline
107 {
108 self.shutdown_and_done("deadline_exceeded").await;
109 return Err(SubRequestError::DeadlineExceeded);
110 }
111
112 let session = self.session.as_mut().expect("session must be present when not done");
113
114 let mut effective_timeout = self.idle_timeout;
116 if let Some(rt) = self.read_timeout {
117 effective_timeout = effective_timeout.min(rt);
118 }
119 if let Some(deadline) = self.stream_deadline {
120 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
121 if remaining.is_zero() {
122 self.shutdown_and_done("deadline_exceeded").await;
123 return Err(SubRequestError::DeadlineExceeded);
124 }
125 effective_timeout = effective_timeout.min(remaining);
126 }
127
128 let read_result = tokio::time::timeout(effective_timeout, session.read_response_body()).await;
130
131 match read_result {
132 Ok(Ok(Some(chunk))) => {
133 self.received_bytes += chunk.len();
134 self.chunk_count += 1;
135
136 if let Some(limit) = self.max_total_bytes
138 && self.received_bytes > limit
139 {
140 self.shutdown_and_done("byte_limit").await;
141 return Err(SubRequestError::ResponseTooLarge {
142 actual: self.received_bytes,
143 limit,
144 });
145 }
146
147 match check_clean_completion(self.session.as_mut().expect("session present")) {
149 Ok(true) => self.release_session().await,
150 Ok(false) => {},
151 Err(e) => {
152 self.shutdown_and_done("h2_error").await;
153 return Err(e);
154 },
155 }
156
157 Ok(Some(chunk))
158 },
159 Ok(Ok(None)) => {
160 match check_clean_completion(self.session.as_mut().expect("session present")) {
162 Ok(true) => {
163 self.release_session().await;
164 Ok(None)
165 },
166 Ok(false) => {
167 self.shutdown_and_done("io_error").await;
168 Err(SubRequestError::Io("upstream closed without clean EOF".to_owned()))
169 },
170 Err(e) => {
171 self.shutdown_and_done("io_error").await;
172 Err(e)
173 },
174 }
175 },
176 Ok(Err(e)) => {
177 self.shutdown_and_done("io_error").await;
178 Err(SubRequestError::Io(e.to_string()))
179 },
180 Err(_elapsed) => {
181 if let Some(deadline) = self.stream_deadline
183 && tokio::time::Instant::now() >= deadline
184 {
185 self.shutdown_and_done("deadline_exceeded").await;
186 return Err(SubRequestError::DeadlineExceeded);
187 }
188 if let Some(rt) = self.read_timeout
189 && rt < self.idle_timeout
190 {
191 self.shutdown_and_done("read_timeout").await;
192 return Err(SubRequestError::Io("upstream read timeout".to_owned()));
193 }
194 let idle_timeout = self.idle_timeout;
195 self.shutdown_and_done("idle_timeout").await;
196 Err(SubRequestError::StreamIdleTimeout { idle_timeout })
197 },
198 }
199 }
200
201 async fn shutdown_and_done(&mut self, termination: &str) {
208 debug!(
209 termination,
210 duration_s = self.stream_started_at.elapsed().as_secs_f64(),
211 bytes = self.received_bytes,
212 chunks = self.chunk_count,
213 "sub-request: stream terminated"
214 );
215 self.record_stream_metrics(termination);
216 self.done = true;
217 if let Some(session) = self.session.take() {
218 let permit = self.permit.take();
219 let peer_ref = self.peer.as_ref();
220 let conn_ref = self
221 .connector
222 .as_ref()
223 .map(super::internals::SubRequestConnector::connector);
224 Box::pin(async move {
225 dispose_session_abnormal(session, peer_ref, conn_ref).await;
226 drop(permit);
227 })
228 .await;
229 }
230 self.peer.take();
231 self.connector.take();
232 }
233
234 async fn release_session(&mut self) {
241 debug!(
242 duration_s = self.stream_started_at.elapsed().as_secs_f64(),
243 bytes = self.received_bytes,
244 chunks = self.chunk_count,
245 "sub-request: stream completed"
246 );
247 self.record_stream_metrics("eof");
248 self.done = true;
249 if let (Some(session), Some(peer), Some(connector)) =
250 (self.session.take(), self.peer.as_ref(), self.connector.as_ref())
251 {
252 let permit = self.permit.take();
253 Box::pin(async move {
254 connector.connector().release_http_session(session, peer, None).await;
255 drop(permit);
256 })
257 .await;
258 }
259 self.peer.take();
260 self.connector.take();
261 }
262
263 pub async fn cancel(mut self) {
269 if !self.done {
270 self.shutdown_and_done("cancel").await;
271 }
272 }
273
274 fn record_stream_metrics(&self, termination: &str) {
276 let elapsed = self.stream_started_at.elapsed().as_secs_f64();
277 counter!(
278 SUBREQUEST_STREAMS_TOTAL,
279 "termination" => termination.to_owned(),
280 )
281 .increment(1);
282 histogram!(SUBREQUEST_STREAM_DURATION_SECONDS).record(elapsed);
283 counter!(SUBREQUEST_STREAM_BYTES_TOTAL).increment(self.received_bytes as u64);
284 }
285}
286
287impl Drop for SubResponseBody {
288 fn drop(&mut self) {
289 if self.done || self.session.is_none() {
290 return;
291 }
292 warn!(
293 duration_s = self.stream_started_at.elapsed().as_secs_f64(),
294 bytes = self.received_bytes,
295 chunks = self.chunk_count,
296 "sub-request: stream dropped without cancel"
297 );
298 self.record_stream_metrics("drop");
299 let Some(session) = self.session.take() else {
300 return;
301 };
302 let peer = self.peer.take();
303 let connector = self.connector.take();
304 let permit = self.permit.take();
305 if let Ok(handle) = tokio::runtime::Handle::try_current() {
306 handle.spawn(async move {
307 Box::pin(dispose_session_abnormal(
308 session,
309 peer.as_ref(),
310 connector.as_ref().map(super::internals::SubRequestConnector::connector),
311 ))
312 .await;
313 drop(permit);
314 });
315 }
316 }
317}