1use nix_compat::{
2 nar::reader::r#async as nar_reader,
3 nixhash::{CAHash, HashAlgo, NixHash, NixHashDigester, Sha256Digester, copy_hashed},
4};
5use snix_castore::{
6 Node, PathBuf,
7 blobservice::BlobService,
8 directoryservice::DirectoryService,
9 import::{
10 IngestionEntry, IngestionError,
11 blobs::{self, ConcurrentBlobUploader},
12 ingest_entries,
13 },
14};
15use tokio::{
16 io::{AsyncBufRead, AsyncRead},
17 sync::mpsc,
18 try_join,
19};
20use tokio_util::io::InspectReader;
21
22#[derive(Debug, thiserror::Error)]
24pub enum NarIngestionError {
25 #[error("{0}")]
26 IngestionError(#[from] IngestionError<Error>),
27
28 #[error("Hash mismatch, expected: {expected}, got: {actual}.")]
29 HashMismatch { expected: NixHash, actual: NixHash },
30
31 #[error("Expected the nar to contain a single file.")]
32 TypeMismatch,
33
34 #[error("Ingestion failed: {0}")]
35 Io(#[from] std::io::Error),
36}
37
38pub async fn ingest_nar_and_hash<R, BS, DS>(
43 blob_service: BS,
44 directory_service: DS,
45 r: &mut R,
46 expected_cahash: &Option<CAHash>,
47) -> Result<(Node, [u8; 32], u64), NarIngestionError>
48where
49 R: AsyncRead + Unpin + Send,
50 BS: BlobService + Clone + 'static,
51 DS: DirectoryService,
52{
53 let mut nar_hash = Sha256Digester::new();
54 let mut nar_size = 0;
55
56 let mut r = tokio_util::io::InspectReader::new(r, |b| {
58 nar_size += b.len() as u64;
59 nar_hash.update(b);
60 });
61
62 match expected_cahash {
63 Some(CAHash::Nar(expected_hash)) => {
64 let (root_node, actual_hash, nar_hash) = if expected_hash.algo() == HashAlgo::Sha256 {
65 let mut r = tokio::io::BufReader::new(&mut r);
67
68 let root_node = ingest_nar(blob_service, directory_service, &mut r).await?;
69 let nar_hash = nar_hash.finalize();
70 (root_node, NixHash::from(nar_hash), nar_hash)
71 } else {
72 let mut digester = NixHashDigester::new(expected_hash.algo());
74 let mut r =
75 tokio::io::BufReader::new(InspectReader::new(r, |data| digester.update(data)));
76
77 let root_node = ingest_nar(blob_service, directory_service, &mut r).await?;
78 (root_node, digester.finalize(), nar_hash.finalize())
79 };
80
81 if actual_hash != *expected_hash {
82 return Err(NarIngestionError::HashMismatch {
83 expected: expected_hash.clone(),
84 actual: actual_hash,
85 });
86 }
87 Ok((root_node, nar_hash.into(), nar_size))
88 }
89 Some(CAHash::Flat(expected_hash)) => {
90 let mut r = tokio::io::BufReader::new(&mut r);
92 let root_node = ingest_nar(blob_service.clone(), directory_service, &mut r).await?;
93
94 if let Node::File { digest, .. } = &root_node {
96 if let Some(mut blob_reader) = blob_service.open_read(digest).await? {
97 let (_, actual_hash) = copy_hashed(
98 &mut blob_reader,
99 &mut tokio::io::sink(),
100 expected_hash.algo(),
101 )
102 .await?;
103
104 if actual_hash != *expected_hash {
105 return Err(NarIngestionError::HashMismatch {
106 expected: expected_hash.clone(),
107 actual: actual_hash,
108 });
109 }
110 Ok((root_node, nar_hash.finalize().into(), nar_size))
111 } else {
112 Err(NarIngestionError::Io(std::io::Error::other(
113 "Ingested data not found",
114 )))
115 }
116 } else {
117 Err(NarIngestionError::TypeMismatch)
118 }
119 }
120 _ => {
126 let mut r = tokio::io::BufReader::new(&mut r);
127 let root_node = ingest_nar(blob_service, directory_service, &mut r).await?;
128 Ok((root_node, nar_hash.finalize().into(), nar_size))
129 }
130 }
131}
132
133pub async fn ingest_nar<R, BS, DS>(
137 blob_service: BS,
138 directory_service: DS,
139 r: &mut R,
140) -> Result<Node, IngestionError<Error>>
141where
142 R: AsyncBufRead + Unpin + Send,
143 BS: BlobService + Clone + 'static,
144 DS: DirectoryService,
145{
146 let root_node = nar_reader::open(r).await.map_err(Error::IO)?;
149
150 let (tx, rx) = mpsc::channel(1);
151 let rx = tokio_stream::wrappers::ReceiverStream::new(rx);
152
153 let produce = async move {
154 let mut blob_uploader = ConcurrentBlobUploader::new(blob_service);
155
156 let res = produce_nar_inner(
157 &mut blob_uploader,
158 root_node,
159 "root".parse().unwrap(), tx.clone(),
161 )
162 .await;
163
164 if let Err(err) = blob_uploader.join().await {
165 tx.send(Err(err.into()))
166 .await
167 .map_err(|e| Error::IO(std::io::Error::new(std::io::ErrorKind::BrokenPipe, e)))?;
168 }
169
170 tx.send(res)
171 .await
172 .map_err(|e| Error::IO(std::io::Error::new(std::io::ErrorKind::BrokenPipe, e)))?;
173
174 Ok(())
175 };
176
177 let consume = ingest_entries(directory_service, rx);
178
179 let (_, node) = try_join!(produce, consume)?;
180
181 Ok(node)
182}
183
184async fn produce_nar_inner<BS>(
185 blob_uploader: &mut ConcurrentBlobUploader<BS>,
186 node: nar_reader::Node<'_, '_>,
187 path: PathBuf,
188 tx: mpsc::Sender<Result<IngestionEntry, Error>>,
189) -> Result<IngestionEntry, Error>
190where
191 BS: BlobService + Clone + 'static,
192{
193 Ok(match node {
194 nar_reader::Node::Symlink { target } => IngestionEntry::Symlink { path, target },
195 nar_reader::Node::File {
196 executable,
197 mut reader,
198 } => {
199 let size = reader.len();
200 let digest = blob_uploader.upload(&path, size, &mut reader).await?;
201
202 IngestionEntry::Regular {
203 path,
204 size,
205 executable,
206 digest,
207 }
208 }
209 nar_reader::Node::Directory(mut dir_reader) => {
210 while let Some(entry) = dir_reader.next().await? {
211 let mut path = path.clone();
212
213 path.try_push(entry.name)
215 .expect("Snix bug: failed to join name");
216
217 let entry = Box::pin(produce_nar_inner(
218 blob_uploader,
219 entry.node,
220 path,
221 tx.clone(),
222 ))
223 .await?;
224
225 tx.send(Ok(entry)).await.map_err(|e| {
226 Error::IO(std::io::Error::new(std::io::ErrorKind::BrokenPipe, e))
227 })?;
228 }
229
230 IngestionEntry::Dir { path }
231 }
232 })
233}
234
235#[derive(Debug, thiserror::Error)]
236pub enum Error {
237 #[error(transparent)]
238 IO(#[from] std::io::Error),
239
240 #[error(transparent)]
241 BlobUpload(#[from] blobs::Error),
242}
243
244#[cfg(test)]
245mod test {
246 use crate::fixtures::{
247 NAR_CONTENTS_COMPLICATED, NAR_CONTENTS_HELLOWORLD, NAR_CONTENTS_SYMLINK,
248 };
249 use crate::nar::{NarIngestionError, ingest_nar, ingest_nar_and_hash};
250 use std::io::Cursor;
251 use std::sync::Arc;
252
253 use hex_literal::hex;
254 use mockall::predicate;
255 use nix_compat::nixhash::{CAHash, NixHash};
256 use rstest::*;
257 use snix_castore::Node;
258 use snix_castore::blobservice::{MockBlobService, TestBlobWriter};
259 use snix_castore::directoryservice::{MockDirectoryPutter, MockDirectoryService};
260 use snix_castore::fixtures::{
261 DIRECTORY_COMPLICATED, DIRECTORY_WITH_KEEP, EMPTY_BLOB_DIGEST, HELLOWORLD_BLOB_CONTENTS,
262 HELLOWORLD_BLOB_DIGEST,
263 };
264 use snix_castore::utils::gen_test_blob_service;
265
266 #[tokio::test]
267 async fn single_symlink() {
268 let root_node = ingest_nar(
269 Arc::new(MockBlobService::new()),
270 MockDirectoryService::new(),
271 &mut Cursor::new(&NAR_CONTENTS_SYMLINK),
272 )
273 .await
274 .expect("must parse");
275
276 assert_eq!(
277 Node::Symlink {
278 target: "/nix/store/somewhereelse".try_into().unwrap()
279 },
280 root_node
281 );
282 }
283
284 #[tokio::test]
285 async fn single_file() {
286 let mut blob_service = MockBlobService::new();
287 let mut seq = mockall::Sequence::new();
288 blob_service
289 .expect_has()
290 .once()
291 .with(predicate::eq(&*HELLOWORLD_BLOB_DIGEST))
292 .return_once(|_| Ok(false))
293 .in_sequence(&mut seq);
294
295 blob_service
296 .expect_open_write()
297 .once()
298 .return_once(|| Box::new(TestBlobWriter::new()))
299 .in_sequence(&mut seq);
300
301 let root_node = ingest_nar(
302 Arc::new(blob_service),
303 MockDirectoryService::new(),
304 &mut Cursor::new(&NAR_CONTENTS_HELLOWORLD),
305 )
306 .await
307 .expect("must parse");
308
309 assert_eq!(
310 Node::File {
311 digest: *HELLOWORLD_BLOB_DIGEST,
312 size: HELLOWORLD_BLOB_CONTENTS.len() as u64,
313 executable: false,
314 },
315 root_node
316 );
317 }
318
319 #[tokio::test]
320 async fn complicated() {
321 let mut blob_service = MockBlobService::new();
322 let mut seq = mockall::Sequence::new();
323 blob_service
324 .expect_has()
325 .once()
326 .with(predicate::eq(&*EMPTY_BLOB_DIGEST))
327 .return_once(|_| Ok(false))
328 .in_sequence(&mut seq);
329 blob_service
330 .expect_open_write()
331 .once()
332 .return_once(|| Box::new(TestBlobWriter::new()))
333 .in_sequence(&mut seq);
334 blob_service
335 .expect_has()
336 .once()
337 .with(predicate::eq(&*EMPTY_BLOB_DIGEST))
338 .return_once(|_| Ok(true))
339 .in_sequence(&mut seq);
340 let mut directory_service = MockDirectoryService::new();
341 directory_service
342 .expect_put_multiple_start()
343 .once()
344 .return_once(|| {
345 let mut directory_putter = MockDirectoryPutter::new();
346 let mut seq = mockall::Sequence::new();
347 directory_putter
348 .expect_put()
349 .once()
350 .with(predicate::eq(&*DIRECTORY_WITH_KEEP))
351 .returning(|_| Ok(()))
352 .in_sequence(&mut seq);
353 directory_putter
354 .expect_put()
355 .once()
356 .with(predicate::eq(&*DIRECTORY_COMPLICATED))
357 .returning(|_| Ok(()))
358 .in_sequence(&mut seq);
359 directory_putter
360 .expect_close()
361 .once()
362 .returning(|| Ok(DIRECTORY_COMPLICATED.digest()))
363 .in_sequence(&mut seq);
364 Box::new(directory_putter)
365 });
366
367 let root_node = ingest_nar(
368 Arc::new(blob_service),
369 directory_service,
370 &mut Cursor::new(&NAR_CONTENTS_COMPLICATED),
371 )
372 .await
373 .expect("must parse");
374
375 assert_eq!(
376 Node::Directory {
377 digest: DIRECTORY_COMPLICATED.digest(),
378 size: DIRECTORY_COMPLICATED.size()
379 },
380 root_node,
381 );
382 }
383
384 #[rstest]
385 #[case::nar_sha256(Some(CAHash::Nar(NixHash::Sha256(hex!("fbd52279a8df024c9fd5718de4103bf5e760dc7f2cf49044ee7dea87ab16911a")))), NAR_CONTENTS_COMPLICATED.as_slice())]
386 #[case::nar_sha512(Some(CAHash::Nar(NixHash::Sha512(Box::new(hex!("ff5d43941411f35f09211f8596b426ee6e4dd3af1639e0ed2273cbe44b818fc4a59e3af02a057c5b18fbfcf435497de5f1994206c137f469b3df674966a922f0"))))), NAR_CONTENTS_COMPLICATED.as_slice())]
387 #[case::flat_md5(Some(CAHash::Flat(NixHash::Md5(hex!("fd076287532e86365e841e92bfc50d8c")))), NAR_CONTENTS_HELLOWORLD.as_slice() )]
388 #[case::nar_symlink_sha1(Some(CAHash::Nar(NixHash::Sha1(hex!("f24eeaaa9cc016bab030bf007cb1be6483e7ba9e")))), NAR_CONTENTS_SYMLINK.as_slice())]
389 #[tokio::test]
390 async fn ingest_with_cahash_mismatch(
391 #[case] ca_hash: Option<CAHash>,
392 #[case] nar_content: &[u8],
393 ) {
394 use snix_castore::utils::gen_test_directory_service;
395
396 let err = ingest_nar_and_hash(
397 gen_test_blob_service(),
398 gen_test_directory_service(),
399 &mut Cursor::new(nar_content),
400 &ca_hash,
401 )
402 .await
403 .expect_err("Ingestion should have failed");
404 assert!(
405 matches!(err, NarIngestionError::HashMismatch { .. }),
406 "CAHash should have mismatched"
407 );
408 }
409
410 #[rstest]
411 #[case::nar_sha256(Some(CAHash::Nar(NixHash::Sha256(hex!("ebd52279a8df024c9fd5718de4103bf5e760dc7f2cf49044ee7dea87ab16911a")))), &NAR_CONTENTS_COMPLICATED.clone())]
412 #[case::nar_sha512(Some(CAHash::Nar(NixHash::Sha512(Box::new(hex!("1f5d43941411f35f09211f8596b426ee6e4dd3af1639e0ed2273cbe44b818fc4a59e3af02a057c5b18fbfcf435497de5f1994206c137f469b3df674966a922f0"))))), &NAR_CONTENTS_COMPLICATED.clone())]
413 #[case::flat_md5(Some(CAHash::Flat(NixHash::Md5(hex!("ed076287532e86365e841e92bfc50d8c")))), &NAR_CONTENTS_HELLOWORLD.clone())]
414 #[case::nar_symlink_sha1(Some(CAHash::Nar(NixHash::Sha1(hex!("424eeaaa9cc016bab030bf007cb1be6483e7ba9e")))), &NAR_CONTENTS_SYMLINK.clone())]
415 #[tokio::test]
416 async fn ingest_with_cahash_correct(
417 #[case] ca_hash: Option<CAHash>,
418 #[case] nar_content: &[u8],
419 ) {
420 ingest_nar_and_hash(
421 snix_castore::utils::gen_test_blob_service(),
422 snix_castore::utils::gen_test_directory_service(),
423 &mut Cursor::new(nar_content),
424 &ca_hash,
425 )
426 .await
427 .expect("CAHash should have matched");
428 }
429
430 #[rstest]
431 #[case::nar_sha256(Some(CAHash::Flat(NixHash::Sha256(hex!("ebd52279a8df024c9fd5718de4103bf5e760dc7f2cf49044ee7dea87ab16911a")))), &NAR_CONTENTS_COMPLICATED.clone())]
432 #[case::nar_symlink_sha1(Some(CAHash::Flat(NixHash::Sha1(hex!("424eeaaa9cc016bab030bf007cb1be6483e7ba9e")))), &NAR_CONTENTS_SYMLINK.clone())]
433 #[tokio::test]
434 async fn ingest_with_flat_non_file(
435 #[case] ca_hash: Option<CAHash>,
436 #[case] nar_content: &[u8],
437 ) {
438 let err = ingest_nar_and_hash(
439 snix_castore::utils::gen_test_blob_service(),
440 snix_castore::utils::gen_test_directory_service(),
441 &mut Cursor::new(nar_content),
442 &ca_hash,
443 )
444 .await
445 .expect_err("Ingestion should have failed");
446
447 assert!(
448 matches!(err, NarIngestionError::TypeMismatch),
449 "Flat cahash should only be allowed for single file nars"
450 );
451 }
452}