Skip to main content

mkit_server/pipeline/
download.rs

1//! `DownloadPack` as a stream of chunks (SPEC-TRANSPORT-CONNECT ยง6.2).
2//!
3//! The chunks follow [`chunk_plan`]: every chunk but the last is exactly
4//! `download_chunk_max` bytes (800 KiB by default, overview Q17), and an
5//! empty pack is one empty `last` chunk. The stream re-chunks the blob
6//! store's body as it arrives, holding at most one chunk plus one store
7//! piece, never the pack.
8
9use 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/// One `PackChunk` of a download.
24#[derive(Debug, Clone, PartialEq, Eq)]
25pub struct DownloadChunk {
26    /// Byte offset of `data` in the pack.
27    pub offset: u64,
28    /// The bytes.
29    pub data: Bytes,
30    /// Whether this is the final chunk.
31    pub last: bool,
32}
33
34/// A pack download: its length, for the header message, then its chunks.
35/// The stream owns its data (`'static`), so a binding can hand it to
36/// Connect's `Response::stream_ok` without borrowing the pipeline.
37pub struct DownloadStream {
38    /// The pack's length in bytes.
39    pub total_bytes: u64,
40    /// The chunks, contiguous from offset 0; only the final one is `last`.
41    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    /// Chunk `body` into pieces of at most `max` bytes; `outcome` records
54    /// the request at the stream's end or first failure.
55    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
76/// A failed body read: a redacted `internal`, its detail logged.
77fn 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
83/// Re-chunks a blob body along [`chunk_plan`]'s spans.
84struct Rechunk<I> {
85    spans: I,
86    span: Option<ChunkSpan>,
87    source: Option<BoxStream<'static, Result<Bytes, StoreError>>>,
88    /// The unread rest of the current store piece.
89    rest: Bytes,
90    /// The chunk being assembled when it spans store pieces.
91    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                // Success is recorded when the `last` chunk is yielded: a
134                // binding may stop polling there (its client has the
135                // whole pack), and that must not count as `canceled`.
136                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}