1use std::path::{Component, Path, PathBuf};
18use std::sync::{Arc, Mutex};
19use std::vec::IntoIter;
20use std::{fmt, io};
21
22use cap_std::ambient_authority;
23use cap_std::fs::{Dir, DirEntry, OpenOptions, ReadDir};
24use opendal::layers::SimulateLayer;
25use opendal::raw::*;
26use opendal::{
27 Buffer, BytesRange, Capability, EntryMode, Error, ErrorKind, Metadata, MetadataBuilder,
28 OperationContext, Operator, Result,
29};
30use tokio::io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt};
31
32const LIST_BATCH_SIZE: usize = 128;
33
34#[derive(Clone)]
36pub struct SecureFsRoot {
37 dir: Arc<Dir>,
38 path: Arc<PathBuf>,
39}
40
41impl fmt::Debug for SecureFsRoot {
42 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
43 f.debug_struct("SecureFsRoot")
44 .field("path", &self.path)
45 .finish_non_exhaustive()
46 }
47}
48
49impl SecureFsRoot {
50 pub fn open(path: impl AsRef<Path>) -> io::Result<Self> {
54 let path = path.as_ref();
55 std::fs::create_dir_all(path)?;
56 if std::fs::symlink_metadata(path)?.file_type().is_symlink() {
57 return Err(io::Error::new(
58 io::ErrorKind::InvalidInput,
59 "filesystem sandbox root must not be a symbolic link",
60 ));
61 }
62
63 let path = path.canonicalize()?;
64 let dir = Dir::open_ambient_dir(&path, ambient_authority())?;
65 Ok(Self {
66 dir: Arc::new(dir),
67 path: Arc::new(path),
68 })
69 }
70
71 pub fn path(&self) -> &Path {
73 &self.path
74 }
75
76 pub async fn is_file(&self, path: &str) -> Result<bool> {
78 let path = backend_path(path).map_err(new_std_io_error)?;
79 let root = self.clone();
80 common_runtime::spawn_blocking_global(move || {
81 root.dir
82 .symlink_metadata(path)
83 .map(|metadata| metadata.is_file())
84 })
85 .await
86 .map_err(new_task_join_error)?
87 .map_err(new_std_io_error)
88 }
89
90 pub fn open_subdir(&self, path: impl AsRef<Path>) -> io::Result<Self> {
92 let path = normalize_relative_path(path.as_ref())?;
93 if path.as_os_str().is_empty() {
94 return Ok(self.clone());
95 }
96
97 let dir = self.dir.open_dir(&path)?;
98 Ok(Self {
99 dir: Arc::new(dir),
100 path: Arc::new(self.path.join(path)),
101 })
102 }
103
104 pub fn create_subdir(&self, path: impl AsRef<Path>) -> io::Result<Self> {
106 let path = normalize_relative_path(path.as_ref())?;
107 if path.as_os_str().is_empty() {
108 return Ok(self.clone());
109 }
110
111 self.dir.create_dir_all(&path)?;
112 let dir = self.dir.open_dir(&path)?;
113 Ok(Self {
114 dir: Arc::new(dir),
115 path: Arc::new(self.path.join(path)),
116 })
117 }
118
119 pub fn build_operator(&self) -> Operator {
121 Operator::from_parts(
122 OperationContext::default(),
123 Arc::new(SecureFsBackend::new(self.clone())) as Servicer,
124 )
125 .layer(SimulateLayer::default())
126 }
127}
128
129fn normalize_relative_path(path: &Path) -> io::Result<PathBuf> {
130 let mut normalized = PathBuf::new();
131 for component in path.components() {
132 match component {
133 Component::CurDir => {}
134 Component::Normal(value) => normalized.push(value),
135 Component::ParentDir | Component::RootDir | Component::Prefix(_) => {
136 return Err(io::Error::new(
137 io::ErrorKind::PermissionDenied,
138 "path escapes the filesystem sandbox",
139 ));
140 }
141 }
142 }
143 Ok(normalized)
144}
145
146fn backend_path(path: &str) -> io::Result<PathBuf> {
147 let path = path.trim_matches('/');
148 if path.is_empty() {
149 Ok(PathBuf::new())
150 } else {
151 normalize_relative_path(Path::new(path))
152 }
153}
154
155fn parse_write_error(error: io::Error, if_not_exists: bool) -> Error {
156 if if_not_exists && error.kind() == io::ErrorKind::AlreadyExists {
157 Error::new(
158 ErrorKind::ConditionNotMatch,
159 "the file already exists in the filesystem",
160 )
161 .set_source(error)
162 } else {
163 new_std_io_error(error)
164 }
165}
166
167fn metadata_from_fs(metadata: cap_std::fs::Metadata) -> Result<Metadata> {
168 let mut builder = if metadata.is_dir() {
169 MetadataBuilder::dir()
170 } else if metadata.is_file() {
171 MetadataBuilder::file(metadata.len())
172 } else {
173 MetadataBuilder::unknown()
174 };
175
176 builder.last_modified(Timestamp::try_from(
177 metadata.modified().map_err(new_std_io_error)?.into_std(),
178 )?);
179 Ok(builder.build())
180}
181
182#[derive(Clone, Debug)]
183struct SecureFsBackend {
184 root: SecureFsRoot,
185 info: ServiceInfo,
186 capability: Capability,
187}
188
189impl SecureFsBackend {
190 fn new(root: SecureFsRoot) -> Self {
191 let info = ServiceInfo::new("fs", root.path().to_string_lossy(), "");
192 let capability = Capability {
193 stat: true,
194 read: true,
195 write: true,
196 write_can_empty: true,
197 write_can_append: true,
198 write_can_multi: true,
199 write_with_if_not_exists: true,
200 create_dir: true,
201 delete: true,
202 delete_with_recursive: true,
203 list: true,
204 shared: true,
205 ..Default::default()
206 };
207 Self {
208 root,
209 info,
210 capability,
211 }
212 }
213}
214
215impl Service for SecureFsBackend {
216 type Reader = oio::StreamReader<SecureFsReader>;
217 type Writer = SecureFsWriter;
218 type Lister = SecureFsLister;
219 type Deleter = oio::OneShotDeleter<SecureFsDeleter>;
220 type Copier = ();
221 type Composer = ();
222
223 fn info(&self) -> ServiceInfo {
224 self.info.clone()
225 }
226
227 fn capability(&self) -> Capability {
228 self.capability
229 }
230
231 async fn create_dir(
232 &self,
233 _: &OperationContext,
234 path: &str,
235 _: OpCreateDir,
236 ) -> Result<RpCreateDir> {
237 let path = backend_path(path).map_err(new_std_io_error)?;
238 let root = self.root.clone();
239 common_runtime::spawn_blocking_global(move || root.dir.create_dir_all(path))
240 .await
241 .map_err(new_task_join_error)?
242 .map_err(new_std_io_error)?;
243 Ok(RpCreateDir::default())
244 }
245
246 async fn stat(&self, _: &OperationContext, path: &str, _: OpStat) -> Result<RpStat> {
247 let path = backend_path(path).map_err(new_std_io_error)?;
248 let root = self.root.clone();
249 let metadata = common_runtime::spawn_blocking_global(move || {
250 if path.as_os_str().is_empty() {
251 root.dir.dir_metadata()
252 } else {
253 root.dir.metadata(path)
254 }
255 })
256 .await
257 .map_err(new_task_join_error)?
258 .map_err(new_std_io_error)?;
259 Ok(RpStat::new(metadata_from_fs(metadata)?))
260 }
261
262 fn read(&self, _: &OperationContext, path: &str, _: OpRead) -> Result<Self::Reader> {
263 let path = backend_path(path).map_err(new_std_io_error)?;
264 Ok(oio::StreamReader::new(SecureFsReader {
265 root: self.root.clone(),
266 path,
267 }))
268 }
269
270 fn write(&self, _: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
271 let path = backend_path(path).map_err(new_std_io_error)?;
272 Ok(SecureFsWriter {
273 root: self.root.clone(),
274 path,
275 args,
276 file: None,
277 synced: false,
278 })
279 }
280
281 fn delete(&self, _: &OperationContext) -> Result<Self::Deleter> {
282 Ok(oio::OneShotDeleter::new(SecureFsDeleter {
283 root: self.root.clone(),
284 }))
285 }
286
287 fn list(&self, _: &OperationContext, path: &str, _: OpList) -> Result<Self::Lister> {
288 let path = backend_path(path).map_err(new_std_io_error)?;
289 let display_prefix = if path.as_os_str().is_empty() {
290 String::new()
291 } else {
292 format!("{}/", path.to_string_lossy().replace('\\', "/"))
293 };
294 Ok(SecureFsLister {
295 root: self.root.clone(),
296 path,
297 display_prefix,
298 read_dir: None,
299 entries: vec![].into_iter(),
300 seeded: false,
301 done: false,
302 })
303 }
304
305 fn copy(&self, _: &OperationContext, _: &str, _: &str, _: OpCopy) -> Result<Self::Copier> {
306 Err(Error::new(
307 ErrorKind::Unsupported,
308 "operation is not supported",
309 ))
310 }
311
312 async fn rename(
313 &self,
314 _: &OperationContext,
315 _: &str,
316 _: &str,
317 _: OpRename,
318 ) -> Result<RpRename> {
319 Err(Error::new(
320 ErrorKind::Unsupported,
321 "operation is not supported",
322 ))
323 }
324
325 async fn presign(&self, _: &OperationContext, _: &str, _: OpPresign) -> Result<RpPresign> {
326 Err(Error::new(
327 ErrorKind::Unsupported,
328 "operation is not supported",
329 ))
330 }
331}
332
333struct SecureFsReader {
334 root: SecureFsRoot,
335 path: PathBuf,
336}
337
338impl oio::StreamRead for SecureFsReader {
339 async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
340 let path = self.path.clone();
341 let root = self.root.clone();
342 let file = common_runtime::spawn_blocking_global(move || root.dir.open(path))
343 .await
344 .map_err(new_task_join_error)?
345 .map_err(new_std_io_error)?;
346 let mut file = tokio::fs::File::from_std(file.into_std());
347 if range.offset() != 0 {
348 file.seek(io::SeekFrom::Start(range.offset()))
349 .await
350 .map_err(new_std_io_error)?;
351 }
352 Ok((
353 RpRead::default(),
354 Box::new(SecureFsReadStream {
355 file,
356 remaining: range.size().unwrap_or(u64::MAX),
357 }),
358 ))
359 }
360}
361
362struct SecureFsReadStream {
363 file: tokio::fs::File,
364 remaining: u64,
365}
366
367impl oio::ReadStream for SecureFsReadStream {
368 async fn read(&mut self) -> Result<Buffer> {
369 if self.remaining == 0 {
370 return Ok(Buffer::new());
371 }
372
373 let size = self.remaining.min(2 * 1024 * 1024) as usize;
374 let mut buffer = vec![0; size];
375 let read = self
376 .file
377 .read(&mut buffer)
378 .await
379 .map_err(new_std_io_error)?;
380 self.remaining = self.remaining.saturating_sub(read as u64);
381 buffer.truncate(read);
382 Ok(Buffer::from(buffer))
383 }
384}
385
386struct SecureFsWriter {
387 root: SecureFsRoot,
388 path: PathBuf,
389 args: OpWrite,
390 file: Option<tokio::fs::File>,
391 synced: bool,
392}
393
394#[derive(Debug)]
395struct UnsyncedOverwrite {
396 flush_error: Option<Error>,
397}
398
399impl fmt::Display for UnsyncedOverwrite {
400 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
401 write!(f, "overwrite file was not synced")
402 }
403}
404
405impl std::error::Error for UnsyncedOverwrite {
406 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
407 self.flush_error.as_ref().map(|error| error as _)
408 }
409}
410
411pub fn is_unsynced_overwrite_abort(error: &Error) -> bool {
413 std::error::Error::source(error).is_some_and(|source| source.is::<UnsyncedOverwrite>())
414}
415
416impl SecureFsWriter {
417 async fn ensure_file(&mut self) -> Result<&mut tokio::fs::File> {
418 if self.file.is_none() {
419 let path = self.path.clone();
420 let root = self.root.clone();
421 let if_not_exists = self.args.if_not_exists();
422 let append = self.args.append();
423 let file = common_runtime::spawn_blocking_global(move || {
424 if let Some(parent) = path.parent()
425 && !parent.as_os_str().is_empty()
426 {
427 root.dir.create_dir_all(parent).map_err(new_std_io_error)?;
428 }
429
430 let mut options = OpenOptions::new();
431 options.write(true);
432 if if_not_exists {
433 options.create_new(true);
434 } else {
435 options.create(true);
436 }
437 if append {
438 options.append(true);
439 } else {
440 options.truncate(true);
441 }
442 root.dir
443 .open_with(path, &options)
444 .map_err(|error| parse_write_error(error, if_not_exists))
445 })
446 .await
447 .map_err(new_task_join_error)??;
448
449 self.file = Some(tokio::fs::File::from_std(file.into_std()));
450 }
451 Ok(self.file.as_mut().expect("file must be initialized"))
452 }
453}
454
455impl oio::Write for SecureFsWriter {
456 async fn write(&mut self, buffer: Buffer) -> Result<()> {
457 self.ensure_file()
458 .await?
459 .write_all(&buffer.to_bytes())
460 .await
461 .map_err(new_std_io_error)
462 }
463
464 async fn close(&mut self) -> Result<Metadata> {
465 {
466 let file = self.ensure_file().await?;
467 file.flush().await.map_err(new_std_io_error)?;
468 file.sync_all().await.map_err(new_std_io_error)?;
469 }
470 self.synced = true;
471 let metadata = self
472 .ensure_file()
473 .await?
474 .metadata()
475 .await
476 .map_err(new_std_io_error)?;
477 let mut builder = MetadataBuilder::file(metadata.len());
478 builder.last_modified(Timestamp::try_from(
479 metadata.modified().map_err(new_std_io_error)?,
480 )?);
481 Ok(builder.build())
482 }
483
484 async fn abort(&mut self) -> Result<()> {
485 let flush = match self.file.as_mut() {
487 Some(file) => file.flush().await.map_err(new_std_io_error),
488 None => Ok(()),
489 };
490 if self.args.if_not_exists() {
491 if let Some(file) = self.file.take() {
494 drop(file);
495 if !self.synced {
496 let root = self.root.clone();
497 let path = self.path.clone();
498 let cleanup =
499 common_runtime::spawn_blocking_global(move || root.dir.remove_file(path))
500 .await
501 .map_err(new_task_join_error)
502 .and_then(|result| result.map_err(new_std_io_error));
503 return flush.and(cleanup);
504 }
505 }
506 return flush;
507 }
508 if self.file.is_none() {
509 return flush;
510 }
511 let error = Error::new(
512 ErrorKind::Unsupported,
513 "filesystem writes cannot be aborted without atomic writes",
514 );
515 if !self.synced {
516 return Err(error.set_source(UnsyncedOverwrite {
517 flush_error: flush.err(),
518 }));
519 }
520 flush?;
521 Err(error)
522 }
523}
524
525struct SecureFsLister {
526 root: SecureFsRoot,
527 path: PathBuf,
528 display_prefix: String,
529 read_dir: Option<Arc<Mutex<ReadDir>>>,
530 entries: IntoIter<oio::Entry>,
531 seeded: bool,
532 done: bool,
533}
534
535impl oio::List for SecureFsLister {
536 async fn next(&mut self) -> Result<Option<oio::Entry>> {
537 if !self.seeded {
538 self.seeded = true;
539 let path = self.path.clone();
540 let root = self.root.clone();
541 let display_prefix = self.display_prefix.clone();
542 let read_dir = common_runtime::spawn_blocking_global(move || {
543 let result = (|| {
544 let dir = if path.as_os_str().is_empty() {
545 root.dir.open_dir(".")?
546 } else {
547 root.dir.open_dir(&path)?
548 };
549 dir.entries()
550 })();
551
552 match result {
553 Ok(read_dir) => Ok(Some(read_dir)),
554 Err(error)
555 if matches!(
556 error.kind(),
557 io::ErrorKind::NotFound | io::ErrorKind::NotADirectory
558 ) =>
559 {
560 Ok(None)
561 }
562 Err(error) => Err(error),
563 }
564 })
565 .await
566 .map_err(new_task_join_error)?
567 .map_err(new_std_io_error)?;
568
569 let Some(read_dir) = read_dir else {
570 self.done = true;
571 return Ok(None);
572 };
573 self.read_dir = Some(Arc::new(Mutex::new(read_dir)));
574 let current_path = oio::Entry::new(
575 if display_prefix.is_empty() {
576 "/"
577 } else {
578 &display_prefix
579 },
580 MetadataBuilder::dir().build(),
581 );
582 self.entries = vec![current_path].into_iter();
583 }
584
585 if let Some(entry) = self.entries.next() {
586 return Ok(Some(entry));
587 }
588 if self.done {
589 return Ok(None);
590 }
591
592 let Some(read_dir) = self.read_dir.clone() else {
593 self.done = true;
594 return Ok(None);
595 };
596 let display_prefix = self.display_prefix.clone();
597 let (entries, done) = common_runtime::spawn_blocking_global(move || {
598 let mut read_dir = read_dir
599 .lock()
600 .map_err(|_| io::Error::other("filesystem directory iterator lock is poisoned"))?;
601 read_list_batch(&mut read_dir, &display_prefix)
602 })
603 .await
604 .map_err(new_task_join_error)?
605 .map_err(new_std_io_error)?;
606
607 self.entries = entries.into_iter();
608 self.done = done;
609 Ok(self.entries.next())
610 }
611}
612
613fn read_list_batch(
614 read_dir: &mut ReadDir,
615 display_prefix: &str,
616) -> io::Result<(Vec<oio::Entry>, bool)> {
617 let mut entries = Vec::with_capacity(LIST_BATCH_SIZE);
618 while entries.len() < LIST_BATCH_SIZE {
619 let entry = match read_dir.next() {
620 Some(Ok(entry)) => entry,
621 Some(Err(error)) if error.kind() == io::ErrorKind::NotFound => {
622 return Ok((entries, true));
623 }
624 Some(Err(error)) => return Err(error),
625 None => return Ok((entries, true)),
626 };
627
628 if let Some(entry) = read_list_entry(entry, display_prefix)? {
629 entries.push(entry);
630 }
631 }
632 Ok((entries, false))
633}
634
635fn read_list_entry(entry: DirEntry, display_prefix: &str) -> io::Result<Option<oio::Entry>> {
636 let file_type = match entry.file_type() {
637 Ok(file_type) => file_type,
638 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
639 Err(error) => return Err(error),
640 };
641 let name = entry.file_name().to_string_lossy().to_string();
642 let (path, mode) = if file_type.is_dir() {
643 (format!("{display_prefix}{name}/"), EntryMode::DIR)
644 } else if file_type.is_file() {
645 (format!("{display_prefix}{name}"), EntryMode::FILE)
646 } else {
647 (format!("{display_prefix}{name}"), EntryMode::Unknown)
648 };
649 let metadata = if mode == EntryMode::Unknown {
650 MetadataBuilder::unknown().build()
651 } else {
652 match entry.metadata() {
653 Ok(metadata) => match metadata_from_fs(metadata) {
654 Ok(metadata) => metadata,
655 Err(error) if error.kind() == ErrorKind::NotFound => return Ok(None),
656 Err(error) => return Err(io::Error::other(error.to_string())),
657 },
658 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
659 Err(error) => return Err(error),
660 }
661 };
662 Ok(Some(oio::Entry::new(&path, metadata)))
663}
664
665struct SecureFsDeleter {
666 root: SecureFsRoot,
667}
668
669impl oio::OneShotDelete for SecureFsDeleter {
670 async fn delete_once(&self, path: String, args: OpDelete) -> Result<()> {
671 let path = backend_path(&path).map_err(new_std_io_error)?;
672 if path.as_os_str().is_empty() {
673 return Err(Error::new(
674 ErrorKind::Unsupported,
675 "deleting the filesystem sandbox root is not supported",
676 ));
677 }
678 let root = self.root.clone();
679 common_runtime::spawn_blocking_global(move || {
680 let metadata = match root.dir.symlink_metadata(&path) {
681 Ok(metadata) => metadata,
682 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
683 Err(error) => return Err(error),
684 };
685
686 if metadata.is_dir() {
687 if args.recursive() {
688 root.dir.remove_dir_all(path)
689 } else {
690 root.dir.remove_dir(path)
691 }
692 } else {
693 root.dir.remove_file(path)
694 }
695 })
696 .await
697 .map_err(new_task_join_error)?
698 .map_err(new_std_io_error)
699 }
700}
701
702#[cfg(test)]
703mod tests {
704 use bytes::Bytes;
705 use common_test_util::temp_dir::create_temp_dir;
706 use opendal::raw::oio::List;
707 use opendal::raw::{OpList, Service};
708 use opendal::{BytesRange, ErrorKind, OperationContext};
709
710 use super::{LIST_BATCH_SIZE, SecureFsBackend, SecureFsRoot, read_list_entry};
711
712 #[tokio::test]
713 async fn test_operator_suffix_reads_final_bytes() {
714 let temp_dir = create_temp_dir("secure_fs_operator_suffix");
715 std::fs::write(temp_dir.path().join("file"), b"0123456789").unwrap();
716 let operator = SecureFsRoot::open(temp_dir.path())
717 .unwrap()
718 .build_operator();
719
720 assert_eq!(
721 Bytes::from_static(b"789"),
722 operator
723 .read_with("file")
724 .range(7..10)
725 .await
726 .unwrap()
727 .to_bytes()
728 );
729 assert_eq!(
730 Bytes::from_static(b"789"),
731 operator
732 .read_with("file")
733 .range(BytesRange::Suffix { size: 3 })
734 .await
735 .unwrap()
736 .to_bytes()
737 );
738 }
739
740 #[tokio::test]
741 async fn test_lister_streams_entries() {
742 let temp_dir = create_temp_dir("secure_fs_lister_streams_entries");
743 for index in 0..129 {
744 std::fs::write(temp_dir.path().join(format!("{index}.parquet")), []).unwrap();
745 }
746
747 let root = SecureFsRoot::open(temp_dir.path()).unwrap();
748 let backend = SecureFsBackend::new(root);
749 let ctx = OperationContext::default();
750 let mut lister = backend.list(&ctx, "/", OpList::new()).unwrap();
751
752 let mut paths = Vec::new();
753 while let Some(entry) = lister.next().await.unwrap() {
754 paths.push(entry.path().to_string());
755 assert!(lister.entries.len() <= LIST_BATCH_SIZE);
756 }
757 assert_eq!(130, paths.len());
758 assert!(paths.iter().any(|path| path == "/"));
759 assert!(paths.iter().any(|path| path == "128.parquet"));
760 }
761
762 #[tokio::test]
763 async fn test_if_not_exists_returns_condition_not_match() {
764 let temp_dir = create_temp_dir("secure_fs_if_not_exists");
765 let operator = SecureFsRoot::open(temp_dir.path())
766 .unwrap()
767 .build_operator();
768 operator
769 .write("existing", Bytes::from_static(b"original"))
770 .await
771 .unwrap();
772
773 let error = operator
774 .write_with("existing", Bytes::from_static(b"replacement"))
775 .if_not_exists(true)
776 .await
777 .unwrap_err();
778
779 assert_eq!(ErrorKind::ConditionNotMatch, error.kind());
780 assert_eq!(
781 Bytes::from_static(b"original"),
782 operator.read("existing").await.unwrap().to_bytes()
783 );
784 }
785
786 #[tokio::test]
787 async fn test_if_not_exists_does_not_remap_parent_directory_error() {
788 let temp_dir = create_temp_dir("secure_fs_if_not_exists_parent_error");
789 std::fs::write(temp_dir.path().join("parent"), []).unwrap();
790 let operator = SecureFsRoot::open(temp_dir.path())
791 .unwrap()
792 .build_operator();
793
794 let error = operator
795 .write_with("parent/file", Bytes::new())
796 .if_not_exists(true)
797 .await
798 .unwrap_err();
799
800 assert_eq!(ErrorKind::AlreadyExists, error.kind());
801 }
802
803 #[tokio::test]
804 async fn test_list_missing_or_non_directory_is_empty() {
805 let temp_dir = create_temp_dir("secure_fs_list_missing_or_non_directory");
806 std::fs::write(temp_dir.path().join("file"), []).unwrap();
807 let operator = SecureFsRoot::open(temp_dir.path())
808 .unwrap()
809 .build_operator();
810
811 assert!(operator.list("missing/").await.unwrap().is_empty());
812 assert!(operator.list("file/").await.unwrap().is_empty());
813 }
814
815 #[test]
824 fn test_lister_skips_entry_removed_during_iteration() {
825 let temp_dir = create_temp_dir("secure_fs_lister_removed_entry");
826 let path = temp_dir.path().join("removed");
827 std::fs::write(&path, []).unwrap();
828 let root = SecureFsRoot::open(temp_dir.path()).unwrap();
829 let mut read_dir = root.dir.entries().unwrap();
830 let entry = read_dir.next().unwrap().unwrap();
831 std::fs::remove_file(path).unwrap();
832
833 let result = read_list_entry(entry, "").unwrap();
834
835 if let Some(entry) = &result {
836 assert_eq!("removed", entry.path());
837 }
838
839 #[cfg(not(windows))]
840 assert!(result.is_none());
841 }
842
843 #[tokio::test]
844 async fn test_delete_root_is_unsupported() {
845 let temp_dir = create_temp_dir("secure_fs_delete_root");
846 let operator = SecureFsRoot::open(temp_dir.path())
847 .unwrap()
848 .build_operator();
849 operator
850 .write("nested/file", Bytes::from_static(b"data"))
851 .await
852 .unwrap();
853
854 let error = operator.delete_with("/").recursive(true).await.unwrap_err();
855
856 assert_eq!(ErrorKind::Unsupported, error.kind());
857 assert!(temp_dir.path().join("nested/file").exists());
858
859 operator
860 .delete_with("nested/")
861 .recursive(true)
862 .await
863 .unwrap();
864 assert!(!temp_dir.path().join("nested").exists());
865 }
866
867 #[tokio::test]
868 async fn test_conditional_abort_only_removes_owned_partial_file() {
869 let temp_dir = create_temp_dir("secure_fs_conditional_abort");
870 let operator = SecureFsRoot::open(temp_dir.path())
871 .unwrap()
872 .build_operator();
873 for started in [false, true] {
874 let mut writer = operator
875 .writer_with("partial")
876 .if_not_exists(true)
877 .await
878 .unwrap();
879 if started {
880 writer.write(Bytes::from_static(b"partial")).await.unwrap();
881 }
882 writer.abort().await.unwrap();
883 assert!(!operator.exists("partial").await.unwrap());
884 }
885 }
886
887 #[tokio::test]
888 async fn test_overwrite_abort_preserves_unopened_destination() {
889 use std::path::PathBuf;
890
891 use opendal::raw::OpWrite;
892 use opendal::raw::oio::Write;
893
894 let temp_dir = create_temp_dir("secure_fs_unopened_abort");
895 std::fs::write(temp_dir.path().join("existing"), b"original").unwrap();
896 let mut writer = super::SecureFsWriter {
897 root: SecureFsRoot::open(temp_dir.path()).unwrap(),
898 path: PathBuf::from("existing"),
899 args: OpWrite::default(),
900 file: None,
901 synced: false,
902 };
903
904 writer.abort().await.unwrap();
905 assert_eq!(
906 std::fs::read(temp_dir.path().join("existing")).unwrap(),
907 b"original"
908 );
909 }
910
911 #[cfg(target_os = "linux")]
912 #[tokio::test]
913 async fn test_abort_after_background_write_error() {
914 use std::path::PathBuf;
915
916 use opendal::options::WriteOptions;
917 use opendal::raw::OpWrite;
918 use opendal::raw::oio::Write;
919 use tokio::io::AsyncWriteExt;
920
921 for if_not_exists in [true, false] {
922 let temp_dir = create_temp_dir("secure_fs_flush_error_abort");
923 let root = SecureFsRoot::open(temp_dir.path()).unwrap();
924 let (args, _) = OpWrite::from_options(
925 &root.build_operator().info().capability(),
926 WriteOptions {
927 if_not_exists,
928 ..Default::default()
929 },
930 )
931 .unwrap();
932 let mut writer = super::SecureFsWriter {
933 root,
934 path: PathBuf::from("partial"),
935 args,
936 file: None,
937 synced: false,
938 };
939 writer.ensure_file().await.unwrap();
940 let mut failing_file = tokio::fs::OpenOptions::new()
941 .write(true)
942 .open("/dev/full")
943 .await
944 .unwrap();
945 failing_file.write_all(b"partial").await.unwrap();
946 writer.file = Some(failing_file);
947
948 let error = writer.abort().await.unwrap_err();
949 if if_not_exists {
950 assert!(!temp_dir.path().join("partial").exists());
951 } else {
952 assert_eq!(error.kind(), opendal::ErrorKind::Unsupported);
953 assert!(std::error::Error::source(&error).is_some());
954 }
955 }
956 }
957
958 #[cfg(target_os = "linux")]
959 #[tokio::test]
960 async fn test_conditional_abort_after_close_flush_error() {
961 use std::path::PathBuf;
962
963 use opendal::options::WriteOptions;
964 use opendal::raw::OpWrite;
965 use opendal::raw::oio::Write;
966 use tokio::io::AsyncWriteExt;
967
968 let temp_dir = create_temp_dir("secure_fs_close_flush_error");
969 let root = SecureFsRoot::open(temp_dir.path()).unwrap();
970 let (args, _) = OpWrite::from_options(
971 &root.build_operator().info().capability(),
972 WriteOptions {
973 if_not_exists: true,
974 ..Default::default()
975 },
976 )
977 .unwrap();
978 let mut writer = super::SecureFsWriter {
979 root,
980 path: PathBuf::from("partial"),
981 args,
982 file: None,
983 synced: false,
984 };
985 writer.ensure_file().await.unwrap();
986 let mut failing_file = tokio::fs::OpenOptions::new()
987 .write(true)
988 .open("/dev/full")
989 .await
990 .unwrap();
991 failing_file.write_all(b"partial").await.unwrap();
992 writer.file = Some(failing_file);
993
994 assert!(writer.close().await.is_err());
995 writer.abort().await.unwrap();
996 assert!(!temp_dir.path().join("partial").exists());
997 }
998
999 #[test]
1000 fn test_writer_abort_drains_before_reporting_unsupported() {
1001 use std::path::PathBuf;
1002
1003 use opendal::Buffer;
1004 use opendal::raw::OpWrite;
1005 use opendal::raw::oio::Write;
1006
1007 use super::SecureFsWriter;
1008
1009 tokio::runtime::Builder::new_current_thread()
1010 .enable_all()
1011 .max_blocking_threads(1)
1012 .build()
1013 .unwrap()
1014 .block_on(async {
1015 let temp_dir = create_temp_dir("secure_fs_writer_abort");
1016 let mut writer = SecureFsWriter {
1017 root: SecureFsRoot::open(temp_dir.path()).unwrap(),
1018 path: PathBuf::from("partial"),
1019 args: OpWrite::default(),
1020 file: None,
1021 synced: false,
1022 };
1023 writer.ensure_file().await.unwrap();
1024 let (started, ready) = tokio::sync::oneshot::channel();
1025 let (release, blocked) = std::sync::mpsc::channel();
1026 let blocking = tokio::task::spawn_blocking(move || {
1027 started.send(()).unwrap();
1028 blocked.recv().unwrap();
1029 });
1030 ready.await.unwrap();
1031 writer.write(Buffer::from("partial")).await.unwrap();
1032 let abort = writer.abort();
1033 tokio::pin!(abort);
1034 let pending = futures::poll!(&mut abort).is_pending();
1035 release.send(()).unwrap();
1036 assert!(pending);
1037 assert_eq!(ErrorKind::Unsupported, abort.await.unwrap_err().kind());
1038 blocking.await.unwrap();
1039 assert_eq!(
1040 std::fs::read(temp_dir.path().join("partial")).unwrap(),
1041 b"partial"
1042 );
1043 });
1044 }
1045}