frame_stream/
lib.rs

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); // 4 字节长度前缀
20        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}