use std::panic;
use futures::future;
use futures::stream;
use futures::stream::Stream;
use futures::stream::TryStreamExt;
use std::task::Poll;
use bytes::Bytes;
use crate::{error, result};
use crate::solicit::header::Headers;
use crate::data_or_headers::DataOrHeaders;
use crate::data_or_headers_with_flag::DataOrHeadersWithFlag;
use crate::data_or_headers_with_flag::DataOrHeadersWithFlagStream;
use crate::misc::any_to_string;
use crate::solicit::end_stream::EndStream;
use futures::stream::StreamExt;
use futures::task::Context;
use std::pin::Pin;
pub enum DataOrTrailers {
Data(Bytes, EndStream),
Trailers(Headers),
}
impl DataOrTrailers {
pub fn intermediate_data(data: Bytes) -> Self {
DataOrTrailers::Data(data, EndStream::No)
}
pub fn into_part(self) -> DataOrHeadersWithFlag {
match self {
DataOrTrailers::Data(data, end_stream) => DataOrHeadersWithFlag {
content: DataOrHeaders::Data(data),
last: end_stream == EndStream::Yes,
},
DataOrTrailers::Trailers(headers) => DataOrHeadersWithFlag {
content: DataOrHeaders::Headers(headers),
last: true,
},
}
}
}
pub struct HttpStreamAfterHeaders(
pub Pin<Box<dyn Stream<Item = result::Result<DataOrTrailers>> + Send + 'static>>,
);
impl HttpStreamAfterHeaders {
pub fn new<S>(s: S) -> HttpStreamAfterHeaders
where
S: Stream<Item = result::Result<DataOrTrailers>> + Send + 'static,
{
HttpStreamAfterHeaders(Box::pin(s))
}
pub(crate) fn from_parts<S>(s: S) -> HttpStreamAfterHeaders
where
S: Stream<Item = result::Result<DataOrHeadersWithFlag>> + Send + 'static,
{
HttpStreamAfterHeaders::new(s.map_ok(DataOrHeadersWithFlag::into_after_headers))
}
pub fn empty() -> HttpStreamAfterHeaders {
HttpStreamAfterHeaders::new(stream::empty())
}
pub fn bytes<S>(bytes: S) -> HttpStreamAfterHeaders
where
S: Stream<Item = result::Result<Bytes>> + Send + 'static,
{
HttpStreamAfterHeaders::new(bytes.map_ok(DataOrTrailers::intermediate_data))
}
pub fn once(part: DataOrHeaders) -> HttpStreamAfterHeaders {
let part = match part {
DataOrHeaders::Data(data) => DataOrTrailers::Data(data, EndStream::Yes),
DataOrHeaders::Headers(header) => DataOrTrailers::Trailers(header),
};
HttpStreamAfterHeaders::new(stream::once(future::ok(part)))
}
pub fn once_bytes<B>(bytes: B) -> HttpStreamAfterHeaders
where
B: Into<Bytes>,
{
HttpStreamAfterHeaders::once(DataOrHeaders::Data(bytes.into()))
}
pub fn filter_data(self) -> impl Stream<Item = result::Result<Bytes>> + Send {
self.try_filter_map(|p| {
future::ok(match p {
DataOrTrailers::Data(data, ..) => Some(data),
DataOrTrailers::Trailers(..) => None,
})
})
}
pub(crate) fn into_flag_stream(
self,
) -> impl Stream<Item = result::Result<DataOrHeadersWithFlag>> + Send {
TryStreamExt::map_ok(self.0, |r| DataOrTrailers::into_part(r))
}
pub(crate) fn _into_part_stream(self) -> DataOrHeadersWithFlagStream {
DataOrHeadersWithFlagStream::new(self.into_flag_stream())
}
pub fn catch_unwind(self) -> HttpStreamAfterHeaders {
HttpStreamAfterHeaders::new(panic::AssertUnwindSafe(self.0).catch_unwind().then(|r| {
future::ready(match r {
Ok(r) => r,
Err(e) => {
let e = any_to_string(e);
warn!("handler panicked: {}", e);
Err(error::Error::HandlerPanicked(e))
}
})
}))
}
}
impl Stream for HttpStreamAfterHeaders {
type Item = result::Result<DataOrTrailers>;
fn poll_next(mut self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<Self::Item>> {
Pin::new(&mut self.0).poll_next(context)
}
}