Skip to main content

object_store/
local.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! An object store implementation for a local filesystem
19use std::fs::{File, Metadata, OpenOptions, metadata, symlink_metadata};
20use std::io::{ErrorKind, Read, Seek, SeekFrom, Write};
21use std::ops::Range;
22#[cfg(target_family = "unix")]
23use std::os::unix::fs::FileExt;
24#[cfg(target_family = "windows")]
25use std::os::windows::fs::FileExt;
26use std::sync::Arc;
27use std::time::SystemTime;
28use std::{collections::BTreeSet, io};
29use std::{collections::VecDeque, path::PathBuf};
30
31use async_trait::async_trait;
32use bytes::Bytes;
33use chrono::{DateTime, Utc};
34use futures_util::{FutureExt, TryStreamExt};
35use futures_util::{StreamExt, stream::BoxStream};
36use parking_lot::Mutex;
37use url::Url;
38use walkdir::{DirEntry, WalkDir};
39
40use crate::{
41    Attributes, GetOptions, GetResult, GetResultPayload, ListResult, MultipartUpload, ObjectMeta,
42    ObjectStore, PutMode, PutMultipartOptions, PutOptions, PutPayload, PutResult, Result,
43    UploadPart, maybe_spawn_blocking,
44    path::{Path, absolute_path_to_url},
45    util::{InvalidGetRange, merge_ranges},
46};
47use crate::{CopyMode, CopyOptions, RenameOptions, RenameTargetMode};
48
49/// A specialized `Error` for filesystem object store-related errors
50#[derive(Debug, thiserror::Error)]
51pub(crate) enum Error {
52    #[error("Unable to walk dir: {}", source)]
53    UnableToWalkDir { source: walkdir::Error },
54
55    #[error("Unable to access metadata for {}: {}", path, source)]
56    Metadata {
57        source: Box<dyn std::error::Error + Send + Sync + 'static>,
58        path: String,
59    },
60
61    #[error("Unable to copy data to file: {}", source)]
62    UnableToCopyDataToFile { source: io::Error },
63
64    #[error("Unable to rename file: {}", source)]
65    UnableToRenameFile { source: io::Error },
66
67    #[error("Unable to create dir {}: {}", path.display(), source)]
68    UnableToCreateDir { source: io::Error, path: PathBuf },
69
70    #[error("Unable to create file {}: {}", path.display(), source)]
71    UnableToCreateFile { source: io::Error, path: PathBuf },
72
73    #[error("Unable to delete file {}: {}", path.display(), source)]
74    UnableToDeleteFile { source: io::Error, path: PathBuf },
75
76    #[error("Unable to open file {}: {}", path.display(), source)]
77    UnableToOpenFile { source: io::Error, path: PathBuf },
78
79    #[error("Unable to read data from file {}: {}", path.display(), source)]
80    UnableToReadBytes { source: io::Error, path: PathBuf },
81
82    #[error("Out of range of file {}, expected: {}, actual: {}", path.display(), expected, actual)]
83    OutOfRange {
84        path: PathBuf,
85        expected: u64,
86        actual: u64,
87    },
88
89    #[error("Requested range was invalid")]
90    InvalidRange { source: InvalidGetRange },
91
92    #[error("Unable to copy file from {} to {}: {}", from.display(), to.display(), source)]
93    UnableToCopyFile {
94        from: PathBuf,
95        to: PathBuf,
96        source: io::Error,
97    },
98
99    #[error("NotFound")]
100    NotFound { path: PathBuf, source: io::Error },
101
102    #[error("Error seeking file {}: {}", path.display(), source)]
103    Seek { source: io::Error, path: PathBuf },
104
105    #[error("Unable to convert URL \"{}\" to filesystem path", url)]
106    InvalidUrl { url: Url },
107
108    #[error("AlreadyExists")]
109    AlreadyExists { path: String, source: io::Error },
110
111    #[error("Unable to canonicalize filesystem root: {}", path.display())]
112    UnableToCanonicalize { path: PathBuf, source: io::Error },
113
114    #[error("Filenames containing trailing '/#\\d+/' are not supported: {}", path)]
115    InvalidPath { path: String },
116
117    #[error("Unable to sync data to disk for {}: {}", path.display(), source)]
118    UnableToSyncFile { source: io::Error, path: PathBuf },
119
120    #[error("Upload aborted")]
121    Aborted,
122}
123
124impl From<Error> for super::Error {
125    fn from(source: Error) -> Self {
126        match source {
127            Error::NotFound { path, source } => Self::NotFound {
128                path: path.to_string_lossy().to_string(),
129                source: source.into(),
130            },
131            Error::AlreadyExists { path, source } => Self::AlreadyExists {
132                path,
133                source: source.into(),
134            },
135            _ => Self::Generic {
136                store: "LocalFileSystem",
137                source: Box::new(source),
138            },
139        }
140    }
141}
142
143/// Explicitly close a file, checking for errors that would be silently ignored by Rust's `File::drop()`.
144///
145/// On network filesystems (e.g. NFS), `close()` can fail and indicate data loss.
146fn close_file(file: File) -> std::result::Result<(), io::Error> {
147    #[cfg(target_family = "unix")]
148    {
149        use std::os::fd::IntoRawFd;
150
151        close_fd(file.into_raw_fd())
152    }
153    #[cfg(target_family = "windows")]
154    {
155        use std::os::windows::io::IntoRawHandle;
156
157        close_handle(file.into_raw_handle())
158    }
159    #[cfg(not(any(target_family = "unix", target_family = "windows")))]
160    {
161        drop(file);
162        Ok(())
163    }
164}
165
166#[cfg(target_family = "unix")]
167fn close_fd(fd: std::os::fd::RawFd) -> std::result::Result<(), io::Error> {
168    nix::unistd::close(fd).map_err(|e| e.into())
169}
170
171#[cfg(target_family = "windows")]
172fn close_handle(handle: std::os::windows::io::RawHandle) -> std::result::Result<(), io::Error> {
173    // SAFETY: `handle` must be an owned handle, except when testing invalid handle errors.
174    match unsafe { windows_sys::Win32::Foundation::CloseHandle(handle) } {
175        0 => Err(io::Error::last_os_error()),
176        _ => Ok(()),
177    }
178}
179
180/// Local filesystem storage providing an [`ObjectStore`] interface to files on
181/// local disk. Can optionally be created with a directory prefix
182///
183/// # Path Semantics
184///
185/// This implementation follows the [file URI] scheme outlined in [RFC 3986]. In
186/// particular paths are delimited by `/`
187///
188/// [file URI]: https://en.wikipedia.org/wiki/File_URI_scheme
189/// [RFC 3986]: https://www.rfc-editor.org/rfc/rfc3986
190///
191/// # Path Semantics
192///
193/// [`LocalFileSystem`] will expose the path semantics of the underlying filesystem, which may
194/// have additional restrictions beyond those enforced by [`Path`].
195///
196/// For example:
197///
198/// * Windows forbids certain filenames, e.g. `COM0`,
199/// * Windows forbids folders with trailing `.`
200/// * Windows forbids certain ASCII characters, e.g. `<` or `|`
201/// * OS X forbids filenames containing `:`
202/// * Leading `-` are discouraged on Unix systems where they may be interpreted as CLI flags
203/// * Filesystems may have restrictions on the maximum path or path segment length
204/// * Filesystem support for non-ASCII characters is inconsistent
205///
206/// Additionally some filesystems, such as NTFS, are case-insensitive, whilst others like
207/// FAT don't preserve case at all. Further some filesystems support non-unicode character
208/// sequences, such as unpaired UTF-16 surrogates, and [`LocalFileSystem`] will error on
209/// encountering such sequences.
210///
211/// Finally, filenames matching the regex `/.*#\d+/`, e.g. `foo.parquet#123`, are not supported
212/// by [`LocalFileSystem`] as they are used to provide atomic writes. Such files will be ignored
213/// for listing operations, and attempting to address such a file will error.
214///
215/// # Tokio Compatibility
216///
217/// Tokio discourages performing blocking IO on a tokio worker thread, however,
218/// no major operating systems have stable async file APIs. Therefore if called from
219/// a tokio context, this will use [`tokio::runtime::Handle::spawn_blocking`] to dispatch
220/// IO to a blocking thread pool, much like `tokio::fs` does under-the-hood.
221///
222/// If not called from a tokio context, this will perform IO on the current thread with
223/// no additional complexity or overheads
224///
225/// # Symlinks
226///
227/// [`LocalFileSystem`] will follow symlinks as normal, however, it is worth noting:
228///
229/// * Broken symlinks will be silently ignored by listing operations
230/// * No effort is made to prevent breaking symlinks when deleting files
231/// * Symlinks that resolve to paths outside the root **will** be followed
232/// * Mutating a file through one or more symlinks will mutate the underlying file
233/// * Deleting a path that resolves to a symlink will only delete the symlink
234///
235/// # Cross-Filesystem Copy
236///
237/// [`LocalFileSystem::copy_opts`] is implemented using [`std::fs::hard_link`], and therefore
238/// does not support copying across filesystem boundaries.
239///
240#[derive(Clone, Debug)]
241pub struct LocalFileSystem {
242    config: Arc<Config>,
243    // if you want to delete empty directories when deleting files
244    automatic_cleanup: bool,
245    // if true, fsync written files and their parent directories after writes
246    fsync: bool,
247}
248
249#[derive(Debug)]
250struct Config {
251    root: Url,
252}
253
254impl std::fmt::Display for LocalFileSystem {
255    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
256        write!(f, "LocalFileSystem({})", self.config.root)
257    }
258}
259
260impl Default for LocalFileSystem {
261    fn default() -> Self {
262        Self::new()
263    }
264}
265
266impl LocalFileSystem {
267    /// Create new filesystem storage with no prefix
268    pub fn new() -> Self {
269        Self {
270            config: Arc::new(Config {
271                root: Url::parse("file:///").unwrap(),
272            }),
273            automatic_cleanup: false,
274            fsync: false,
275        }
276    }
277
278    /// Create new filesystem storage with `prefix` applied to all paths
279    ///
280    /// Returns an error if the path does not exist
281    ///
282    pub fn new_with_prefix(prefix: impl AsRef<std::path::Path>) -> Result<Self> {
283        let path = std::fs::canonicalize(&prefix).map_err(|source| {
284            let path = prefix.as_ref().into();
285            Error::UnableToCanonicalize { source, path }
286        })?;
287
288        Ok(Self {
289            config: Arc::new(Config {
290                root: absolute_path_to_url(path)?,
291            }),
292            automatic_cleanup: false,
293            fsync: false,
294        })
295    }
296
297    /// Return an absolute filesystem path of the given file location
298    pub fn path_to_filesystem(&self, location: &Path) -> Result<PathBuf> {
299        self.config.path_to_filesystem(location)
300    }
301
302    /// Enable automatic cleanup of empty directories when deleting files
303    pub fn with_automatic_cleanup(mut self, automatic_cleanup: bool) -> Self {
304        self.automatic_cleanup = automatic_cleanup;
305        self
306    }
307
308    /// Enable `fsync` after writes for durability
309    ///
310    /// When enabled, [`LocalFileSystem`] calls [`File::sync_all`] on written files and fsyncs
311    /// the affected parent directories before a write operation
312    /// ([`put_opts`](ObjectStore::put_opts), [`copy_opts`](ObjectStore::copy_opts),
313    /// [`rename_opts`](ObjectStore::rename_opts), and multipart upload completion) returns
314    /// success. This guarantees that both the file contents and the directory entries pointing
315    /// to them are durable on stable storage, matching the implicit durability contract of
316    /// remote object stores such as S3 or GCS.
317    ///
318    /// This trades write throughput for durability and is **disabled by default**.
319    ///
320    /// Note that directory fsync is only performed on Unix; on other platforms (e.g. Windows)
321    /// it is a no-op, as directories cannot be portably opened and synced.
322    pub fn with_fsync(mut self, fsync: bool) -> Self {
323        self.fsync = fsync;
324        self
325    }
326}
327
328impl Config {
329    /// Return an absolute filesystem path of the given location
330    fn prefix_to_filesystem(&self, location: &Path) -> Result<PathBuf> {
331        let mut url = self.root.clone();
332        url.path_segments_mut()
333            .expect("url path")
334            // technically not necessary as Path ignores empty segments
335            // but avoids creating paths with "//" which look odd in error messages.
336            .pop_if_empty()
337            .extend(location.parts());
338
339        url.to_file_path()
340            .map_err(|_| Error::InvalidUrl { url }.into())
341    }
342
343    /// Return an absolute filesystem path of the given file location
344    fn path_to_filesystem(&self, location: &Path) -> Result<PathBuf> {
345        if !is_valid_file_path(location) {
346            let path = location.as_ref().into();
347            let error = Error::InvalidPath { path };
348            return Err(error.into());
349        }
350
351        let path = self.prefix_to_filesystem(location)?;
352
353        #[cfg(target_os = "windows")]
354        let path = {
355            let path = path.to_string_lossy();
356
357            // Assume the first char is the drive letter and the next is a colon.
358            let mut out = String::new();
359            let drive = &path[..2]; // The drive letter and colon (e.g., "C:")
360            let filepath = &path[2..].replace(':', "%3A"); // Replace subsequent colons
361            out.push_str(drive);
362            out.push_str(filepath);
363            PathBuf::from(out)
364        };
365
366        Ok(path)
367    }
368
369    /// Resolves the provided absolute filesystem path to a [`Path`] prefix
370    fn filesystem_to_path(&self, location: &std::path::Path) -> Result<Path> {
371        Ok(Path::from_absolute_path_with_base(
372            location,
373            Some(&self.root),
374        )?)
375    }
376}
377
378fn is_valid_file_path(path: &Path) -> bool {
379    match path.filename() {
380        Some(p) => match p.split_once('#') {
381            Some((_, suffix)) if !suffix.is_empty() => {
382                // Valid if contains non-digits
383                !suffix.as_bytes().iter().all(|x| x.is_ascii_digit())
384            }
385            _ => true,
386        },
387        None => false,
388    }
389}
390
391#[async_trait]
392impl ObjectStore for LocalFileSystem {
393    async fn put_opts(
394        &self,
395        location: &Path,
396        payload: PutPayload,
397        opts: PutOptions,
398    ) -> Result<PutResult> {
399        if matches!(opts.mode, PutMode::Update(_)) {
400            return Err(crate::Error::NotImplemented {
401                operation: "`put_opts` with mode `PutMode::Update`".into(),
402                implementer: self.to_string(),
403            });
404        }
405
406        if !opts.attributes.is_empty() {
407            return Err(crate::Error::NotImplemented {
408                operation: "`put_opts` with `opts.attributes` specified".into(),
409                implementer: self.to_string(),
410            });
411        }
412
413        let path = self.path_to_filesystem(location)?;
414        let fsync = self.fsync;
415        maybe_spawn_blocking(move || {
416            let (mut file, staging_path) = new_staged_upload(&path, fsync)?;
417            let mut e_tag = None;
418
419            let err = match payload.iter().try_for_each(|x| file.write_all(x)) {
420                Ok(_) => {
421                    let metadata = file.metadata().map_err(|e| Error::Metadata {
422                        source: e.into(),
423                        path: path.to_string_lossy().to_string(),
424                    })?;
425                    e_tag = Some(get_etag(&metadata));
426                    // Atomically publish the staged file. When fsync is enabled the publish
427                    // helpers flush the file's contents and the destination's parent directory to
428                    // disk first, so a successful return is durable; the fsync calls are bundled
429                    // into the helpers so a file-system modification can never be left unsynced.
430                    match opts.mode {
431                        PutMode::Overwrite => {
432                            finish_staged_rename(file, &staging_path, &path, fsync).err()
433                        }
434                        PutMode::Create => {
435                            finish_staged_hard_link(file, &staging_path, &path, fsync).err()
436                        }
437                        PutMode::Update(_) => unreachable!(),
438                    }
439                }
440                Err(source) => Some(Error::UnableToCopyDataToFile { source }.into()),
441            };
442
443            if let Some(err) = err {
444                let _ = std::fs::remove_file(&staging_path); // Attempt to cleanup
445                return Err(err);
446            }
447
448            Ok(PutResult {
449                e_tag,
450                version: None,
451                extensions: Default::default(),
452            })
453        })
454        .await
455    }
456
457    async fn put_multipart_opts(
458        &self,
459        location: &Path,
460        opts: PutMultipartOptions,
461    ) -> Result<Box<dyn MultipartUpload>> {
462        if !opts.attributes.is_empty() {
463            return Err(crate::Error::NotImplemented {
464                operation: "`put_multipart_opts` with `opts.attributes` specified".into(),
465                implementer: self.to_string(),
466            });
467        }
468
469        let dest = self.path_to_filesystem(location)?;
470        let (file, src) = new_staged_upload(&dest, self.fsync)?;
471        Ok(Box::new(LocalUpload::new(src, dest, file, self.fsync)))
472    }
473
474    async fn get_opts(&self, location: &Path, options: GetOptions) -> Result<GetResult> {
475        let location = location.clone();
476        let path = self.path_to_filesystem(&location)?;
477        maybe_spawn_blocking(move || {
478            let file = open_file(&path)?;
479            let metadata = open_metadata(&file, &path)?;
480            let meta = convert_metadata(metadata, location);
481            options.check_preconditions(&meta)?;
482
483            let range = match options.range {
484                Some(r) => r
485                    .as_range(meta.size)
486                    .map_err(|source| Error::InvalidRange { source })?,
487                None => 0..meta.size,
488            };
489
490            Ok(GetResult {
491                payload: GetResultPayload::File(file, path),
492                attributes: Attributes::default(),
493                range,
494                meta,
495                extensions: Default::default(),
496            })
497        })
498        .await
499    }
500
501    async fn get_ranges(&self, location: &Path, ranges: &[Range<u64>]) -> Result<Vec<Bytes>> {
502        let path = self.path_to_filesystem(location)?;
503        let ranges = ranges.to_vec();
504        maybe_spawn_blocking(move || {
505            // We do not read the metadata here, but error in `read_range` if necessary
506            let mut file = File::open(&path).map_err(|e| map_open_error(e, &path))?;
507
508            // Coalesce only contiguous/overlapping ranges (gap 0), so we never
509            // read bytes that weren't requested while still collapsing adjacent
510            // ranges into fewer reads.
511            let fetch_ranges = merge_ranges(&ranges, 0);
512
513            // Fast path: nothing coalesced, read each range directly in order.
514            if fetch_ranges.len() == ranges.len() {
515                return ranges
516                    .iter()
517                    .map(|r| read_range(&mut file, &path, r.clone()))
518                    .collect();
519            }
520
521            let fetched = fetch_ranges
522                .iter()
523                .map(|r| read_range(&mut file, &path, r.clone()))
524                .collect::<Result<Vec<_>>>()?;
525
526            Ok(ranges
527                .iter()
528                .map(|range| {
529                    let idx = fetch_ranges.partition_point(|v| v.start <= range.start) - 1;
530                    let fetch_range = &fetch_ranges[idx];
531                    let fetch_bytes = &fetched[idx];
532
533                    let start = (range.start - fetch_range.start) as usize;
534                    let end = (range.end - fetch_range.start) as usize;
535                    fetch_bytes.slice(start..end.min(fetch_bytes.len()))
536                })
537                .collect())
538        })
539        .await
540    }
541
542    fn delete_stream(
543        &self,
544        locations: BoxStream<'static, Result<Path>>,
545    ) -> BoxStream<'static, Result<Path>> {
546        let config = Arc::clone(&self.config);
547        let automatic_cleanup = self.automatic_cleanup;
548        locations
549            .map(move |location| {
550                let config = Arc::clone(&config);
551                maybe_spawn_blocking(move || {
552                    let location = location?;
553                    // `with_fsync` does not apply to standalone deletes; only create-mode rename
554                    // fsyncs its internal source removal as part of the durable copy-and-delete.
555                    Self::delete_location(config, automatic_cleanup, &location, false)?;
556                    Ok(location)
557                })
558            })
559            .buffered(10)
560            .boxed()
561    }
562
563    fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>> {
564        Self::list_with_maybe_offset(Arc::clone(&self.config), prefix, None)
565    }
566
567    fn list_with_offset(
568        &self,
569        prefix: Option<&Path>,
570        offset: &Path,
571    ) -> BoxStream<'static, Result<ObjectMeta>> {
572        Self::list_with_maybe_offset(Arc::clone(&self.config), prefix, Some(offset))
573    }
574
575    async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result<ListResult> {
576        let config = Arc::clone(&self.config);
577
578        let prefix = prefix.cloned().unwrap_or_default();
579        let resolved_prefix = config.prefix_to_filesystem(&prefix)?;
580
581        maybe_spawn_blocking(move || {
582            let walkdir = WalkDir::new(&resolved_prefix)
583                .min_depth(1)
584                .max_depth(1)
585                .follow_links(true);
586
587            let mut common_prefixes = BTreeSet::new();
588            let mut objects = Vec::new();
589
590            for entry_res in walkdir.into_iter().map(convert_walkdir_result) {
591                if let Some(entry) = entry_res? {
592                    let is_directory = entry.file_type().is_dir();
593                    let entry_location = config.filesystem_to_path(entry.path())?;
594                    if !is_directory && !is_valid_file_path(&entry_location) {
595                        continue;
596                    }
597
598                    let mut parts = match entry_location.prefix_match(&prefix) {
599                        Some(parts) => parts,
600                        None => continue,
601                    };
602
603                    let common_prefix = match parts.next() {
604                        Some(p) => p,
605                        None => continue,
606                    };
607
608                    drop(parts);
609
610                    if is_directory {
611                        common_prefixes.insert(prefix.clone().join(common_prefix));
612                    } else if let Some(metadata) = convert_entry(entry, entry_location)? {
613                        objects.push(metadata);
614                    }
615                }
616            }
617
618            Ok(ListResult {
619                common_prefixes: common_prefixes.into_iter().collect(),
620                objects,
621                extensions: Default::default(),
622            })
623        })
624        .await
625    }
626
627    async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> Result<()> {
628        let CopyOptions {
629            mode,
630            extensions: _,
631        } = options;
632
633        let from = self.path_to_filesystem(from)?;
634        let to = self.path_to_filesystem(to)?;
635        let fsync = self.fsync;
636
637        match mode {
638            CopyMode::Overwrite => {
639                let mut id = 0;
640                // In order to make this atomic we:
641                //
642                // - hard link to a hidden temporary file
643                // - atomically rename this temporary file into place
644                //
645                // This is necessary because hard_link returns an error if the destination already exists
646                maybe_spawn_blocking(move || {
647                    loop {
648                        let staged = staged_upload_path(&to, &id.to_string());
649                        // Stage via a temporary hard link; the source is already durable so the
650                        // staging link itself needs no fsync (the publish rename below fsyncs the
651                        // shared parent directory).
652                        match std::fs::hard_link(&from, &staged) {
653                            // `rename` bundles in the fsync of `to`'s parent directory.
654                            Ok(_) => match rename(&staged, &to, fsync) {
655                                Ok(_) => return Ok(()),
656                                Err(source) => {
657                                    let _ = std::fs::remove_file(&staged); // Attempt to clean up
658                                    return Err(Error::UnableToCopyFile { from, to, source }.into());
659                                }
660                            },
661                            Err(source) => match source.kind() {
662                                ErrorKind::AlreadyExists => id += 1,
663                                ErrorKind::NotFound => match from.exists() {
664                                    true => create_parent_dirs(&to, source, fsync)?,
665                                    false => {
666                                        return Err(Error::NotFound { path: from, source }.into());
667                                    }
668                                },
669                                _ => {
670                                    return Err(Error::UnableToCopyFile { from, to, source }.into());
671                                }
672                            },
673                        }
674                    }
675                })
676                .await
677            }
678            CopyMode::Create => {
679                maybe_spawn_blocking(move || {
680                    loop {
681                        // The source is an existing object that is already durable, so no file
682                        // sync is needed; `hard_link` bundles in the fsync of `to`'s parent dir.
683                        match hard_link(&from, &to, fsync) {
684                            Ok(_) => return Ok(()),
685                            Err(source) => match source.kind() {
686                                ErrorKind::AlreadyExists => {
687                                    return Err(Error::AlreadyExists {
688                                        path: to.to_str().unwrap().to_string(),
689                                        source,
690                                    }
691                                    .into());
692                                }
693                                ErrorKind::NotFound => match from.exists() {
694                                    true => create_parent_dirs(&to, source, fsync)?,
695                                    false => {
696                                        return Err(Error::NotFound { path: from, source }.into());
697                                    }
698                                },
699                                _ => {
700                                    return Err(Error::UnableToCopyFile { from, to, source }.into());
701                                }
702                            },
703                        }
704                    }
705                })
706                .await
707            }
708        }
709    }
710
711    async fn rename_opts(&self, from: &Path, to: &Path, options: RenameOptions) -> Result<()> {
712        let RenameOptions {
713            target_mode,
714            extensions,
715        } = options;
716
717        match target_mode {
718            // optimized implementation
719            RenameTargetMode::Overwrite => {
720                let from = self.path_to_filesystem(from)?;
721                let to = self.path_to_filesystem(to)?;
722                let fsync = self.fsync;
723                maybe_spawn_blocking(move || {
724                    loop {
725                        // Unlike multipart `complete`, there is no freshly written file to
726                        // `sync_all` here: `from` is an existing, already-durable object and a
727                        // rename only mutates directory entries. `rename` bundles in the fsync of
728                        // both affected directories (destination, and source if it differs).
729                        match rename(&from, &to, fsync) {
730                            Ok(_) => return Ok(()),
731                            Err(source) => match source.kind() {
732                                ErrorKind::NotFound => match from.exists() {
733                                    true => create_parent_dirs(&to, source, fsync)?,
734                                    false => {
735                                        return Err(Error::NotFound { path: from, source }.into());
736                                    }
737                                },
738                                _ => {
739                                    return Err(Error::UnableToCopyFile { from, to, source }.into());
740                                }
741                            },
742                        }
743                    }
744                })
745                .await
746            }
747            // fall-back to copy & delete
748            RenameTargetMode::Create => {
749                self.copy_opts(
750                    from,
751                    to,
752                    CopyOptions {
753                        mode: CopyMode::Create,
754                        extensions,
755                    },
756                )
757                .await?;
758                let config = Arc::clone(&self.config);
759                let automatic_cleanup = self.automatic_cleanup;
760                let fsync = self.fsync;
761                let from = from.clone();
762                maybe_spawn_blocking(move || {
763                    Self::delete_location(config, automatic_cleanup, &from, fsync)
764                })
765                .await?;
766                Ok(())
767            }
768        }
769    }
770}
771
772impl LocalFileSystem {
773    fn delete_location(
774        config: Arc<Config>,
775        automatic_cleanup: bool,
776        location: &Path,
777        fsync: bool,
778    ) -> Result<()> {
779        let path = config.path_to_filesystem(location)?;
780        if let Err(e) = std::fs::remove_file(&path) {
781            Err(match e.kind() {
782                ErrorKind::NotFound => Error::NotFound { path, source: e }.into(),
783                _ => Error::UnableToDeleteFile { path, source: e }.into(),
784            })
785        } else {
786            if fsync {
787                fsync_parent_dir(&path).map_err(|source| Error::UnableToSyncFile {
788                    source,
789                    path: path.clone(),
790                })?;
791            }
792
793            if !automatic_cleanup {
794                return Ok(());
795            }
796
797            let root = &config.root;
798            let root = root
799                .to_file_path()
800                .map_err(|_| Error::InvalidUrl { url: root.clone() })?;
801
802            // here we will try to traverse up and delete an empty dir if possible until we reach the root or get an error
803            let mut parent = path.parent();
804
805            while let Some(loc) = parent {
806                if loc != root && std::fs::remove_dir(loc).is_ok() {
807                    parent = loc.parent();
808                } else {
809                    break;
810                }
811            }
812
813            Ok(())
814        }
815    }
816
817    fn list_with_maybe_offset(
818        config: Arc<Config>,
819        prefix: Option<&Path>,
820        maybe_offset: Option<&Path>,
821    ) -> BoxStream<'static, Result<ObjectMeta>> {
822        let root_path = match prefix {
823            Some(prefix) => match config.prefix_to_filesystem(prefix) {
824                Ok(path) => path,
825                Err(e) => return futures_util::future::ready(Err(e)).into_stream().boxed(),
826            },
827            None => config.root.to_file_path().unwrap(),
828        };
829
830        let walkdir = WalkDir::new(root_path)
831            // Don't include the root directory itself
832            .min_depth(1)
833            .follow_links(true);
834
835        let maybe_offset = maybe_offset.cloned();
836
837        let s = walkdir.into_iter().flat_map(move |result_dir_entry| {
838            // Apply offset filter before proceeding, to reduce statx file system calls
839            // This matters for NFS mounts
840            if let (Some(offset), Ok(entry)) = (maybe_offset.as_ref(), result_dir_entry.as_ref()) {
841                let location = config.filesystem_to_path(entry.path());
842                match location {
843                    Ok(path) if path <= *offset => return None,
844                    Err(e) => return Some(Err(e)),
845                    _ => {}
846                }
847            }
848
849            let entry = match convert_walkdir_result(result_dir_entry).transpose()? {
850                Ok(entry) => entry,
851                Err(e) => return Some(Err(e)),
852            };
853
854            if !entry.path().is_file() {
855                return None;
856            }
857
858            match config.filesystem_to_path(entry.path()) {
859                Ok(path) => match is_valid_file_path(&path) {
860                    true => convert_entry(entry, path).transpose(),
861                    false => None,
862                },
863                Err(e) => Some(Err(e)),
864            }
865        });
866
867        // If no tokio context, return iterator directly as no
868        // need to perform chunked spawn_blocking reads
869        if tokio::runtime::Handle::try_current().is_err() {
870            return futures_util::stream::iter(s).boxed();
871        }
872
873        // Otherwise list in batches of CHUNK_SIZE
874        const CHUNK_SIZE: usize = 1024;
875
876        let buffer = VecDeque::with_capacity(CHUNK_SIZE);
877        futures_util::stream::try_unfold((s, buffer), |(mut s, mut buffer)| async move {
878            if buffer.is_empty() {
879                (s, buffer) = tokio::task::spawn_blocking(move || {
880                    for _ in 0..CHUNK_SIZE {
881                        match s.next() {
882                            Some(r) => buffer.push_back(r),
883                            None => break,
884                        }
885                    }
886                    (s, buffer)
887                })
888                .await?;
889            }
890
891            match buffer.pop_front() {
892                Some(Err(e)) => Err(e),
893                Some(Ok(meta)) => Ok(Some((meta, (s, buffer)))),
894                None => Ok(None),
895            }
896        })
897        .boxed()
898    }
899}
900
901/// Creates the parent directories of `path` or returns an error based on `source` if no parent
902///
903/// When `fsync` is true, every directory created here is fsynced, up to and including the first
904/// pre-existing ancestor (whose entry list also changed), so the new directory entries are durable.
905fn create_parent_dirs(path: &std::path::Path, source: io::Error, fsync: bool) -> Result<()> {
906    let parent = path.parent().ok_or_else(|| {
907        let path = path.to_path_buf();
908        Error::UnableToCreateFile { path, source }
909    })?;
910
911    // Record the deepest already-existing ancestor *before* creating any directories, so that
912    // afterwards we know exactly which directories are new and need to be fsynced.
913    let first_existing = fsync.then(|| {
914        let mut dir = parent;
915        while !dir.exists() {
916            match dir.parent() {
917                Some(p) => dir = p,
918                None => break,
919            }
920        }
921        dir.to_path_buf()
922    });
923
924    std::fs::create_dir_all(parent).map_err(|source| {
925        let path = parent.into();
926        Error::UnableToCreateDir { source, path }
927    })?;
928
929    if let Some(first_existing) = first_existing {
930        // Walk from `parent` up to `first_existing`, fsyncing each directory whose entries changed.
931        let mut dir = parent;
932        loop {
933            fsync_dir(dir).map_err(|source| Error::UnableToSyncFile {
934                source,
935                path: dir.into(),
936            })?;
937            if dir == first_existing {
938                break;
939            }
940            dir = match dir.parent() {
941                Some(p) => p,
942                None => break,
943            };
944        }
945    }
946    Ok(())
947}
948
949/// Renames `from` to `to`, then — when `fsync` is enabled — fsyncs the destination's parent
950/// directory (and the source's too, if it differs) so the moved directory entries are durable.
951///
952/// The directory fsync is bundled in deliberately: every durable rename goes through here, so the
953/// post-rename fsync can never be forgotten at an individual call site.
954fn rename(from: &std::path::Path, to: &std::path::Path, fsync: bool) -> io::Result<()> {
955    std::fs::rename(from, to)?;
956    if fsync {
957        fsync_parent_dir(to)?;
958        // A cross-directory move also removes an entry from the source directory.
959        if from.parent() != to.parent() {
960            fsync_parent_dir(from)?;
961        }
962    }
963    Ok(())
964}
965
966/// Hard-links `original` to `link`, then — when `fsync` is enabled — fsyncs `link`'s parent
967/// directory so the new directory entry is durable.
968///
969/// As with [`rename`], the directory fsync is bundled in so it cannot be forgotten at a call site.
970fn hard_link(original: &std::path::Path, link: &std::path::Path, fsync: bool) -> io::Result<()> {
971    std::fs::hard_link(original, link)?;
972    if fsync {
973        fsync_parent_dir(link)?;
974    }
975    Ok(())
976}
977
978/// Durably publishes the freshly-written staging file `file` (located at `src`) to `dest` via a
979/// rename.
980///
981/// When `fsync` is enabled, the file's contents are flushed before — and `dest`'s parent
982/// directory after — the rename, so a successful return is durable. The file is always closed
983/// before the rename (checking for close errors): required for NFS error detection and for some
984/// FUSE filesystems (e.g. Blobfuse) that only commit the data on close.
985fn finish_staged_rename(
986    file: File,
987    src: &std::path::Path,
988    dest: &std::path::Path,
989    fsync: bool,
990) -> Result<()> {
991    sync_and_close(file, src, fsync)?;
992    rename(src, dest, fsync).map_err(|source| Error::UnableToRenameFile { source })?;
993    Ok(())
994}
995
996/// Like [`finish_staged_rename`] but publishes via a hard link (`PutMode::Create` semantics): the
997/// staging file is linked to `dest` and then removed. Returns [`Error::AlreadyExists`] if `dest`
998/// already exists.
999fn finish_staged_hard_link(
1000    file: File,
1001    src: &std::path::Path,
1002    dest: &std::path::Path,
1003    fsync: bool,
1004) -> Result<()> {
1005    sync_and_close(file, src, fsync)?;
1006    match hard_link(src, dest, fsync) {
1007        Ok(()) => {
1008            let _ = std::fs::remove_file(src); // Attempt to cleanup
1009            Ok(())
1010        }
1011        Err(source) => match source.kind() {
1012            ErrorKind::AlreadyExists => Err(Error::AlreadyExists {
1013                path: dest.to_str().unwrap().to_string(),
1014                source,
1015            }
1016            .into()),
1017            _ => Err(Error::UnableToRenameFile { source }.into()),
1018        },
1019    }
1020}
1021
1022/// Flushes the freshly-written `file`'s contents to disk (when `fsync` is enabled) and then
1023/// closes it, checking for close errors that dropping the [`File`] would silently ignore.
1024fn sync_and_close(file: File, path: &std::path::Path, fsync: bool) -> Result<()> {
1025    if fsync {
1026        file.sync_all().map_err(|source| Error::UnableToSyncFile {
1027            source,
1028            path: path.into(),
1029        })?;
1030    }
1031    close_file(file).map_err(|source| Error::UnableToCopyDataToFile { source })?;
1032    Ok(())
1033}
1034
1035/// Fsyncs the parent directory of `path` so a change to its directory entries (e.g. one just made
1036/// by a rename or hard link) is durable. A no-op when `path` has no parent.
1037fn fsync_parent_dir(path: &std::path::Path) -> io::Result<()> {
1038    match path.parent() {
1039        Some(parent) => fsync_dir(parent),
1040        None => Ok(()),
1041    }
1042}
1043
1044/// Fsyncs `dir_path` so that changes to its directory entries are durable.
1045///
1046/// This is only meaningful on Unix; on other platforms (e.g. Windows) directories cannot be
1047/// portably opened as a [`File`] and synced, so this is a no-op.
1048fn fsync_dir(dir_path: &std::path::Path) -> io::Result<()> {
1049    #[cfg(target_family = "unix")]
1050    {
1051        File::open(dir_path)?.sync_all()
1052    }
1053    #[cfg(not(target_family = "unix"))]
1054    {
1055        let _ = dir_path;
1056        Ok(())
1057    }
1058}
1059
1060/// Generates a unique file path `{base}#{suffix}`, returning the opened `File` and `path`
1061///
1062/// Creates any directories if necessary, fsyncing them when `fsync` is enabled
1063fn new_staged_upload(base: &std::path::Path, fsync: bool) -> Result<(File, PathBuf)> {
1064    let mut multipart_id = 1;
1065    loop {
1066        let suffix = multipart_id.to_string();
1067        let path = staged_upload_path(base, &suffix);
1068        let mut options = OpenOptions::new();
1069        match options.read(true).write(true).create_new(true).open(&path) {
1070            Ok(f) => return Ok((f, path)),
1071            Err(source) => match source.kind() {
1072                ErrorKind::AlreadyExists => multipart_id += 1,
1073                ErrorKind::NotFound => create_parent_dirs(&path, source, fsync)?,
1074                _ => return Err(Error::UnableToOpenFile { source, path }.into()),
1075            },
1076        }
1077    }
1078}
1079
1080/// Returns the unique upload for the given path and suffix
1081fn staged_upload_path(dest: &std::path::Path, suffix: &str) -> PathBuf {
1082    let mut staging_path = dest.as_os_str().to_owned();
1083    staging_path.push("#");
1084    staging_path.push(suffix);
1085    staging_path.into()
1086}
1087
1088#[derive(Debug)]
1089struct LocalUpload {
1090    /// The upload state
1091    state: Arc<UploadState>,
1092    /// The location of the temporary file
1093    src: Option<PathBuf>,
1094    /// The next offset to write into the file
1095    offset: u64,
1096    /// Whether to fsync the file and its parent directory on completion
1097    fsync: bool,
1098}
1099
1100#[derive(Debug)]
1101struct UploadState {
1102    dest: PathBuf,
1103    file: Mutex<Option<File>>,
1104}
1105
1106impl LocalUpload {
1107    pub(crate) fn new(src: PathBuf, dest: PathBuf, file: File, fsync: bool) -> Self {
1108        Self {
1109            state: Arc::new(UploadState {
1110                dest,
1111                file: Mutex::new(Some(file)),
1112            }),
1113            src: Some(src),
1114            offset: 0,
1115            fsync,
1116        }
1117    }
1118}
1119
1120#[async_trait]
1121impl MultipartUpload for LocalUpload {
1122    fn put_part(&mut self, data: PutPayload) -> UploadPart {
1123        let offset = self.offset;
1124        self.offset += data.content_length() as u64;
1125
1126        let s = Arc::clone(&self.state);
1127        maybe_spawn_blocking(move || {
1128            let mut guard = s.file.lock();
1129            let file = guard.as_mut().ok_or(Error::Aborted)?;
1130            file.seek(SeekFrom::Start(offset)).map_err(|source| {
1131                let path = s.dest.clone();
1132                Error::Seek { source, path }
1133            })?;
1134
1135            data.iter()
1136                .try_for_each(|x| file.write_all(x))
1137                .map_err(|source| Error::UnableToCopyDataToFile { source })?;
1138
1139            Ok(())
1140        })
1141        .boxed()
1142    }
1143
1144    async fn complete(&mut self) -> Result<PutResult> {
1145        let src = self.src.take().ok_or(Error::Aborted)?;
1146        let s = Arc::clone(&self.state);
1147        let fsync = self.fsync;
1148        maybe_spawn_blocking(move || {
1149            // Ensure no inflight writes
1150            let mut guard = s.file.lock();
1151            let file = guard.take().ok_or(Error::Aborted)?;
1152
1153            let metadata = file.metadata().map_err(|e| Error::Metadata {
1154                source: e.into(),
1155                path: src.to_string_lossy().to_string(),
1156            })?;
1157
1158            // Durably publish the freshly-written staging file: flush its contents, close it, then
1159            // rename it into place and fsync the destination's parent directory (the fsync calls
1160            // are bundled into the helper and only run when fsync is enabled).
1161            finish_staged_rename(file, &src, &s.dest, fsync)?;
1162
1163            Ok(PutResult {
1164                e_tag: Some(get_etag(&metadata)),
1165                version: None,
1166                extensions: Default::default(),
1167            })
1168        })
1169        .await
1170    }
1171
1172    async fn abort(&mut self) -> Result<()> {
1173        let src = self.src.take().ok_or(Error::Aborted)?;
1174        maybe_spawn_blocking(move || {
1175            std::fs::remove_file(&src)
1176                .map_err(|source| Error::UnableToDeleteFile { source, path: src })?;
1177            Ok(())
1178        })
1179        .await
1180    }
1181}
1182
1183impl Drop for LocalUpload {
1184    fn drop(&mut self) {
1185        if let Some(src) = self.src.take() {
1186            // Try to clean up intermediate file ignoring any error
1187            match tokio::runtime::Handle::try_current() {
1188                Ok(r) => drop(r.spawn_blocking(move || std::fs::remove_file(src))),
1189                Err(_) => drop(std::fs::remove_file(src)),
1190            };
1191        }
1192    }
1193}
1194
1195pub(crate) fn chunked_stream(
1196    mut file: File,
1197    path: PathBuf,
1198    range: Range<u64>,
1199    chunk_size: usize,
1200) -> BoxStream<'static, Result<Bytes, super::Error>> {
1201    futures_util::stream::once(async move {
1202        let requested = range.end - range.start;
1203
1204        let (file, path) = maybe_spawn_blocking(move || {
1205            file.seek(SeekFrom::Start(range.start as _))
1206                .map_err(|err| map_seek_error(err, &file, &path, range.start))?;
1207            Ok((file, path))
1208        })
1209        .await?;
1210
1211        let stream = futures_util::stream::try_unfold(
1212            (file, path, requested),
1213            move |(mut file, path, remaining)| {
1214                maybe_spawn_blocking(move || {
1215                    if remaining == 0 {
1216                        return Ok(None);
1217                    }
1218
1219                    let to_read = remaining.min(chunk_size as u64);
1220                    let cap = usize::try_from(to_read).map_err(|_e| Error::InvalidRange {
1221                        source: InvalidGetRange::TooLarge {
1222                            requested: to_read,
1223                            max: usize::MAX as u64,
1224                        },
1225                    })?;
1226                    let mut buffer = Vec::with_capacity(cap);
1227                    let read = (&mut file)
1228                        .take(to_read)
1229                        .read_to_end(&mut buffer)
1230                        .map_err(|e| Error::UnableToReadBytes {
1231                            source: e,
1232                            path: path.clone(),
1233                        })?;
1234
1235                    Ok(Some((buffer.into(), (file, path, remaining - read as u64))))
1236                })
1237            },
1238        );
1239        Ok::<_, super::Error>(stream)
1240    })
1241    .try_flatten()
1242    .boxed()
1243}
1244
1245pub(crate) fn read_range(
1246    file: &mut File,
1247    path: &std::path::Path,
1248    range: Range<u64>,
1249) -> Result<Bytes> {
1250    let requested = range.end - range.start;
1251
1252    let mut buf = Vec::with_capacity(requested as usize);
1253
1254    #[cfg(any(target_family = "unix", target_family = "windows"))]
1255    {
1256        buf.resize(requested as usize, 0_u8);
1257
1258        let mut buf_slice = &mut buf[..];
1259        let mut offset = range.start;
1260
1261        while !buf_slice.is_empty() {
1262            #[cfg(target_family = "unix")]
1263            let read_result = file.read_at(buf_slice, offset);
1264
1265            #[cfg(target_family = "windows")]
1266            let read_result = file.seek_read(buf_slice, offset);
1267
1268            match read_result {
1269                Ok(0) => break,
1270                Ok(n) => {
1271                    let tmp = buf_slice;
1272                    buf_slice = &mut tmp[n..];
1273                    offset += n as u64;
1274                }
1275                // This error is recoverable
1276                Err(e) if e.kind() == ErrorKind::Interrupted => {}
1277                Err(source) => {
1278                    let error = Error::UnableToReadBytes {
1279                        source,
1280                        path: path.into(),
1281                    };
1282
1283                    return Err(error.into());
1284                }
1285            }
1286        }
1287
1288        // If we reached EOF before filling the buffer
1289        if !buf_slice.is_empty() {
1290            let metadata = open_metadata(file, path)?;
1291            let file_len = metadata.len();
1292
1293            // If none of the range is satisfiable we should error, e.g. if the start offset is beyond the
1294            // extents of the file, or if its at the end of the file and wants to read a non-empty range.
1295            // if range.start > file_len || (range.start == file_len && !range.is_empty()) {
1296            if range.start >= file_len {
1297                return Err(Error::InvalidRange {
1298                    source: InvalidGetRange::StartTooLarge {
1299                        requested: range.start,
1300                        length: file_len,
1301                    },
1302                }
1303                .into());
1304            }
1305
1306            let expected = range.end.min(file_len) - range.start;
1307
1308            let error = Error::OutOfRange {
1309                path: path.into(),
1310                expected,
1311                actual: offset - range.start,
1312            };
1313
1314            return Err(error.into());
1315        }
1316    }
1317    #[cfg(all(not(windows), not(unix)))]
1318    {
1319        file.seek(SeekFrom::Start(range.start))
1320            .map_err(|err| map_seek_error(err, file, path, range.start))?;
1321
1322        let read = file.take(requested).read_to_end(&mut buf).map_err(|err| {
1323            // try to read metadata to give a better error in case of directory
1324            if let Err(e) = open_metadata(file, path) {
1325                return e;
1326            }
1327            Error::UnableToReadBytes {
1328                source: err,
1329                path: path.to_path_buf(),
1330            }
1331        })? as u64;
1332
1333        if read != requested {
1334            let metadata = open_metadata(file, path)?;
1335            let file_len = metadata.len();
1336
1337            if range.start >= file_len {
1338                return Err(Error::InvalidRange {
1339                    source: InvalidGetRange::StartTooLarge {
1340                        requested: range.start,
1341                        length: file_len,
1342                    },
1343                }
1344                .into());
1345            }
1346
1347            let expected = range.end.min(file_len) - range.start;
1348            if read != expected {
1349                return Err(Error::OutOfRange {
1350                    path: path.to_path_buf(),
1351                    expected,
1352                    actual: read,
1353                }
1354                .into());
1355            }
1356        }
1357    }
1358
1359    Ok(buf.into())
1360}
1361
1362fn open_file(path: &std::path::Path) -> Result<File, Error> {
1363    File::open(path).map_err(|e| map_open_error(e, path))
1364}
1365
1366fn open_metadata(file: &File, path: &std::path::Path) -> Result<Metadata, Error> {
1367    let metadata = file.metadata().map_err(|e| map_open_error(e, path))?;
1368    if metadata.is_dir() {
1369        Err(Error::NotFound {
1370            path: PathBuf::from(path),
1371            source: io::Error::new(ErrorKind::NotFound, "is directory"),
1372        })
1373    } else {
1374        Ok(metadata)
1375    }
1376}
1377
1378/// Translates errors from opening a file into a more specific [`Error`] when possible
1379fn map_open_error(source: io::Error, path: &std::path::Path) -> Error {
1380    let path = PathBuf::from(path);
1381    match source.kind() {
1382        ErrorKind::NotFound => Error::NotFound { path, source },
1383        _ => Error::UnableToOpenFile { path, source },
1384    }
1385}
1386
1387/// Translates errors from attempting to a file into a more specific [`Error`] when possible
1388fn map_seek_error(source: io::Error, file: &File, path: &std::path::Path, requested: u64) -> Error {
1389    // if we can't seek, check if start is out of bounds to give
1390    // a better error. Don't read metadata before to avoid
1391    // an extra syscall in the common case
1392    let m = match open_metadata(file, path) {
1393        Err(e) => return e,
1394        Ok(m) => m,
1395    };
1396    if requested >= m.len() {
1397        return Error::InvalidRange {
1398            source: InvalidGetRange::StartTooLarge {
1399                requested,
1400                length: m.len(),
1401            },
1402        };
1403    }
1404    Error::Seek {
1405        source,
1406        path: PathBuf::from(path),
1407    }
1408}
1409
1410fn convert_entry(entry: DirEntry, location: Path) -> Result<Option<ObjectMeta>> {
1411    match entry.metadata() {
1412        Ok(metadata) => Ok(Some(convert_metadata(metadata, location))),
1413        Err(e) => {
1414            if let Some(io_err) = e.io_error() {
1415                if io_err.kind() == ErrorKind::NotFound {
1416                    return Ok(None);
1417                }
1418            }
1419            Err(Error::Metadata {
1420                source: e.into(),
1421                path: location.to_string(),
1422            })?
1423        }
1424    }
1425}
1426
1427fn last_modified(metadata: &Metadata) -> DateTime<Utc> {
1428    metadata
1429        .modified()
1430        .expect("Modified file time should be supported on this platform")
1431        .into()
1432}
1433
1434fn get_etag(metadata: &Metadata) -> String {
1435    let inode = get_inode(metadata);
1436    let size = metadata.len();
1437    let mtime = metadata
1438        .modified()
1439        .ok()
1440        .and_then(|mtime| mtime.duration_since(SystemTime::UNIX_EPOCH).ok())
1441        .unwrap_or_default()
1442        .as_micros();
1443
1444    // Use an ETag scheme based on that used by many popular HTTP servers
1445    // <https://httpd.apache.org/docs/2.2/mod/core.html#fileetag>
1446    // <https://stackoverflow.com/questions/47512043/how-etags-are-generated-and-configured>
1447    format!("\"{inode:x}-{mtime:x}-{size:x}\"")
1448}
1449
1450fn convert_metadata(metadata: Metadata, location: Path) -> ObjectMeta {
1451    let last_modified = last_modified(&metadata);
1452
1453    ObjectMeta {
1454        location,
1455        last_modified,
1456        size: metadata.len(),
1457        e_tag: Some(get_etag(&metadata)),
1458        version: None,
1459    }
1460}
1461
1462#[cfg(unix)]
1463/// We include the inode when available to yield an ETag more resistant to collisions
1464/// and as used by popular web servers such as [Apache](https://httpd.apache.org/docs/2.2/mod/core.html#fileetag)
1465fn get_inode(metadata: &Metadata) -> u64 {
1466    std::os::unix::fs::MetadataExt::ino(metadata)
1467}
1468
1469#[cfg(not(unix))]
1470/// On platforms where an inode isn't available, fallback to just relying on size and mtime
1471fn get_inode(_metadata: &Metadata) -> u64 {
1472    0
1473}
1474
1475/// Convert walkdir results and converts not-found errors into `None`.
1476/// Convert broken symlinks to `None`.
1477fn convert_walkdir_result(
1478    res: std::result::Result<DirEntry, walkdir::Error>,
1479) -> Result<Option<DirEntry>> {
1480    match res {
1481        Ok(entry) => {
1482            // To check for broken symlink: call symlink_metadata() - it does not traverse symlinks);
1483            // if ok: check if entry is symlink; and try to read it by calling metadata().
1484            match symlink_metadata(entry.path()) {
1485                Ok(attr) => {
1486                    if attr.is_symlink() {
1487                        let target_metadata = metadata(entry.path());
1488                        match target_metadata {
1489                            Ok(_) => {
1490                                // symlink is valid
1491                                Ok(Some(entry))
1492                            }
1493                            Err(_) => {
1494                                // this is a broken symlink, return None
1495                                Ok(None)
1496                            }
1497                        }
1498                    } else {
1499                        Ok(Some(entry))
1500                    }
1501                }
1502                Err(_) => Ok(None),
1503            }
1504        }
1505
1506        Err(walkdir_err) => match walkdir_err.io_error() {
1507            Some(io_err) => match io_err.kind() {
1508                ErrorKind::NotFound => Ok(None),
1509                _ => Err(Error::UnableToWalkDir {
1510                    source: walkdir_err,
1511                }
1512                .into()),
1513            },
1514            None => Err(Error::UnableToWalkDir {
1515                source: walkdir_err,
1516            }
1517            .into()),
1518        },
1519    }
1520}
1521
1522#[cfg(test)]
1523mod tests {
1524    use std::fs;
1525
1526    use futures_util::TryStreamExt;
1527    use tempfile::TempDir;
1528
1529    #[cfg(target_family = "unix")]
1530    use std::os::unix::fs::PermissionsExt;
1531
1532    #[cfg(target_family = "unix")]
1533    use tempfile::NamedTempFile;
1534
1535    use crate::{ObjectStoreExt, integration::*};
1536
1537    use super::*;
1538
1539    #[tokio::test]
1540    #[cfg(target_family = "unix")]
1541    async fn file_test() {
1542        let root = TempDir::new().unwrap();
1543        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
1544
1545        put_get_delete_list(&integration).await;
1546        list_with_offset_exclusivity(&integration).await;
1547        get_opts(&integration).await;
1548        list_uses_directories_correctly(&integration).await;
1549        list_with_delimiter(&integration).await;
1550        rename_and_copy(&integration).await;
1551        copy_if_not_exists(&integration).await;
1552        copy_rename_nonexistent_object(&integration).await;
1553        stream_get(&integration).await;
1554        put_opts(&integration, false).await;
1555    }
1556
1557    #[tokio::test]
1558    #[cfg(target_family = "unix")]
1559    async fn file_test_fsync() {
1560        // Run the full integration suite with fsync enabled to ensure the durability code
1561        // paths (file sync + directory fsync on put/copy/rename/multipart, including recursive
1562        // directory creation) behave identically to the default.
1563        let root = TempDir::new().unwrap();
1564        let integration = LocalFileSystem::new_with_prefix(root.path())
1565            .unwrap()
1566            .with_fsync(true);
1567
1568        put_get_delete_list(&integration).await;
1569        list_with_offset_exclusivity(&integration).await;
1570        get_opts(&integration).await;
1571        list_uses_directories_correctly(&integration).await;
1572        list_with_delimiter(&integration).await;
1573        rename_and_copy(&integration).await;
1574        copy_if_not_exists(&integration).await;
1575        copy_rename_nonexistent_object(&integration).await;
1576        stream_get(&integration).await;
1577        put_opts(&integration, false).await;
1578    }
1579
1580    #[tokio::test]
1581    async fn fsync_creates_nested_dirs() {
1582        // Exercises the recursive directory fsync in `create_parent_dirs`: every directory
1583        // component is newly created, so each must be synced up to the pre-existing root.
1584        let root = TempDir::new().unwrap();
1585        let integration = LocalFileSystem::new_with_prefix(root.path())
1586            .unwrap()
1587            .with_fsync(true);
1588
1589        let data = Bytes::from("arbitrary data");
1590
1591        // `put` (overwrite) into a deeply nested, non-existent directory tree
1592        let location = Path::from("a/b/c/d/put_file");
1593        integration
1594            .put(&location, data.clone().into())
1595            .await
1596            .unwrap();
1597        let read = integration
1598            .get(&location)
1599            .await
1600            .unwrap()
1601            .bytes()
1602            .await
1603            .unwrap();
1604        assert_eq!(read, data);
1605
1606        // multipart upload into another nested tree
1607        let location = Path::from("e/f/g/multipart_file");
1608        let mut upload = integration.put_multipart(&location).await.unwrap();
1609        upload.put_part(data.clone().into()).await.unwrap();
1610        upload.complete().await.unwrap();
1611        let read = integration
1612            .get(&location)
1613            .await
1614            .unwrap()
1615            .bytes()
1616            .await
1617            .unwrap();
1618        assert_eq!(read, data);
1619    }
1620
1621    #[tokio::test]
1622    #[cfg(target_family = "unix")]
1623    async fn fsync_rename_if_not_exists_propagates_source_delete_sync_error() {
1624        let root = TempDir::new().unwrap();
1625        let integration = LocalFileSystem::new_with_prefix(root.path())
1626            .unwrap()
1627            .with_fsync(true);
1628
1629        let source = Path::from("source_dir/source_file");
1630        let dest = Path::from("dest_dir/dest_file");
1631        integration.put(&source, "data".into()).await.unwrap();
1632
1633        let source_dir = root.path().join("source_dir");
1634        let original_permissions = fs::metadata(&source_dir).unwrap().permissions();
1635        // Allow unlinking the file from the directory, but prevent opening the
1636        // directory for fsync. Without the source-parent fsync, rename succeeds;
1637        // with it, the sync error is propagated after the source is deleted.
1638        fs::set_permissions(&source_dir, fs::Permissions::from_mode(0o300)).unwrap();
1639
1640        let result = integration.rename_if_not_exists(&source, &dest).await;
1641        fs::set_permissions(&source_dir, original_permissions).unwrap();
1642
1643        match result {
1644            Err(crate::Error::Generic { source, .. }) => {
1645                assert!(source.to_string().contains("Unable to sync data to disk"));
1646            }
1647            _ => panic!("expected source parent fsync to fail"),
1648        }
1649
1650        let read = integration.get(&dest).await.unwrap().bytes().await.unwrap();
1651        assert_eq!(read, Bytes::from("data"));
1652        assert!(matches!(
1653            integration.get(&source).await.unwrap_err(),
1654            crate::Error::NotFound { .. }
1655        ));
1656    }
1657
1658    #[test]
1659    #[cfg(target_family = "unix")]
1660    fn test_non_tokio() {
1661        let root = TempDir::new().unwrap();
1662        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
1663        futures_executor::block_on(async move {
1664            put_get_delete_list(&integration).await;
1665            list_uses_directories_correctly(&integration).await;
1666            list_with_delimiter(&integration).await;
1667
1668            // Can't use stream_get test as WriteMultipart uses a tokio JoinSet
1669            let p = Path::from("manual_upload");
1670            let mut upload = integration.put_multipart(&p).await.unwrap();
1671            upload.put_part("123".into()).await.unwrap();
1672            upload.put_part("45678".into()).await.unwrap();
1673            let r = upload.complete().await.unwrap();
1674
1675            let get = integration.get(&p).await.unwrap();
1676            assert_eq!(get.meta.e_tag.as_ref().unwrap(), r.e_tag.as_ref().unwrap());
1677            let actual = get.bytes().await.unwrap();
1678            assert_eq!(actual.as_ref(), b"12345678");
1679        });
1680    }
1681
1682    #[tokio::test]
1683    async fn creates_dir_if_not_present() {
1684        let root = TempDir::new().unwrap();
1685        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
1686
1687        let location = Path::from("nested/file/test_file");
1688
1689        let data = Bytes::from("arbitrary data");
1690
1691        integration
1692            .put(&location, data.clone().into())
1693            .await
1694            .unwrap();
1695
1696        let read_data = integration
1697            .get(&location)
1698            .await
1699            .unwrap()
1700            .bytes()
1701            .await
1702            .unwrap();
1703        assert_eq!(&*read_data, data);
1704    }
1705
1706    #[tokio::test]
1707    async fn unknown_length() {
1708        let root = TempDir::new().unwrap();
1709        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
1710
1711        let location = Path::from("some_file");
1712
1713        let data = Bytes::from("arbitrary data");
1714
1715        integration
1716            .put(&location, data.clone().into())
1717            .await
1718            .unwrap();
1719
1720        let read_data = integration
1721            .get(&location)
1722            .await
1723            .unwrap()
1724            .bytes()
1725            .await
1726            .unwrap();
1727        assert_eq!(&*read_data, data);
1728    }
1729
1730    #[tokio::test]
1731    async fn range_request_start_beyond_end_of_file() {
1732        let root = TempDir::new().unwrap();
1733        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
1734
1735        let location = Path::from("some_file");
1736
1737        let data = Bytes::from("arbitrary data");
1738
1739        integration
1740            .put(&location, data.clone().into())
1741            .await
1742            .unwrap();
1743
1744        integration
1745            .get_range(&location, 100..200)
1746            .await
1747            .expect_err("Should error with start range beyond end of file");
1748    }
1749
1750    #[tokio::test]
1751    async fn range_request_beyond_end_of_file() {
1752        let root = TempDir::new().unwrap();
1753        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
1754
1755        let location = Path::from("some_file");
1756
1757        let data = Bytes::from("arbitrary data");
1758
1759        integration
1760            .put(&location, data.clone().into())
1761            .await
1762            .unwrap();
1763
1764        let read_data = integration.get_range(&location, 0..100).await.unwrap();
1765        assert_eq!(&*read_data, data);
1766    }
1767
1768    #[tokio::test]
1769    async fn get_ranges_coalesces() {
1770        let root = TempDir::new().unwrap();
1771        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
1772        let location = Path::from("some_file");
1773
1774        // Position-dependent content so a misrouted slice is detectable.
1775        let data: Bytes = (0..4096u32)
1776            .map(|i| (i % 251) as u8)
1777            .collect::<Vec<_>>()
1778            .into();
1779        integration
1780            .put(&location, data.clone().into())
1781            .await
1782            .unwrap();
1783
1784        let cases: Vec<Vec<Range<u64>>> = vec![
1785            vec![0..16, 16..32, 32..48],                    // adjacent -> coalesced
1786            vec![0..16, 64..80, 1000..1016],                // gapped -> not coalesced
1787            vec![100..120, 0..50, 40..60, 0..50, 110..130], // out-of-order, overlapping, duplicate
1788            vec![],                                         // empty
1789        ];
1790
1791        for ranges in cases {
1792            let got = integration.get_ranges(&location, &ranges).await.unwrap();
1793            assert_eq!(got.len(), ranges.len());
1794            for (range, bytes) in ranges.iter().zip(got) {
1795                assert_eq!(
1796                    bytes,
1797                    data.slice(range.start as usize..range.end as usize),
1798                    "mismatch for range {range:?}"
1799                );
1800            }
1801        }
1802    }
1803
1804    #[tokio::test]
1805    #[cfg(target_family = "unix")]
1806    // Fails on github actions runner (which runs the tests as root)
1807    #[ignore]
1808    async fn bubble_up_io_errors() {
1809        use std::{fs::set_permissions, os::unix::prelude::PermissionsExt};
1810
1811        let root = TempDir::new().unwrap();
1812
1813        // make non-readable
1814        let metadata = root.path().metadata().unwrap();
1815        let mut permissions = metadata.permissions();
1816        permissions.set_mode(0o000);
1817        set_permissions(root.path(), permissions).unwrap();
1818
1819        let store = LocalFileSystem::new_with_prefix(root.path()).unwrap();
1820
1821        let mut stream = store.list(None);
1822        let mut any_err = false;
1823        while let Some(res) = stream.next().await {
1824            if res.is_err() {
1825                any_err = true;
1826            }
1827        }
1828        assert!(any_err);
1829
1830        // `list_with_delimiter
1831        assert!(store.list_with_delimiter(None).await.is_err());
1832    }
1833
1834    const NON_EXISTENT_NAME: &str = "nonexistentname";
1835
1836    #[tokio::test]
1837    async fn get_nonexistent_location() {
1838        let root = TempDir::new().unwrap();
1839        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
1840
1841        let location = Path::from(NON_EXISTENT_NAME);
1842
1843        let err = get_nonexistent_object(&integration, Some(location))
1844            .await
1845            .unwrap_err();
1846        if let crate::Error::NotFound { path, source } = err {
1847            let source_variant = source.downcast_ref::<std::io::Error>();
1848            assert!(
1849                matches!(source_variant, Some(std::io::Error { .. }),),
1850                "got: {source_variant:?}"
1851            );
1852            assert!(path.ends_with(NON_EXISTENT_NAME), "{}", path);
1853        } else {
1854            panic!("unexpected error type: {err:?}");
1855        }
1856    }
1857
1858    #[tokio::test]
1859    async fn root() {
1860        let integration = LocalFileSystem::new();
1861
1862        let canonical = std::path::Path::new("Cargo.toml").canonicalize().unwrap();
1863        let url = Url::from_directory_path(&canonical).unwrap();
1864        let path = Path::parse(url.path()).unwrap();
1865
1866        let roundtrip = integration.path_to_filesystem(&path).unwrap();
1867
1868        // Needed as on Windows canonicalize returns extended length path syntax
1869        // C:\Users\circleci -> \\?\C:\Users\circleci
1870        let roundtrip = roundtrip.canonicalize().unwrap();
1871
1872        assert_eq!(roundtrip, canonical);
1873
1874        integration.head(&path).await.unwrap();
1875    }
1876
1877    #[tokio::test]
1878    #[cfg(target_family = "windows")]
1879    async fn test_list_root() {
1880        let fs = LocalFileSystem::new();
1881        let r = fs.list_with_delimiter(None).await.unwrap_err().to_string();
1882
1883        assert!(
1884            r.contains("Unable to convert URL \"file:///\" to filesystem path"),
1885            "{}",
1886            r
1887        );
1888    }
1889
1890    #[tokio::test]
1891    #[cfg(target_os = "linux")]
1892    async fn test_list_root() {
1893        let fs = LocalFileSystem::new();
1894        fs.list_with_delimiter(None).await.unwrap();
1895    }
1896
1897    #[cfg(target_family = "unix")]
1898    async fn check_list(integration: &LocalFileSystem, prefix: Option<&Path>, expected: &[&str]) {
1899        let result: Vec<_> = integration.list(prefix).try_collect().await.unwrap();
1900
1901        let mut strings: Vec<_> = result.iter().map(|x| x.location.as_ref()).collect();
1902        strings.sort_unstable();
1903        assert_eq!(&strings, expected)
1904    }
1905
1906    #[tokio::test]
1907    #[cfg(target_family = "unix")]
1908    async fn test_symlink() {
1909        let root = TempDir::new().unwrap();
1910        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
1911
1912        let subdir = root.path().join("a");
1913        std::fs::create_dir(&subdir).unwrap();
1914        let file = subdir.join("file.parquet");
1915        std::fs::write(file, "test").unwrap();
1916
1917        check_list(&integration, None, &["a/file.parquet"]).await;
1918        integration
1919            .head(&Path::from("a/file.parquet"))
1920            .await
1921            .unwrap();
1922
1923        // Follow out of tree symlink
1924        let other = NamedTempFile::new().unwrap();
1925        std::os::unix::fs::symlink(other.path(), root.path().join("test.parquet")).unwrap();
1926
1927        // Should return test.parquet even though out of tree
1928        check_list(&integration, None, &["a/file.parquet", "test.parquet"]).await;
1929
1930        // Can fetch test.parquet
1931        integration.head(&Path::from("test.parquet")).await.unwrap();
1932
1933        // Follow in tree symlink
1934        std::os::unix::fs::symlink(&subdir, root.path().join("b")).unwrap();
1935        check_list(
1936            &integration,
1937            None,
1938            &["a/file.parquet", "b/file.parquet", "test.parquet"],
1939        )
1940        .await;
1941        check_list(&integration, Some(&Path::from("b")), &["b/file.parquet"]).await;
1942
1943        // Can fetch through symlink
1944        integration
1945            .head(&Path::from("b/file.parquet"))
1946            .await
1947            .unwrap();
1948
1949        // Ignore broken symlink
1950        std::os::unix::fs::symlink(root.path().join("foo.parquet"), root.path().join("c")).unwrap();
1951
1952        check_list(
1953            &integration,
1954            None,
1955            &["a/file.parquet", "b/file.parquet", "test.parquet"],
1956        )
1957        .await;
1958
1959        let mut r = integration.list_with_delimiter(None).await.unwrap();
1960        r.common_prefixes.sort_unstable();
1961        assert_eq!(r.common_prefixes.len(), 2);
1962        assert_eq!(r.common_prefixes[0].as_ref(), "a");
1963        assert_eq!(r.common_prefixes[1].as_ref(), "b");
1964        assert_eq!(r.objects.len(), 1);
1965        assert_eq!(r.objects[0].location.as_ref(), "test.parquet");
1966
1967        let r = integration
1968            .list_with_delimiter(Some(&Path::from("a")))
1969            .await
1970            .unwrap();
1971        assert_eq!(r.common_prefixes.len(), 0);
1972        assert_eq!(r.objects.len(), 1);
1973        assert_eq!(r.objects[0].location.as_ref(), "a/file.parquet");
1974
1975        // Deleting a symlink doesn't delete the source file
1976        integration
1977            .delete(&Path::from("test.parquet"))
1978            .await
1979            .unwrap();
1980        assert!(other.path().exists());
1981
1982        check_list(&integration, None, &["a/file.parquet", "b/file.parquet"]).await;
1983
1984        // Deleting through a symlink deletes both files
1985        integration
1986            .delete(&Path::from("b/file.parquet"))
1987            .await
1988            .unwrap();
1989
1990        check_list(&integration, None, &[]).await;
1991
1992        // Adding a file through a symlink creates in both paths
1993        integration
1994            .put(&Path::from("b/file.parquet"), vec![0, 1, 2].into())
1995            .await
1996            .unwrap();
1997
1998        check_list(&integration, None, &["a/file.parquet", "b/file.parquet"]).await;
1999    }
2000
2001    #[tokio::test]
2002    async fn invalid_path() {
2003        let root = TempDir::new().unwrap();
2004        let root = root.path().join("🙀");
2005        std::fs::create_dir(root.clone()).unwrap();
2006
2007        // Invalid paths supported above root of store
2008        let integration = LocalFileSystem::new_with_prefix(root.clone()).unwrap();
2009
2010        let directory = Path::from("directory");
2011        let object = directory.clone().join("child.txt");
2012        let data = Bytes::from("arbitrary");
2013        integration.put(&object, data.clone().into()).await.unwrap();
2014        integration.head(&object).await.unwrap();
2015        let result = integration.get(&object).await.unwrap();
2016        assert_eq!(result.bytes().await.unwrap(), data);
2017
2018        flatten_list_stream(&integration, None).await.unwrap();
2019        flatten_list_stream(&integration, Some(&directory))
2020            .await
2021            .unwrap();
2022
2023        let result = integration
2024            .list_with_delimiter(Some(&directory))
2025            .await
2026            .unwrap();
2027        assert_eq!(result.objects.len(), 1);
2028        assert!(result.common_prefixes.is_empty());
2029        assert_eq!(result.objects[0].location, object);
2030
2031        let emoji = root.join("💀");
2032        std::fs::write(emoji, "foo").unwrap();
2033
2034        // Can list illegal file
2035        let mut paths = flatten_list_stream(&integration, None).await.unwrap();
2036        paths.sort_unstable();
2037
2038        assert_eq!(
2039            paths,
2040            vec![
2041                Path::parse("directory/child.txt").unwrap(),
2042                Path::parse("💀").unwrap()
2043            ]
2044        );
2045    }
2046
2047    #[tokio::test]
2048    async fn list_hides_incomplete_uploads() {
2049        let root = TempDir::new().unwrap();
2050        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
2051        let location = Path::from("some_file");
2052
2053        let data = PutPayload::from("arbitrary data");
2054        let mut u1 = integration.put_multipart(&location).await.unwrap();
2055        u1.put_part(data.clone()).await.unwrap();
2056
2057        let mut u2 = integration.put_multipart(&location).await.unwrap();
2058        u2.put_part(data).await.unwrap();
2059
2060        let list = flatten_list_stream(&integration, None).await.unwrap();
2061        assert_eq!(list.len(), 0);
2062
2063        assert_eq!(
2064            integration
2065                .list_with_delimiter(None)
2066                .await
2067                .unwrap()
2068                .objects
2069                .len(),
2070            0
2071        );
2072    }
2073
2074    #[tokio::test]
2075    async fn test_path_with_offset() {
2076        let root = TempDir::new().unwrap();
2077        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
2078
2079        let root_path = root.path();
2080        for i in 0..5 {
2081            let filename = format!("test{i}.parquet");
2082            let file = root_path.join(filename);
2083            std::fs::write(file, "test").unwrap();
2084        }
2085        let filter_str = "test";
2086        let filter = String::from(filter_str);
2087        let offset_str = filter + "1";
2088        let offset = Path::from(offset_str.clone());
2089
2090        // Use list_with_offset to retrieve files
2091        let res = integration.list_with_offset(None, &offset);
2092        let offset_paths: Vec<_> = res.map_ok(|x| x.location).try_collect().await.unwrap();
2093        let mut offset_files: Vec<_> = offset_paths
2094            .iter()
2095            .map(|x| String::from(x.filename().unwrap()))
2096            .collect();
2097
2098        // Check result with direct filesystem read
2099        let files = fs::read_dir(root_path).unwrap();
2100        let filtered_files = files
2101            .filter_map(Result::ok)
2102            .filter_map(|d| {
2103                d.file_name().to_str().and_then(|f| {
2104                    if f.contains(filter_str) {
2105                        Some(String::from(f))
2106                    } else {
2107                        None
2108                    }
2109                })
2110            })
2111            .collect::<Vec<_>>();
2112
2113        let mut expected_offset_files: Vec<_> = filtered_files
2114            .iter()
2115            .filter(|s| **s > offset_str)
2116            .cloned()
2117            .collect();
2118
2119        fn do_vecs_match<T: PartialEq>(a: &[T], b: &[T]) -> bool {
2120            let matching = a.iter().zip(b.iter()).filter(|&(a, b)| a == b).count();
2121            matching == a.len() && matching == b.len()
2122        }
2123
2124        offset_files.sort();
2125        expected_offset_files.sort();
2126
2127        // println!("Expected Offset Files: {:?}", expected_offset_files);
2128        // println!("Actual Offset Files: {:?}", offset_files);
2129
2130        assert_eq!(offset_files.len(), expected_offset_files.len());
2131        assert!(do_vecs_match(&expected_offset_files, &offset_files));
2132    }
2133
2134    #[tokio::test]
2135    async fn filesystem_filename_with_percent() {
2136        let temp_dir = TempDir::new().unwrap();
2137        let integration = LocalFileSystem::new_with_prefix(temp_dir.path()).unwrap();
2138        let filename = "L%3ABC.parquet";
2139
2140        std::fs::write(temp_dir.path().join(filename), "foo").unwrap();
2141
2142        let res: Vec<_> = integration.list(None).try_collect().await.unwrap();
2143        assert_eq!(res.len(), 1);
2144        assert_eq!(res[0].location.as_ref(), filename);
2145
2146        let res = integration.list_with_delimiter(None).await.unwrap();
2147        assert_eq!(res.objects.len(), 1);
2148        assert_eq!(res.objects[0].location.as_ref(), filename);
2149    }
2150
2151    #[tokio::test]
2152    async fn relative_paths() {
2153        LocalFileSystem::new_with_prefix(".").unwrap();
2154        LocalFileSystem::new_with_prefix("..").unwrap();
2155        LocalFileSystem::new_with_prefix("../..").unwrap();
2156
2157        let integration = LocalFileSystem::new();
2158        let path = Path::from_filesystem_path(".").unwrap();
2159        integration.list_with_delimiter(Some(&path)).await.unwrap();
2160    }
2161
2162    #[test]
2163    fn test_valid_path() {
2164        let cases = [
2165            ("foo#123/test.txt", true),
2166            ("foo#123/test#23.txt", true),
2167            ("foo#123/test#34", false),
2168            ("foo😁/test#34", false),
2169            ("foo/test#😁34", true),
2170        ];
2171
2172        for (case, expected) in cases {
2173            let path = Path::parse(case).unwrap();
2174            assert_eq!(is_valid_file_path(&path), expected);
2175        }
2176    }
2177
2178    #[tokio::test]
2179    async fn test_intermediate_files() {
2180        let root = TempDir::new().unwrap();
2181        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
2182
2183        let a = Path::parse("foo#123/test.txt").unwrap();
2184        integration.put(&a, "test".into()).await.unwrap();
2185
2186        let list = flatten_list_stream(&integration, None).await.unwrap();
2187        assert_eq!(list, vec![a.clone()]);
2188
2189        std::fs::write(root.path().join("bar#123"), "test").unwrap();
2190
2191        // Should ignore file
2192        let list = flatten_list_stream(&integration, None).await.unwrap();
2193        assert_eq!(list, vec![a.clone()]);
2194
2195        let b = Path::parse("bar#123").unwrap();
2196        let err = integration.get(&b).await.unwrap_err().to_string();
2197        assert_eq!(
2198            err,
2199            "Generic LocalFileSystem error: Filenames containing trailing '/#\\d+/' are not supported: bar#123"
2200        );
2201
2202        let c = Path::parse("foo#123.txt").unwrap();
2203        integration.put(&c, "test".into()).await.unwrap();
2204
2205        let mut list = flatten_list_stream(&integration, None).await.unwrap();
2206        list.sort_unstable();
2207        assert_eq!(list, vec![c, a]);
2208    }
2209
2210    #[tokio::test]
2211    #[cfg(target_os = "windows")]
2212    async fn filesystem_filename_with_colon() {
2213        let root = TempDir::new().unwrap();
2214        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
2215        let path = Path::parse("file%3Aname.parquet").unwrap();
2216        let location = Path::parse("file:name.parquet").unwrap();
2217
2218        integration.put(&location, "test".into()).await.unwrap();
2219        let list = flatten_list_stream(&integration, None).await.unwrap();
2220        assert_eq!(list, vec![path.clone()]);
2221
2222        let result = integration
2223            .get(&location)
2224            .await
2225            .unwrap()
2226            .bytes()
2227            .await
2228            .unwrap();
2229        assert_eq!(result, Bytes::from("test"));
2230    }
2231
2232    #[tokio::test]
2233    async fn delete_dirs_automatically() {
2234        let root = TempDir::new().unwrap();
2235        let integration = LocalFileSystem::new_with_prefix(root.path())
2236            .unwrap()
2237            .with_automatic_cleanup(true);
2238        let location = Path::from("nested/file/test_file");
2239        let data = Bytes::from("arbitrary data");
2240
2241        integration
2242            .put(&location, data.clone().into())
2243            .await
2244            .unwrap();
2245
2246        let read_data = integration
2247            .get(&location)
2248            .await
2249            .unwrap()
2250            .bytes()
2251            .await
2252            .unwrap();
2253
2254        assert_eq!(&*read_data, data);
2255        assert!(fs::read_dir(root.path()).unwrap().count() > 0);
2256        integration.delete(&location).await.unwrap();
2257        assert!(fs::read_dir(root.path()).unwrap().count() == 0);
2258    }
2259
2260    #[test]
2261    #[cfg(target_family = "unix")]
2262    fn test_close_file_detects_error_unix() {
2263        let err = super::close_fd(-1).unwrap_err();
2264        assert_eq!(err.raw_os_error(), Some(nix::libc::EBADF), "got: {err:?}");
2265    }
2266
2267    #[test]
2268    #[cfg(target_family = "windows")]
2269    fn test_close_file_detects_error_windows() {
2270        let err = super::close_handle(std::ptr::null_mut()).unwrap_err();
2271        assert_eq!(
2272            err.raw_os_error(),
2273            Some(windows_sys::Win32::Foundation::ERROR_INVALID_HANDLE as i32),
2274            "got: {err:?}"
2275        );
2276    }
2277}
2278
2279#[cfg(not(target_arch = "wasm32"))]
2280#[cfg(test)]
2281mod not_wasm_tests {
2282    use std::time::Duration;
2283    use tempfile::TempDir;
2284
2285    use crate::local::LocalFileSystem;
2286    use crate::{ObjectStoreExt, Path, PutPayload};
2287
2288    #[tokio::test]
2289    async fn test_cleanup_intermediate_files() {
2290        let root = TempDir::new().unwrap();
2291        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
2292
2293        let location = Path::from("some_file");
2294        let data = PutPayload::from_static(b"hello");
2295        let mut upload = integration.put_multipart(&location).await.unwrap();
2296        upload.put_part(data).await.unwrap();
2297
2298        let file_count = std::fs::read_dir(root.path()).unwrap().count();
2299        assert_eq!(file_count, 1);
2300        drop(upload);
2301
2302        for _ in 0..100 {
2303            tokio::time::sleep(Duration::from_millis(1)).await;
2304            let file_count = std::fs::read_dir(root.path()).unwrap().count();
2305            if file_count == 0 {
2306                return;
2307            }
2308        }
2309        panic!("Failed to cleanup file in 100ms")
2310    }
2311}
2312
2313#[cfg(target_family = "unix")]
2314#[cfg(test)]
2315mod unix_test {
2316    use std::fs::OpenOptions;
2317
2318    use nix::sys::stat;
2319    use nix::unistd;
2320    use tempfile::TempDir;
2321
2322    use crate::local::LocalFileSystem;
2323    use crate::{ObjectStoreExt, Path};
2324
2325    #[tokio::test]
2326    async fn test_fifo() {
2327        let filename = "some_file";
2328        let root = TempDir::new().unwrap();
2329        let integration = LocalFileSystem::new_with_prefix(root.path()).unwrap();
2330        let path = root.path().join(filename);
2331        unistd::mkfifo(&path, stat::Mode::S_IRWXU).unwrap();
2332
2333        // Need to open read and write side in parallel
2334        let spawned =
2335            tokio::task::spawn_blocking(|| OpenOptions::new().write(true).open(path).unwrap());
2336
2337        let location = Path::from(filename);
2338        integration.head(&location).await.unwrap();
2339        integration.get(&location).await.unwrap();
2340
2341        spawned.await.unwrap();
2342    }
2343}