mkit_server/pipeline/
download.rs1use core::fmt;
10use core::pin::Pin;
11use core::task::{Context, Poll, ready};
12
13use bytes::{Bytes, BytesMut};
14use futures_core::Stream;
15
16use super::outcome::Outcome;
17use crate::download::{ChunkSpan, chunk_plan};
18use crate::error::ServerError;
19use crate::rt::BoxStream;
20use crate::storage_error::{StorageOp, describe_and_map};
21use crate::store::{BlobBody, StoreError};
22
23#[derive(Debug, Clone, PartialEq, Eq)]
25pub struct DownloadChunk {
26 pub offset: u64,
28 pub data: Bytes,
30 pub last: bool,
32}
33
34pub struct DownloadStream {
38 pub total_bytes: u64,
40 pub chunks: BoxStream<'static, Result<DownloadChunk, ServerError>>,
42}
43
44impl fmt::Debug for DownloadStream {
45 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
46 f.debug_struct("DownloadStream")
47 .field("total_bytes", &self.total_bytes)
48 .finish_non_exhaustive()
49 }
50}
51
52impl DownloadStream {
53 pub(crate) fn new(body: BlobBody, max: usize, outcome: Option<Outcome>) -> Self {
56 let (total, rest, source) = match body {
57 BlobBody::Bytes(bytes) => (bytes.len() as u64, bytes, None),
58 BlobBody::Stream { len, stream } => (len, Bytes::new(), Some(stream)),
59 };
60 let chunks = Rechunk {
61 spans: chunk_plan(total, max),
62 span: None,
63 source,
64 rest,
65 buf: BytesMut::new(),
66 done: false,
67 outcome,
68 };
69 Self {
70 total_bytes: total,
71 chunks: Box::pin(chunks),
72 }
73 }
74}
75
76fn read_error(detail: impl fmt::Display) -> ServerError {
78 let (line, err) = describe_and_map(StorageOp::BlobRead, detail);
79 tracing::warn!(detail = %line, "storage failure");
80 err
81}
82
83struct Rechunk<I> {
85 spans: I,
86 span: Option<ChunkSpan>,
87 source: Option<BoxStream<'static, Result<Bytes, StoreError>>>,
88 rest: Bytes,
90 buf: BytesMut,
92 done: bool,
93 outcome: Option<Outcome>,
94}
95
96impl<I: Iterator<Item = ChunkSpan> + Unpin> Rechunk<I> {
97 fn fail(&mut self, err: ServerError) -> Poll<Option<Result<DownloadChunk, ServerError>>> {
98 self.done = true;
99 if let Some(outcome) = &mut self.outcome {
100 outcome.record(Err(&err));
101 }
102 Poll::Ready(Some(Err(err)))
103 }
104}
105
106impl<I: Iterator<Item = ChunkSpan> + Unpin> Stream for Rechunk<I> {
107 type Item = Result<DownloadChunk, ServerError>;
108
109 fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
110 let this = self.get_mut();
111 if this.done {
112 return Poll::Ready(None);
113 }
114 loop {
115 let Some(span) = this.span.or_else(|| this.spans.next()) else {
116 this.done = true;
117 if let Some(outcome) = &mut this.outcome {
118 outcome.record(Ok(()));
119 }
120 return Poll::Ready(None);
121 };
122 this.span = Some(span);
123 let need = span.len - this.buf.len();
124 let data = if need == 0 {
125 Some(this.buf.split().freeze())
126 } else if this.buf.is_empty() && this.rest.len() >= need {
127 Some(this.rest.split_to(need))
128 } else {
129 None
130 };
131 if let Some(data) = data {
132 this.span = None;
133 if span.last
137 && let Some(outcome) = &mut this.outcome
138 {
139 outcome.record(Ok(()));
140 }
141 return Poll::Ready(Some(Ok(DownloadChunk {
142 offset: span.offset,
143 data,
144 last: span.last,
145 })));
146 }
147 if !this.rest.is_empty() {
148 let piece = this.rest.split_to(need.min(this.rest.len()));
149 this.buf.extend_from_slice(&piece);
150 continue;
151 }
152 let next = match this.source.as_mut() {
153 Some(source) => ready!(source.as_mut().poll_next(cx)),
154 None => None,
155 };
156 match next {
157 Some(Ok(piece)) => this.rest = piece,
158 Some(Err(e)) => return this.fail(read_error(e)),
159 None => return this.fail(read_error("blob body ended before its length")),
160 }
161 }
162 }
163}