Skip to main content

snix_store/nar/
import.rs

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/// Represents errors that can happen during nar ingestion.
23#[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
38/// Ingests the contents from a [AsyncRead] providing NAR into the snix store,
39/// interacting with a [BlobService] and [DirectoryService].
40/// Returns the castore root node, as well as the sha256 and size of the NAR
41/// contents ingested.
42pub 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    // Assemble NarHash and NarSize as we read bytes.
57    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                // If this is the required CAHash, we're already computing excatly this in `nar_hash` above.
66                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                // For the other algos, wrap the reader with another digester.
73                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            // ingest as NAR
91            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            // The resulting root node must be Node::File, else CAHash::Flat is not applicable
95            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        // We either got CAHash::Text, or no CAHash at all, so we just don't do any additional
121        // hash calculation/validation.
122        // FUTUREWORK: We should figure out what to do with CAHash::Text, according to nix-cpp
123        // they don't handle it either:
124        // https://github.com/NixOS/nix/blob/3e9cc78eb5e5c4f1e762e201856273809fd92e71/src/libstore/local-store.cc#L1099-L1133
125        _ => {
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
133/// Ingests the contents from a [AsyncRead] providing NAR into the snix store,
134/// interacting with a [BlobService] and [DirectoryService].
135/// It returns the castore root node or an error.
136pub 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    // open the NAR for reading.
147    // The NAR reader emits nodes in DFS preorder.
148    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(), // HACK: the root node sent to ingest_entries may not be ROOT.
160            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                // valid NAR names are valid castore names
214                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}