sark_core/http/response/
stream.rs1use std::pin::Pin;
2use std::task::Poll;
3
4use dope_fiber::{Context, Fiber};
5use http::StatusCode;
6use o3::buffer::Shared;
7
8use super::wire_emit::{HeadWrite, TransferEncodingChunked};
9
10pub const CHUNK_TERMINATOR: &[u8; 5] = b"0\r\n\r\n";
11
12pub struct Stream<S> {
13 status: StatusCode,
14 wire_headers: Vec<Shared>,
15 stream: S,
16}
17
18pub struct IterStream<I> {
19 iter: I,
20}
21
22impl<'d, I> Fiber<'d> for IterStream<I>
23where
24 I: Iterator<Item = Shared> + Unpin + 'static,
25{
26 type Output = Option<Shared>;
27 fn poll(self: Pin<&mut Self>, _cx: Pin<&mut Context<'_, 'd>>) -> Poll<Option<Shared>> {
28 Poll::Ready(self.get_mut().iter.next())
29 }
30}
31
32impl<S> Stream<S> {
33 pub fn new(stream: S) -> Self {
34 Self {
35 status: StatusCode::OK,
36 wire_headers: Vec::new(),
37 stream,
38 }
39 }
40
41 pub fn header(mut self, name: &[u8], value: &[u8]) -> Self {
42 let mut buf = Vec::new();
43 buf.extend_from_slice(name);
44 buf.extend_from_slice(b": ");
45 buf.extend_from_slice(value);
46 buf.extend_from_slice(b"\r\n");
47 self.wire_headers.push(Shared::from(buf));
48 self
49 }
50
51 pub fn write_head_stream(self, out: &mut [u8], date: &[u8; 29]) -> Option<(usize, S)> {
52 let status_str = self.status.as_str().as_bytes();
53 let reason = self
54 .status
55 .canonical_reason()
56 .map(str::as_bytes)
57 .unwrap_or(b"");
58
59 let head = HeadWrite {
60 status_str,
61 reason,
62 headers: self.wire_headers.as_slice(),
63 framing: TransferEncodingChunked,
64 };
65 if out.len() < head.wire_len() {
66 return None;
67 }
68
69 let written = head.write(out, date);
70 Some((written.len, self.stream))
71 }
72}
73
74impl<II> Stream<IterStream<II>>
75where
76 II: Iterator<Item = Shared> + Unpin + 'static,
77{
78 pub fn from_chunks<I>(iter: I) -> Self
79 where
80 I: IntoIterator<Item = Shared, IntoIter = II>,
81 {
82 Self::new(IterStream {
83 iter: iter.into_iter(),
84 })
85 }
86}