Skip to main content

vortex_file/multi/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright the Vortex contributors
3
4//! Builder for constructing a [`MultiLayoutDataSource`] from multiple Vortex files.
5
6mod 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
35/// A builder that discovers multiple Vortex files from glob patterns and constructs a
36/// [`MultiLayoutDataSource`] to scan them as a single data source.
37///
38/// The primary interface is [`Self::with_glob`], which accepts a glob pattern and an optional
39/// filesystem. For non-local filesystems (S3, GCS, etc.), callers must provide a [`FileSystemRef`].
40/// For local files, pass `None` and a local filesystem will be created automatically.
41///
42/// # Examples
43///
44/// ```ignore
45/// // Local files — filesystem is auto-created:
46/// let ds = MultiFileDataSource::new(session)
47///     .with_glob("/data/warehouse/*.vortex", None)
48///     .build()
49///     .await?;
50///
51/// // S3 — caller provides the filesystem:
52/// let ds = MultiFileDataSource::new(session)
53///     .with_glob("prefix/*.vortex", Some(s3_fs))
54///     .build()
55///     .await?;
56///
57/// // Mixed filesystems — multiple globs with different filesystems:
58/// let ds = MultiFileDataSource::new(session)
59///     .with_glob("bucket-a/*.vortex", Some(s3_fs.clone()))
60///     .with_glob("bucket-b/*.vortex", Some(s3_fs))
61///     .with_glob("gcs-bucket/*.vortex", Some(gcs_fs))
62///     .build()
63///     .await?;
64/// ```
65pub struct MultiFileDataSource {
66    session: VortexSession,
67    /// List of (glob, optional filesystem) pairs to resolve.
68    /// When the filesystem is None, a local filesystem will be created in build().
69    glob_sources: Vec<(String, Option<FileSystemRef>)>,
70    open_options_fn: Arc<dyn Fn(VortexOpenOptions) -> VortexOpenOptions + Send + Sync>,
71}
72
73/// In-flight glob resolutions in [`MultiFileDataSource::build`]. Callers like the JNI data
74/// source add one exact path per glob source, where each resolution is a single remote
75/// metadata lookup; resolving them concurrently avoids one round trip of latency per file.
76const GLOB_RESOLUTION_CONCURRENCY: usize = 16;
77
78impl MultiFileDataSource {
79    /// Create a new [`MultiFileDataSource`] builder.
80    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    /// Add a path glob for file discovery.
89    ///
90    /// The glob path should be relative to the filesystem's base URL. Pass `None` for the
91    /// filesystem to use the local filesystem (auto-created in [`Self::build`]).
92    ///
93    /// Relative paths are resolved against the process working directory.
94    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    /// Customize [`VortexOpenOptions`] applied to each file.
110    ///
111    /// Use this to configure segment caches, metrics registries, or other per-file options.
112    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    /// Build the [`DataSource`].
121    ///
122    /// Discovers files via glob, opens the first file eagerly to determine the schema,
123    /// and creates lazy factories for the remaining files.
124    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        // Create local filesystem lazily if needed (only if any glob lacks a filesystem).
130        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        // Resolve glob sources concurrently while preserving their order, since the order
140        // determines partition indices and which file is opened eagerly for the schema.
141        let resolved: Vec<Vec<(FileListing, FileSystemRef)>> =
142            stream::iter(self.glob_sources.into_iter().map(|(glob, maybe_fs)| {
143                // Use the provided filesystem, or fall back to the local filesystem.
144                // We know local_fs is Some when maybe_fs is None (by construction above).
145                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        // Open first file eagerly for dtype.
173        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/// Creates a local filesystem backed by `object_store::local::LocalFileSystem`.
206// TODO(ngates): create a native file system without an object_store dependency.
207//  Turns out it's not a trivial change because we have always used object_store with its own
208//  coalescing and concurrency configs, so we need to re-tune for local disk.
209#[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
227/// Open a single Vortex file, checking the session's footer cache.
228async 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
241/// Open a Vortex file and cache its footer on the session.
242/// Subsequent calls to this function will reuse the footer from cache.
243///
244/// "key" is the optional cache key provided by user. If it's not found,
245/// source.uri() is probed. If there's no uri(), open_cached errors.
246pub 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
277/// A [`LayoutReaderFactory`] that lazily opens a single Vortex file and returns its layout reader.
278struct 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    // Distinct files carry distinct metadata; opening both through the multi-file
448    // cache must not let one file's metadata leak into another (URI isolation).
449    #[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}