Skip to main content

snix_store/nar/renderer/seekable/
segments.rs

1use bytes::Bytes;
2use futures::{StreamExt, TryStreamExt};
3use nix_compat::nar::writer::sync as nar_writer;
4use snix_castore::{B3Digest, Directory, Node, blobservice::BlobService};
5use std::{
6    collections::HashMap,
7    io::{self},
8    sync::atomic::AtomicU64,
9};
10use tokio::io::{AsyncBufRead, AsyncReadExt};
11use tokio_util::io::{InspectReader, ReaderStream, StreamReader};
12
13/// Contains a list of [Segment]s.
14/// For each segment, stores its offset to allow binary search seeking.
15pub struct Segments {
16    segments: Vec<(u64, Segment)>,
17    total_len: u64,
18}
19
20impl Segments {
21    /// Construct segments that can be used to assemble a NAR,
22    /// using the given root_node and directories.
23    /// Panics if the directory closure is not complete.
24    pub fn from_root_node_and_directories(
25        root_node: &Node,
26        directories: &HashMap<B3Digest, Directory>,
27    ) -> Self {
28        let mut segments = Self {
29            segments: vec![],
30            total_len: 0,
31        };
32
33        let mut buf: Vec<u8> = vec![];
34        let nar_node = nar_writer::open(&mut buf).expect("Snix bug: failed to open nar_writer");
35
36        walk_node(&mut segments, directories, root_node, nar_node);
37
38        // Flush the final segment
39        segments.push(Segment::Literal(std::mem::take(&mut buf)));
40        segments
41    }
42
43    /// Provides a [AsyncBufRead] into a [Segments] from the given offset.
44    ///
45    /// A [BlobService] to read blobs from needs to be specified, as well as
46    /// the desired concurrency (for segments, not blobs).
47    ///
48    /// This does not implement AsyncSeek, seeking can be done by calling this
49    /// again with another offset.
50    pub fn reader_for_offset<'bs, BS: BlobService + Clone + 'bs>(
51        &self,
52        offset: u64,
53        segment_concurrency: usize,
54        blob_service: BS,
55    ) -> Box<dyn AsyncBufRead + Send + Unpin + 'bs> {
56        // find the segment at the selected offset, or right before it
57        let (idx, skip_in_segment) = match self
58            .segments
59            .binary_search_by_key(&offset, |(offset, _)| *offset)
60        {
61            Ok(idx) => (idx, 0),
62            Err(idx) => (idx - 1, offset - self.segments[idx - 1].0),
63        };
64
65        let segments: Vec<_> = self.segments[idx..]
66            .iter()
67            .map(|(_offset, segment)| segment.clone())
68            .collect();
69
70        // We might need to skip something from the first segment.
71        // Turn the iterator of segments into bytes to skip at the beginning and the Segment itself.
72        let items = segments
73            .into_iter()
74            .zip(std::iter::once((skip_in_segment) as usize).chain(std::iter::repeat(0)));
75
76        // produce a stream of byte chunks
77        let bytes_stream = tokio_stream::iter(items)
78            .map(move |(segment, skip_in_segment)| {
79                let blob_service = blob_service.clone();
80                async move {
81                    let segment_len = segment.len();
82                    if skip_in_segment > 0 {
83                        debug_assert!(
84                            segment_len > skip_in_segment as u64,
85                            "Snix bug: segment size is smaller than bytes to skip"
86                        );
87                    }
88                    match segment {
89                        Segment::Literal(mut data) => futures::stream::once(async move {
90                            if skip_in_segment > 0 {
91                                data.rotate_left(skip_in_segment);
92                                data.truncate(data.len() - skip_in_segment);
93                            }
94                            Ok::<_, io::Error>(Bytes::from(data))
95                        })
96                        .boxed(),
97                        // skip over empty blobs
98                        Segment::BlobRef { size: 0, .. } => futures::stream::empty().boxed(),
99                        Segment::BlobRef { digest, size } => {
100                            async_stream::try_stream! {
101                                let blob_reader = blob_service.open_read(&digest).await
102                                    .map_err(io::Error::other)?
103                                    .ok_or_else(|| {
104                                        io::Error::new(
105                                            io::ErrorKind::NotFound,
106                                            format!("blob {0} not found", &digest),
107                                        )
108                                    })?;
109
110                                let bytes_read: AtomicU64 = AtomicU64::new(0);
111                                let blob_reader = InspectReader::new(blob_reader, |d| {
112                                    bytes_read.fetch_add(d.len() as u64, std::sync::atomic::Ordering::Relaxed);
113                                });
114
115                                // discard skip_in_segment bytes from the reader
116                                let blob_reader = if skip_in_segment > 0 {
117                                    let mut limited_rd = blob_reader.take(skip_in_segment as u64);
118                                    tokio::io::copy(
119                                        &mut limited_rd,
120                                        &mut tokio::io::sink(),
121                                    ).await?;
122
123                                    limited_rd.into_inner()
124                                } else {
125                                    blob_reader
126                                };
127
128                                // construct a ReaderStream for the rest
129                                let mut stream = ReaderStream::new(blob_reader);
130                                while let Some(bytes) = stream.try_next().await? {
131                                    // `bytes.len()` has already been added to the `bytes_read` counter.
132                                    if (bytes_read.load(std::sync::atomic::Ordering::Relaxed)) > size {
133                                        Err(io::Error::new(io::ErrorKind::InvalidData, "got more bytes from BlobReader than expected"))?
134                                    }
135                                    yield bytes;
136                                }
137
138                                drop(stream);
139
140                                if bytes_read.into_inner() < size {
141                                    Err(io::Error::new(io::ErrorKind::InvalidData, "got less bytes from BlobReader than expected"))?
142                                }
143                            }
144                            .boxed()
145                        }
146                    }
147                }
148            })
149            .buffered(segment_concurrency)
150            .flatten();
151
152        Box::new(StreamReader::new(bytes_stream))
153    }
154
155    pub fn total_len(&self) -> u64 {
156        self.total_len
157    }
158
159    fn push(&mut self, elem: Segment) {
160        let new_total_len = self.total_len() + elem.len();
161        self.segments.push((self.total_len(), elem));
162        self.total_len = new_total_len;
163    }
164}
165
166#[derive(Clone, Debug)]
167pub enum Segment {
168    // Literal bytes, aka 'NAR framing' around blob pointers
169    Literal(Vec<u8>),
170    // A pointer to a blob, by its digest. Also stores the size, so we can compute the segment length.
171    BlobRef {
172        digest: B3Digest,
173        // FUTUREWORK: maybe drop this and Segment::len(),
174        // where we want BlobRef size we can peek at the offset in the next segment.
175        size: u64,
176    },
177}
178
179impl Segment {
180    pub fn len(&self) -> u64 {
181        match self {
182            Segment::Literal(data) => data.len() as u64,
183            Segment::BlobRef { size, .. } => *size,
184        }
185    }
186}
187
188/// Used during construction.
189/// Recursively walks the node and its children, and pushes new segments.
190///
191/// The function is infallible, as:
192///  - we only write to buffers
193///  - all castore `PathComponent` and `SymlinkTarget` can be expressed in NAR
194///  - the passed directory closure is complete
195fn walk_node(
196    segments: &mut Segments,
197    directories: &HashMap<B3Digest, Directory>,
198    node: &Node,
199    // Includes a reference to the current segment's buffer
200    nar_node: nar_writer::Node<'_, Vec<u8>>,
201) {
202    match node {
203        snix_castore::Node::Symlink { target } => {
204            nar_node
205                .symlink(target.as_ref())
206                .expect("Snix bug: failed to write symlink as NAR");
207        }
208        snix_castore::Node::File {
209            digest,
210            size,
211            executable,
212        } => {
213            let (buf, skip) = nar_node
214                .file_manual_write(*executable, *size)
215                .expect("Snix bug: failed to write framing before file node as NAR");
216
217            // Flush the segment up until the beginning of the blob
218            segments.push(Segment::Literal(std::mem::take(buf)));
219
220            // Insert the blob segment
221            segments.push(Segment::BlobRef {
222                digest: *digest,
223                size: *size,
224            });
225
226            // Close the file node
227            // We **intentionally** do not write the file contents anywhere.
228            // Instead we have stored the blob reference in a Data::Blob segment,
229            // and the poll_read implementation will take care of serving the
230            // appropriate blob at this offset.
231            skip.close(buf)
232                .expect("Snix bug: failed to close NAR file node");
233        }
234        snix_castore::Node::Directory { digest, .. } => {
235            let directory = directories
236                .get(digest)
237                .expect("Snix bug: referenced directory missing from directory closure");
238
239            // start a directory node
240            let mut nar_node_directory = nar_node
241                .directory()
242                .expect("Snix bug: failed to write directory node as NAR");
243
244            // for each node in the directory, create a new entry with its name,
245            // and then recurse on that entry.
246            for (name, node) in directory.nodes() {
247                let child_node = nar_node_directory
248                    .entry(name.as_ref())
249                    .expect("Snix bug: failed to write NAR entry");
250
251                walk_node(segments, directories, node, child_node);
252            }
253
254            // close the directory
255            nar_node_directory
256                .close()
257                .expect("Snix bug: failed to close directory node");
258        }
259    }
260}