Skip to main content

praxis_core/subrequest/
body.rs

1// SPDX-License-Identifier: MIT
2// Copyright (c) 2026 Praxis Contributors
3
4use 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
19/// Dispose of a session after abnormal termination (cancel, timeout,
20/// error). H2 streams are released through the connector so the
21/// multiplexed connection survives. H1/Custom sessions are shut down
22/// without pooling because unread response bytes would corrupt the
23/// next response on a reused connection.
24pub(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
38// ---------------------------------------------------------------------------
39// SubResponseBody
40// ---------------------------------------------------------------------------
41
42impl SubResponseBody {
43    /// Create a body in the completed state (for header-time completion).
44    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    /// Whether this body has completed (EOF, error, or cancel).
62    pub fn is_done(&self) -> bool {
63        self.done
64    }
65
66    /// Total bytes received so far.
67    pub fn received_bytes(&self) -> usize {
68        self.received_bytes
69    }
70
71    /// Number of chunks received so far.
72    pub fn chunk_count(&self) -> u64 {
73        self.chunk_count
74    }
75
76    /// Pull the next body chunk from the upstream.
77    ///
78    /// Returns:
79    /// - `Ok(Some(chunk))` — a data chunk.
80    /// - `Ok(None)` — clean EOF; the session has been released to the pool.
81    ///
82    /// # Errors
83    ///
84    /// - `Err(StreamIdleTimeout)` — upstream stalled for `idle_timeout`.
85    /// - `Err(DeadlineExceeded)` — `max_stream_duration` expired.
86    /// - `Err(ResponseTooLarge)` — cumulative bytes exceeded `max_total_bytes`.
87    /// - `Err(Io)` — transport error or unclean EOF.
88    ///
89    /// # Panics
90    ///
91    /// Panics if the session is `None` when `done` is `false` (internal
92    /// invariant violation).
93    #[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        // Check stream deadline before reading.
105        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        // Compute effective timeout: min(read_timeout, idle_timeout, remaining stream deadline).
115        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        // Read one chunk with timeout.
129        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                // Check byte limit.
137                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                // Check for immediate completion after this chunk.
148                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                // No more data. Check clean completion.
161                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                // Distinguish stream deadline, read timeout, and idle timeout.
182                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    /// Shut down the session and mark the body as done.
202    ///
203    /// `done` is set before the first `.await`, so cancellation of
204    /// the cleanup future leaves the body in a valid terminal state.
205    /// The permit is moved into the boxed future so cancellation
206    /// releases it immediately.
207    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    /// Release the session to the connection pool and mark done.
235    ///
236    /// `done` is set before the first `.await`, so cancellation of
237    /// the release future leaves the body in a valid terminal state.
238    /// The permit is moved into the boxed future so cancellation
239    /// releases it immediately.
240    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    /// Explicitly cancel the streaming response.
264    ///
265    /// Shuts down the upstream session and releases the admission
266    /// permit. Consumes `self`. No-op if the body is already done
267    /// (EOF, error, or prior cancel).
268    pub async fn cancel(mut self) {
269        if !self.done {
270            self.shutdown_and_done("cancel").await;
271        }
272    }
273
274    /// Record stream termination metrics.
275    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}