1use 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#[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
143fn 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 match unsafe { windows_sys::Win32::Foundation::CloseHandle(handle) } {
175 0 => Err(io::Error::last_os_error()),
176 _ => Ok(()),
177 }
178}
179
180#[derive(Clone, Debug)]
241pub struct LocalFileSystem {
242 config: Arc<Config>,
243 automatic_cleanup: bool,
245 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 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 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 pub fn path_to_filesystem(&self, location: &Path) -> Result<PathBuf> {
299 self.config.path_to_filesystem(location)
300 }
301
302 pub fn with_automatic_cleanup(mut self, automatic_cleanup: bool) -> Self {
304 self.automatic_cleanup = automatic_cleanup;
305 self
306 }
307
308 pub fn with_fsync(mut self, fsync: bool) -> Self {
323 self.fsync = fsync;
324 self
325 }
326}
327
328impl Config {
329 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 .pop_if_empty()
337 .extend(location.parts());
338
339 url.to_file_path()
340 .map_err(|_| Error::InvalidUrl { url }.into())
341 }
342
343 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 let mut out = String::new();
359 let drive = &path[..2]; let filepath = &path[2..].replace(':', "%3A"); out.push_str(drive);
362 out.push_str(filepath);
363 PathBuf::from(out)
364 };
365
366 Ok(path)
367 }
368
369 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 !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 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); 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 let mut file = File::open(&path).map_err(|e| map_open_error(e, &path))?;
507
508 let fetch_ranges = merge_ranges(&ranges, 0);
512
513 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 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 maybe_spawn_blocking(move || {
647 loop {
648 let staged = staged_upload_path(&to, &id.to_string());
649 match std::fs::hard_link(&from, &staged) {
653 Ok(_) => match rename(&staged, &to, fsync) {
655 Ok(_) => return Ok(()),
656 Err(source) => {
657 let _ = std::fs::remove_file(&staged); 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 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 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 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 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 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 .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 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 tokio::runtime::Handle::try_current().is_err() {
870 return futures_util::stream::iter(s).boxed();
871 }
872
873 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
901fn 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 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 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
949fn 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 if from.parent() != to.parent() {
960 fsync_parent_dir(from)?;
961 }
962 }
963 Ok(())
964}
965
966fn 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
978fn 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
996fn 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); 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
1022fn 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
1035fn 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
1044fn 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
1060fn 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
1080fn 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 state: Arc<UploadState>,
1092 src: Option<PathBuf>,
1094 offset: u64,
1096 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 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 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 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 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 !buf_slice.is_empty() {
1290 let metadata = open_metadata(file, path)?;
1291 let file_len = metadata.len();
1292
1293 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 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
1378fn 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
1387fn map_seek_error(source: io::Error, file: &File, path: &std::path::Path, requested: u64) -> Error {
1389 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 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)]
1463fn get_inode(metadata: &Metadata) -> u64 {
1466 std::os::unix::fs::MetadataExt::ino(metadata)
1467}
1468
1469#[cfg(not(unix))]
1470fn get_inode(_metadata: &Metadata) -> u64 {
1472 0
1473}
1474
1475fn convert_walkdir_result(
1478 res: std::result::Result<DirEntry, walkdir::Error>,
1479) -> Result<Option<DirEntry>> {
1480 match res {
1481 Ok(entry) => {
1482 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 Ok(Some(entry))
1492 }
1493 Err(_) => {
1494 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 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 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 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 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 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 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 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], vec![0..16, 64..80, 1000..1016], vec![100..120, 0..50, 40..60, 0..50, 110..130], vec![], ];
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 #[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 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 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 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 let other = NamedTempFile::new().unwrap();
1925 std::os::unix::fs::symlink(other.path(), root.path().join("test.parquet")).unwrap();
1926
1927 check_list(&integration, None, &["a/file.parquet", "test.parquet"]).await;
1929
1930 integration.head(&Path::from("test.parquet")).await.unwrap();
1932
1933 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 integration
1945 .head(&Path::from("b/file.parquet"))
1946 .await
1947 .unwrap();
1948
1949 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 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 integration
1986 .delete(&Path::from("b/file.parquet"))
1987 .await
1988 .unwrap();
1989
1990 check_list(&integration, None, &[]).await;
1991
1992 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 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 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 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 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 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 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 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}