1mod session;
7mod uri;
8
9use std::sync::Arc;
10
11use async_trait::async_trait;
12use futures::StreamExt;
13use futures::TryStreamExt;
14use futures::stream;
15pub use session::MultiFileSession;
16use session::MultiFileSessionExt;
17use tracing::debug;
18pub use uri::parse_uri_or_path;
19use vortex_error::VortexError;
20use vortex_error::VortexResult;
21use vortex_error::vortex_bail;
22use vortex_io::VortexReadAt;
23use vortex_io::filesystem::FileListing;
24use vortex_io::filesystem::FileSystemRef;
25use vortex_layout::LayoutReaderRef;
26use vortex_layout::scan::multi::LayoutReaderFactory;
27use vortex_layout::scan::multi::MultiLayoutDataSource;
28use vortex_scan::DataSource;
29use vortex_session::VortexSession;
30
31use crate::OpenOptionsSessionExt;
32use crate::VortexFile;
33use crate::VortexOpenOptions;
34
35pub struct MultiFileDataSource {
66 session: VortexSession,
67 glob_sources: Vec<(String, Option<FileSystemRef>)>,
70 open_options_fn: Arc<dyn Fn(VortexOpenOptions) -> VortexOpenOptions + Send + Sync>,
71}
72
73const GLOB_RESOLUTION_CONCURRENCY: usize = 16;
77
78impl MultiFileDataSource {
79 pub fn new(session: VortexSession) -> Self {
81 Self {
82 session,
83 glob_sources: Vec::new(),
84 open_options_fn: Arc::new(|opts| opts),
85 }
86 }
87
88 pub fn with_glob(mut self, glob: impl Into<String>, fs: Option<FileSystemRef>) -> Self {
95 let glob = glob.into();
96 let glob = if fs.is_none() && std::path::Path::new(&glob).is_relative() {
97 std::env::current_dir()
98 .map(|cwd| cwd.join(&glob).to_string_lossy().into_owned())
99 .unwrap_or(glob)
100 .trim_start_matches('/')
101 .to_string()
102 } else {
103 glob.trim_start_matches('/').to_string()
104 };
105 self.glob_sources.push((glob, fs));
106 self
107 }
108
109 pub fn with_open_options(
113 mut self,
114 f: impl Fn(VortexOpenOptions) -> VortexOpenOptions + Send + Sync + 'static,
115 ) -> Self {
116 self.open_options_fn = Arc::new(f);
117 self
118 }
119
120 pub async fn build(self) -> VortexResult<MultiLayoutDataSource> {
125 if self.glob_sources.is_empty() {
126 vortex_bail!("MultiFileDataSource requires at least one glob pattern");
127 }
128
129 let local_fs: Option<FileSystemRef> = self
131 .glob_sources
132 .iter()
133 .any(|(_, fs)| fs.is_none())
134 .then(|| create_local_filesystem(&self.session))
135 .transpose()?;
136
137 let globs: Vec<String> = self.glob_sources.iter().map(|(g, _)| g.clone()).collect();
138
139 let resolved: Vec<Vec<(FileListing, FileSystemRef)>> =
142 stream::iter(self.glob_sources.into_iter().map(|(glob, maybe_fs)| {
143 let fs = maybe_fs
146 .or_else(|| local_fs.as_ref().map(Arc::clone))
147 .unwrap_or_else(|| {
148 unreachable!("local_fs is set when any glob lacks a filesystem")
149 });
150 async move {
151 let files: Vec<FileListing> = fs.glob(&glob)?.try_collect().await?;
152 Ok::<_, VortexError>(
153 files
154 .into_iter()
155 .map(|file| (file, Arc::clone(&fs)))
156 .collect(),
157 )
158 }
159 }))
160 .buffered(GLOB_RESOLUTION_CONCURRENCY)
161 .try_collect()
162 .await?;
163 let all_files: Vec<(FileListing, FileSystemRef)> = resolved.into_iter().flatten().collect();
164
165 if all_files.is_empty() {
166 vortex_bail!("No files matched the glob pattern(s): {:?}", globs);
167 }
168
169 let file_count = all_files.len();
170 debug!(file_count, glob = ?globs, "discovered files");
171
172 let (first_file_listing, first_fs) = &all_files[0];
174 let open_fn = self.open_options_fn.as_ref();
175 let first_file = open_file(first_fs, first_file_listing, &self.session, open_fn).await?;
176 let first_reader = first_file.layout_reader()?;
177
178 let byte_sizes: Vec<Option<u64>> = all_files.iter().map(|(file, _)| file.size).collect();
179
180 let factories: Vec<Arc<dyn LayoutReaderFactory>> = all_files[1..]
181 .iter()
182 .map(|(file, fs)| {
183 Arc::new(VortexFileReaderFactory {
184 fs: Arc::clone(fs),
185 file: file.clone(),
186 session: self.session.clone(),
187 open_options_fn: Arc::clone(&self.open_options_fn),
188 }) as Arc<dyn LayoutReaderFactory>
189 })
190 .collect();
191
192 let inner = MultiLayoutDataSource::new_with_first(
193 first_reader,
194 factories,
195 byte_sizes,
196 &self.session,
197 );
198
199 debug!(file_count, dtype = %inner.dtype(), "built MultiFileDataSource");
200
201 Ok(inner)
202 }
203}
204
205#[cfg(feature = "object_store")]
210fn create_local_filesystem(session: &VortexSession) -> VortexResult<FileSystemRef> {
211 use vortex_io::object_store::ObjectStoreFileSystem;
212 use vortex_io::session::RuntimeSessionExt;
213
214 let store = Arc::new(object_store::local::LocalFileSystem::default());
215 let fs: FileSystemRef = Arc::new(ObjectStoreFileSystem::new(store, session.handle()));
216 Ok(fs)
217}
218
219#[cfg(not(feature = "object_store"))]
220fn create_local_filesystem(_session: &VortexSession) -> VortexResult<FileSystemRef> {
221 vortex_bail!(
222 "The 'object_store' feature is required for automatic local filesystem creation. \
223 Either enable the feature or provide a filesystem via .with_filesystem()."
224 );
225}
226
227async fn open_file(
229 fs: &FileSystemRef,
230 file: &FileListing,
231 session: &VortexSession,
232 open_options_fn: &(dyn Fn(VortexOpenOptions) -> VortexOpenOptions + Send + Sync),
233) -> VortexResult<VortexFile> {
234 tracing::trace!(path = %file.path, "opening vortex file");
235
236 let source = fs.open_read(&file.path).await?;
237 let key = source.uri().is_none().then_some(file.path.as_str());
238 open_cached(session, key, source, file.size, open_options_fn).await
239}
240
241pub async fn open_cached(
247 session: &VortexSession,
248 mut key: Option<&str>,
249 source: Arc<dyn VortexReadAt>,
250 file_size: Option<u64>,
251 open_options_fn: &(dyn Fn(VortexOpenOptions) -> VortexOpenOptions + Send + Sync),
252) -> VortexResult<VortexFile> {
253 let uri = source.uri().cloned();
254 if key.is_none() {
255 key = uri.as_deref();
256 }
257 let Some(key) = key else {
258 vortex_bail!("Missing cache key");
259 };
260
261 let mut options = open_options_fn(session.open_options());
262 if let Some(size) = file_size {
263 options = options.with_file_size(size);
264 }
265
266 {
267 if let Some(footer) = session.multi_file().get_footer(key) {
268 options = options.with_footer(footer);
269 }
270 }
271
272 let file = options.open(source).await?;
273 session.multi_file().put_footer(key, file.footer().clone());
274 Ok(file)
275}
276
277struct VortexFileReaderFactory {
279 fs: FileSystemRef,
280 file: FileListing,
281 session: VortexSession,
282 open_options_fn: Arc<dyn Fn(VortexOpenOptions) -> VortexOpenOptions + Send + Sync>,
283}
284
285#[async_trait]
286impl LayoutReaderFactory for VortexFileReaderFactory {
287 async fn open(&self) -> VortexResult<Option<LayoutReaderRef>> {
288 let file = open_file(
289 &self.fs,
290 &self.file,
291 &self.session,
292 self.open_options_fn.as_ref(),
293 )
294 .await?;
295
296 Ok(Some(file.layout_reader()?))
297 }
298}
299
300#[cfg(test)]
301mod tests {
302 use std::fmt;
303
304 use vortex_array::IntoArray;
305 use vortex_array::array_session;
306 use vortex_buffer::ByteBuffer;
307 use vortex_buffer::ByteBufferMut;
308 use vortex_buffer::buffer;
309 use vortex_io::VortexReadAt;
310 use vortex_io::filesystem::FileSystem;
311 use vortex_io::session::RuntimeSession;
312 use vortex_layout::session::LayoutSession;
313
314 use super::*;
315 use crate::WriteOptionsSessionExt;
316
317 struct MetadataFileSystem {
318 bytes: ByteBuffer,
319 }
320
321 impl fmt::Debug for MetadataFileSystem {
322 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
323 f.debug_struct("MetadataFileSystem").finish()
324 }
325 }
326
327 #[async_trait]
328 impl FileSystem for MetadataFileSystem {
329 fn list(&self, _prefix: &str) -> stream::BoxStream<'_, VortexResult<FileListing>> {
330 stream::empty().boxed()
331 }
332
333 async fn head(&self, path: &str) -> VortexResult<Option<FileListing>> {
334 Ok((path == "metadata.vortex").then_some(FileListing {
335 path: path.to_string(),
336 size: Some(self.bytes.len() as u64),
337 }))
338 }
339
340 async fn open_read(&self, _path: &str) -> VortexResult<Arc<dyn VortexReadAt>> {
341 Ok(Arc::new(self.bytes.clone()))
342 }
343
344 async fn delete(&self, _path: &str) -> VortexResult<()> {
345 Ok(())
346 }
347 }
348
349 async fn assert_cached_footer_metadata_order(default_first: bool) -> VortexResult<()> {
350 let session = array_session()
351 .with::<LayoutSession>()
352 .with::<RuntimeSession>()
353 .with::<MultiFileSession>();
354 crate::register_default_encodings(&session);
355 crate::enable_all_registered_array_encodings(&session);
356
357 let expected = ByteBuffer::copy_from(b"cached metadata");
358 let mut output = ByteBufferMut::empty();
359 session
360 .write_options()
361 .with_metadata_segment("key", expected.clone())
362 .write(&mut output, buffer![1u32].into_array().to_array_stream())
363 .await?;
364 let fs: FileSystemRef = Arc::new(MetadataFileSystem {
365 bytes: ByteBuffer::from(output),
366 });
367 let listing = FileListing {
368 path: "metadata.vortex".to_string(),
369 size: None,
370 };
371 let default_options = |options: VortexOpenOptions| options;
372 let metadata_options = |options: VortexOpenOptions| options.include_metadata();
373
374 if default_first {
375 let default = open_file(&fs, &listing, &session, &default_options).await?;
376 assert!(default.metadata_segment("key").is_none());
377 let included = open_file(&fs, &listing, &session, &metadata_options).await?;
378 assert_eq!(
379 included.metadata_segment("key").map(ByteBuffer::as_slice),
380 Some(expected.as_slice())
381 );
382 } else {
383 let included = open_file(&fs, &listing, &session, &metadata_options).await?;
384 assert_eq!(
385 included.metadata_segment("key").map(ByteBuffer::as_slice),
386 Some(expected.as_slice())
387 );
388 let default = open_file(&fs, &listing, &session, &default_options).await?;
389 assert!(default.metadata_segment("key").is_none());
390 }
391
392 Ok(())
393 }
394
395 #[tokio::test]
396 async fn test_cached_footer_default_then_metadata_is_order_independent() -> VortexResult<()> {
397 assert_cached_footer_metadata_order(true).await
398 }
399
400 #[tokio::test]
401 async fn test_cached_footer_metadata_then_default_is_order_independent() -> VortexResult<()> {
402 assert_cached_footer_metadata_order(false).await
403 }
404
405 struct TwoFileSystem {
406 a: ByteBuffer,
407 b: ByteBuffer,
408 }
409
410 impl fmt::Debug for TwoFileSystem {
411 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
412 f.debug_struct("TwoFileSystem").finish()
413 }
414 }
415
416 #[async_trait]
417 impl FileSystem for TwoFileSystem {
418 fn list(&self, _prefix: &str) -> stream::BoxStream<'_, VortexResult<FileListing>> {
419 stream::empty().boxed()
420 }
421
422 async fn head(&self, path: &str) -> VortexResult<Option<FileListing>> {
423 let size = match path {
424 "a.vortex" => self.a.len(),
425 "b.vortex" => self.b.len(),
426 _ => return Ok(None),
427 };
428 Ok(Some(FileListing {
429 path: path.to_string(),
430 size: Some(size as u64),
431 }))
432 }
433
434 async fn open_read(&self, path: &str) -> VortexResult<Arc<dyn VortexReadAt>> {
435 let bytes = match path {
436 "b.vortex" => self.b.clone(),
437 _ => self.a.clone(),
438 };
439 Ok(Arc::new(bytes))
440 }
441
442 async fn delete(&self, _path: &str) -> VortexResult<()> {
443 Ok(())
444 }
445 }
446
447 #[tokio::test]
450 async fn test_multi_file_metadata_isolated_per_uri() -> VortexResult<()> {
451 let session = array_session()
452 .with::<LayoutSession>()
453 .with::<RuntimeSession>()
454 .with::<MultiFileSession>();
455 crate::register_default_encodings(&session);
456 crate::enable_all_registered_array_encodings(&session);
457
458 let mut out_a = ByteBufferMut::empty();
459 session
460 .write_options()
461 .with_metadata_segment("k", ByteBuffer::copy_from(b"file-a"))
462 .write(&mut out_a, buffer![1u32].into_array().to_array_stream())
463 .await?;
464 let mut out_b = ByteBufferMut::empty();
465 session
466 .write_options()
467 .with_metadata_segment("k", ByteBuffer::copy_from(b"file-b"))
468 .write(&mut out_b, buffer![2u32].into_array().to_array_stream())
469 .await?;
470
471 let fs: FileSystemRef = Arc::new(TwoFileSystem {
472 a: ByteBuffer::from(out_a),
473 b: ByteBuffer::from(out_b),
474 });
475 let include = |options: VortexOpenOptions| options.include_metadata();
476
477 let a_listing = FileListing {
478 path: "a.vortex".to_string(),
479 size: None,
480 };
481 let b_listing = FileListing {
482 path: "b.vortex".to_string(),
483 size: None,
484 };
485 let fa = open_file(&fs, &a_listing, &session, &include).await?;
486 let fb = open_file(&fs, &b_listing, &session, &include).await?;
487
488 assert_eq!(
489 fa.metadata_segment("k").map(ByteBuffer::as_slice),
490 Some(b"file-a".as_slice())
491 );
492 assert_eq!(
493 fb.metadata_segment("k").map(ByteBuffer::as_slice),
494 Some(b"file-b".as_slice())
495 );
496 Ok(())
497 }
498}