Skip to main content

snix_castore/import/
fs.rs

1//! Import from a real filesystem.
2
3use futures::StreamExt;
4use futures::TryStreamExt;
5use futures::stream::BoxStream;
6use std::fs::FileType;
7use std::os::unix::ffi::OsStringExt;
8use std::os::unix::fs::MetadataExt;
9use std::os::unix::fs::PermissionsExt;
10use tokio::io::BufReader;
11use tokio_util::io::InspectReader;
12use tracing::Span;
13use tracing::{Instrument, info_span, instrument, trace_span};
14use tracing_indicatif::span_ext::IndicatifSpanExt;
15use walkdir::DirEntry;
16use walkdir::WalkDir;
17
18use crate::blobservice::BlobService;
19use crate::directoryservice::DirectoryService;
20use crate::refscan::{ReferenceReader, ReferenceScanner};
21use crate::{B3Digest, Node};
22
23use super::IngestionEntry;
24use super::IngestionError;
25use super::ingest_entries;
26
27/// Ingests the contents at a given path into the snix store, interacting with a [BlobService] and
28/// [DirectoryService]. It returns the root node or an error.
29///
30/// It does not follow symlinks at the root, they will be ingested as actual symlinks.
31///
32/// This function will walk the filesystem using `walkdir` and will consume
33/// `O(#number of entries)` space.
34#[instrument(
35    skip(blob_service, directory_service, reference_scanner),
36    fields(path, indicatif.pb_show = tracing::field::Empty),
37    err
38)]
39pub async fn ingest_path<BS, DS, P, P2>(
40    blob_service: BS,
41    directory_service: DS,
42    path: P,
43    reference_scanner: Option<&ReferenceScanner<P2>>,
44) -> Result<Node, IngestionError<Error>>
45where
46    P: AsRef<std::path::Path>,
47    BS: BlobService + Clone,
48    DS: DirectoryService,
49    P2: AsRef<[u8]> + Send + Sync,
50{
51    let span = Span::current();
52    span.pb_set_style(&snix_tracing::PB_SPINNER_LONG_STYLE);
53    span.pb_set_message(&format!("Ingesting {}", path.as_ref().display()));
54
55    let iter = WalkDir::new(path.as_ref())
56        .follow_links(false)
57        .follow_root_links(false)
58        .contents_first(true)
59        .into_iter();
60
61    ingest_entries(
62        directory_service,
63        dir_entries_to_ingestion_stream(blob_service, iter, path.as_ref(), reference_scanner)
64            .inspect_ok(|ingestion_entry| {
65                if matches!(ingestion_entry, IngestionEntry::Regular { .. }) {
66                    span.pb_inc(1);
67                }
68            }),
69    )
70    .await
71}
72
73/// Converts an iterator of [walkdir::DirEntry]s into a stream of ingestion entries.
74/// This can then be fed into [ingest_entries] to ingest all the entries into the castore.
75///
76/// The produced stream is buffered, so uploads can happen concurrently.
77///
78/// The root is the [std::path::Path] in the filesystem that is being ingested
79/// into castore.
80pub fn dir_entries_to_ingestion_stream<'a, BS, I, P>(
81    blob_service: BS,
82    walkdir_direntries: I,
83    root: &'a std::path::Path,
84    reference_scanner: Option<&'a ReferenceScanner<P>>,
85) -> BoxStream<'a, Result<IngestionEntry, Error>>
86where
87    BS: BlobService + Clone + 'a,
88    I: Iterator<Item = Result<DirEntry, walkdir::Error>> + Send + 'a,
89    P: AsRef<[u8]> + Send + Sync,
90{
91    let prefix = root.parent().unwrap_or_else(|| std::path::Path::new(""));
92
93    futures::stream::iter(walkdir_direntries)
94        .map(move |x| {
95            let blob_service = blob_service.clone();
96            async move {
97                match x {
98                    Ok(dir_entry) => {
99                        dir_entry_to_ingestion_entry(
100                            blob_service,
101                            &dir_entry,
102                            prefix,
103                            reference_scanner,
104                        )
105                        .await
106                    }
107                    Err(e) => Err(Error::Stat(
108                        prefix.to_path_buf(),
109                        e.into_io_error().expect("walkdir err must be some"),
110                    )),
111                }
112            }
113            .instrument(trace_span!("process_walkdir_direntry"))
114        })
115        .buffered(50)
116        .boxed()
117}
118
119/// Converts a [walkdir::DirEntry] into an [IngestionEntry], uploading blobs to the
120/// provided [BlobService].
121///
122/// The prefix path is stripped from the path of each entry. This is usually the parent path
123/// of the path being ingested so that the last element of the stream only has one component.
124pub async fn dir_entry_to_ingestion_entry<BS, P>(
125    blob_service: BS,
126    walkdir_direntry: &DirEntry,
127    prefix: &std::path::Path,
128    reference_scanner: Option<&ReferenceScanner<P>>,
129) -> Result<IngestionEntry, Error>
130where
131    BS: BlobService,
132    P: AsRef<[u8]>,
133{
134    let file_type = walkdir_direntry.file_type();
135
136    let fs_path = walkdir_direntry
137        .path()
138        .strip_prefix(prefix)
139        .expect("Snix bug: failed to strip root path prefix");
140
141    // convert to castore PathBuf
142    let path = crate::path::PathBuf::from_host_path(fs_path, false)
143        .unwrap_or_else(|e| panic!("Snix bug: walkdir direntry cannot be parsed: {e}"));
144
145    if file_type.is_dir() {
146        Ok(IngestionEntry::Dir { path })
147    } else if file_type.is_symlink() {
148        let target = tokio::fs::read_link(walkdir_direntry.path())
149            .await
150            .map_err(|e| Error::Stat(walkdir_direntry.path().to_path_buf(), e))?
151            .into_os_string()
152            .into_vec();
153
154        if let Some(reference_scanner) = &reference_scanner {
155            reference_scanner.scan(&target);
156        }
157
158        Ok(IngestionEntry::Symlink { path, target })
159    } else if file_type.is_file() {
160        let metadata = walkdir_direntry
161            .metadata()
162            .map_err(|e| Error::Stat(walkdir_direntry.path().to_path_buf(), e.into()))?;
163
164        let digest = upload_blob(blob_service, walkdir_direntry.path(), reference_scanner).await?;
165
166        Ok(IngestionEntry::Regular {
167            path,
168            size: metadata.size(),
169            // If it's executable by the user, it'll become executable.
170            // This matches nix's dump() function behaviour.
171            executable: metadata.permissions().mode() & 64 != 0,
172            digest,
173        })
174    } else {
175        Err(Error::FileType(fs_path.to_path_buf(), file_type))
176    }
177}
178
179/// Uploads the file at the provided [std::path::Path] to the [BlobService].
180#[instrument(skip_all, fields(blob.path=%path.as_ref().display()), err)]
181async fn upload_blob<BS, P>(
182    blob_service: BS,
183    path: impl AsRef<std::path::Path>,
184    reference_scanner: Option<&ReferenceScanner<P>>,
185) -> Result<B3Digest, Error>
186where
187    BS: BlobService,
188    P: AsRef<[u8]>,
189{
190    let progress_span = info_span!("upload_blobs", "indicatif.pb_show" = tracing::field::Empty);
191    progress_span.pb_set_style(&snix_tracing::PB_TRANSFER_STYLE);
192    progress_span.pb_start();
193    progress_span.pb_set_message(&format!("Uploading blob at {:?}", path.as_ref()));
194
195    let file = tokio::fs::File::open(path.as_ref())
196        .await
197        .map_err(|e| Error::BlobRead(path.as_ref().to_path_buf(), e))?;
198
199    let metadata = file
200        .metadata()
201        .await
202        .map_err(|e| Error::Stat(path.as_ref().to_path_buf(), e))?;
203
204    progress_span.pb_set_length(metadata.len());
205    let reader = InspectReader::new(file, |d| {
206        progress_span.pb_inc(d.len() as u64);
207    });
208
209    let mut writer = blob_service.open_write().await;
210    let mut reader = BufReader::with_capacity(128 * 1024, reader);
211    if let Some(reference_scanner) = reference_scanner {
212        let mut reader = ReferenceReader::new(reference_scanner, reader);
213        tokio::io::copy_buf(&mut reader, &mut writer)
214            .await
215            .map_err(|e| Error::BlobRead(path.as_ref().to_path_buf(), e))?;
216    } else {
217        tokio::io::copy_buf(&mut reader, &mut writer)
218            .await
219            .map_err(|e| Error::BlobRead(path.as_ref().to_path_buf(), e))?;
220    }
221
222    let digest = writer
223        .close()
224        .await
225        .map_err(|e| Error::BlobFinalize(path.as_ref().to_path_buf(), e))?;
226
227    Ok(digest)
228}
229
230#[derive(Debug, thiserror::Error)]
231pub enum Error {
232    #[error("unsupported file type at {0}: {1:?}")]
233    FileType(std::path::PathBuf, FileType),
234
235    #[error("unable to stat {0}: {1}")]
236    Stat(std::path::PathBuf, std::io::Error),
237
238    #[error("unable to open {0}: {1}")]
239    Open(std::path::PathBuf, std::io::Error),
240
241    #[error("unable to read {0}: {1}")]
242    BlobRead(std::path::PathBuf, std::io::Error),
243
244    // TODO: proper error for blob finalize
245    #[error("unable to finalize blob {0}: {1}")]
246    BlobFinalize(std::path::PathBuf, std::io::Error),
247}