#[cfg(feature = "axum")]
use crate::canonical::CanonChunk;
#[cfg(feature = "axum")]
use crate::error::ProxyError;
#[cfg(feature = "axum")]
use crate::error::error_response;
#[cfg(feature = "axum")]
use axum::body::Body;
#[cfg(feature = "axum")]
use axum::http::StatusCode;
#[cfg(feature = "axum")]
use axum::response::Response;
#[cfg(feature = "axum")]
use futures::{Stream, StreamExt};
#[cfg(feature = "axum")]
pub fn sse_response<S, F, G>(inner: S, frame: F, finish: G) -> Response
where
S: Stream<Item = Result<CanonChunk, ProxyError>> + Unpin + Send + 'static,
F: FnMut(&mut Vec<String>, Result<CanonChunk, ProxyError>) + Send + 'static,
G: FnOnce(&mut Vec<String>) + Send + 'static,
{
let mut frame = frame;
let stream = inner.map(move |item| {
let mut out: Vec<String> = Vec::new();
frame(&mut out, item);
if out.is_empty() {
None
} else {
Some(
out.into_iter()
.map(Ok::<_, std::convert::Infallible>)
.collect::<Vec<_>>(),
)
}
});
let stream = stream.chain(futures::stream::once(async move {
let mut out: Vec<String> = Vec::new();
finish(&mut out);
if out.is_empty() {
None
} else {
Some(
out.into_iter()
.map(Ok::<_, std::convert::Infallible>)
.collect::<Vec<_>>(),
)
}
}));
let body = Body::from_stream(
stream
.filter_map(|v| async move { v })
.flat_map(futures::stream::iter),
);
Response::builder()
.status(StatusCode::OK)
.header("content-type", "text/event-stream")
.header("cache-control", "no-cache")
.header("x-accel-buffering", "no")
.body(body)
.unwrap_or_else(|e| error_response(&ProxyError::Internal(e.into())))
}