snix_store/nar/renderer/seekable/
segments.rs1use 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
13pub struct Segments {
16 segments: Vec<(u64, Segment)>,
17 total_len: u64,
18}
19
20impl Segments {
21 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 segments.push(Segment::Literal(std::mem::take(&mut buf)));
40 segments
41 }
42
43 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 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 let items = segments
73 .into_iter()
74 .zip(std::iter::once((skip_in_segment) as usize).chain(std::iter::repeat(0)));
75
76 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 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 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 let mut stream = ReaderStream::new(blob_reader);
130 while let Some(bytes) = stream.try_next().await? {
131 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(Vec<u8>),
170 BlobRef {
172 digest: B3Digest,
173 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
188fn walk_node(
196 segments: &mut Segments,
197 directories: &HashMap<B3Digest, Directory>,
198 node: &Node,
199 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 segments.push(Segment::Literal(std::mem::take(buf)));
219
220 segments.push(Segment::BlobRef {
222 digest: *digest,
223 size: *size,
224 });
225
226 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 let mut nar_node_directory = nar_node
241 .directory()
242 .expect("Snix bug: failed to write directory node as NAR");
243
244 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 nar_node_directory
256 .close()
257 .expect("Snix bug: failed to close directory node");
258 }
259 }
260}