1#![feature(doc_auto_cfg)]
2#![feature(doc_cfg)]
3
4use std::convert::Infallible;
5
6use bytes::{BufMut, Bytes, BytesMut};
7use futures_lite::stream::{Stream, unfold};
8use kanal::AsyncReceiver;
9
10pub fn frame_stream<B: AsRef<[u8]>>(
11 receiver: AsyncReceiver<B>,
12) -> impl Stream<Item = Result<Bytes, Infallible>> {
13 unfold(receiver, |rx| async move {
14 match rx.recv().await {
15 Ok(chunk) => {
16 let chunk = chunk.as_ref();
17 let len = chunk.len();
18 let mut framed_chunk = BytesMut::with_capacity(4 + len);
19 framed_chunk.put_u32_le(len as u32); framed_chunk.put_slice(&chunk);
21
22 let item = Ok(framed_chunk.freeze());
23 let next_state = rx;
24 Some((item, next_state))
25 }
26 Err(_) => None,
27 }
28 })
29}
30
31#[cfg(feature = "axum")]
32pub fn response<B: AsRef<[u8]> + Send + 'static>(
33 receiver: AsyncReceiver<B>,
34) -> axum::response::Response {
35 use axum::{body::Body, http::header::CONTENT_TYPE};
36
37 axum::response::Response::builder()
38 .header(CONTENT_TYPE, "application/octet-stream")
39 .body(Body::from_stream(frame_stream(receiver)))
40 .unwrap()
41}