Skip to main content

ntex/http/h1/
payload.rs

1use std::task::{Context, Poll};
2use std::{cell::RefCell, ops::Deref, pin::Pin, rc::Rc};
3
4use crate::channel::bstream::{self, Receiver};
5use crate::http::{HeaderMap, error::PayloadError};
6use crate::util::{Bytes, Stream};
7
8/// A buffered stream of an HTTP/1 request body's decoded bytes.
9///
10/// Each item is either a body chunk or a [`PayloadError`].
11/// Normal body completion closes the stream, after which receiving returns
12/// `None`. A payload error is yielded once before the stream terminates.
13///
14/// The HTTP/1 dispatcher stops reading body data while this stream's buffer is
15/// full. Consuming items therefore releases transport-level backpressure.
16/// Dropping the stream before the complete body has been decoded prevents the
17/// connection from being reused and causes the dispatcher to disconnect it.
18///
19/// Dereferences to the underlying [`bstream::Receiver`].
20#[derive(Debug)]
21pub struct Payload {
22    rx: Receiver<PayloadError>,
23    trailers: Rc<RefCell<Option<HeaderMap>>>,
24}
25
26impl Payload {
27    /// Creates a payload stream and its sender.
28    pub(crate) fn create() -> (PayloadSender, Payload) {
29        let (tx, rx) = bstream::channel();
30        let trailers = Rc::new(RefCell::new(None));
31        (
32            PayloadSender {
33                tx,
34                trailers: trailers.clone(),
35            },
36            Payload { rx, trailers },
37        )
38    }
39
40    /// Returns the trailer fields of a chunked payload.
41    ///
42    /// Trailers are available once the payload is complete, `None` is
43    /// returned if the payload is not complete or has no trailers.
44    pub fn trailers(&self) -> Option<HeaderMap> {
45        self.trailers.borrow().clone()
46    }
47}
48
49impl From<Receiver<PayloadError>> for Payload {
50    fn from(rx: Receiver<PayloadError>) -> Self {
51        Payload {
52            rx,
53            trailers: Rc::default(),
54        }
55    }
56}
57
58impl Deref for Payload {
59    type Target = Receiver<PayloadError>;
60
61    fn deref(&self) -> &Self::Target {
62        &self.rx
63    }
64}
65
66impl Stream for Payload {
67    type Item = Result<Bytes, PayloadError>;
68
69    fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
70        self.rx.poll_read(cx)
71    }
72}
73
74/// Sender side of the HTTP/1 payload stream.
75#[derive(Debug)]
76pub(crate) struct PayloadSender {
77    tx: bstream::Sender<PayloadError>,
78    trailers: Rc<RefCell<Option<HeaderMap>>>,
79}
80
81impl PayloadSender {
82    /// Stores the trailer fields, the payload is completed with `feed_eof()`.
83    pub(crate) fn feed_trailers(&self, trailers: HeaderMap) {
84        *self.trailers.borrow_mut() = Some(trailers);
85    }
86}
87
88impl Deref for PayloadSender {
89    type Target = bstream::Sender<PayloadError>;
90
91    fn deref(&self) -> &Self::Target {
92        &self.tx
93    }
94}