1use 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#[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
73pub 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
119pub 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 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 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#[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 #[error("unable to finalize blob {0}: {1}")]
246 BlobFinalize(std::path::PathBuf, std::io::Error),
247}