1use std::sync::Arc;
11
12use crate::copy::AsyncCopyFailure;
13use crate::copy::AsyncCopyOperation;
14use crate::copy::CopyAssessment;
15use crate::copy::CopyFailureState;
16use crate::copy::CopyOptions;
17use crate::copy::CopyOutcome;
18use crate::copy::CopyStats;
19use crate::directory::AsyncDirectoryOperation;
20use crate::directory::AsyncDirectoryStream;
21use crate::directory::CreateDirectoryOptions;
22use crate::directory::CreateDirectoryOutcome;
23use crate::directory::DeleteOptions;
24use crate::directory::DeleteOutcome;
25use crate::directory::ListOptions;
26use crate::directory::ListScope;
27use crate::error::FsError;
28use crate::error::FsErrorKind;
29use crate::error::FsOperation;
30use crate::error::FsResult;
31use crate::error::OpenFailure;
32use crate::error::OpenFailureStage;
33use crate::facade::facade_core::FacadeCore;
34use crate::metadata::FileMetadata;
35use crate::metadata::FileSystemCapability;
36use crate::metadata::FileSystemProperties;
37use crate::path::Path;
38use crate::read::AsyncFileReader;
39use crate::read::AsyncReadOperation;
40use crate::read::ReadOptions;
41use crate::rename::RenameFailure;
42use crate::rename::RenameFailureState;
43use crate::rename::RenameOptions;
44use crate::rename::RenameOutcome;
45use crate::rename::validate_rename_outcome;
46use crate::spi::AsyncFileSystemSpi;
47use crate::spi::CreateDirectoryRequest;
48use crate::spi::DeleteDirectoryRequest;
49use crate::spi::DeleteFileRequest;
50use crate::spi::OpenReaderRequest;
51use crate::spi::OpenWriterRequest;
52use crate::spi::RenameRequest;
53use crate::spi::ResolvedCreateDirectoryOptions;
54use crate::spi::ResolvedDeleteOptions;
55use crate::spi::ResolvedReadOptions;
56use crate::spi::ResolvedRenameOptions;
57use crate::spi::ResolvedWriteOptions;
58use crate::spi::StatRequest;
59use crate::temp::AsyncTempDirectory;
60use crate::temp::AsyncTempFile;
61use crate::temp::PersistOptions;
62use crate::temp::RejectedAsyncTempResource;
63use crate::temp::TempOptions;
64use crate::write::AsyncFileWriter;
65use crate::write::AsyncWriteAllOperation;
66use crate::write::AsyncWriteAllOperationFailure;
67use crate::write::RejectedAsyncWriter;
68use crate::write::WriteOptions;
69
70#[derive(Clone)]
94pub struct AsyncFileSystem {
95 spi: Arc<dyn AsyncFileSystemSpi>,
97 core: Arc<FacadeCore>,
99}
100
101impl AsyncFileSystem {
102 #[inline]
104 pub fn from_spi<S>(spi: S) -> FsResult<Self>
105 where
106 S: AsyncFileSystemSpi + 'static,
107 {
108 Self::from_shared_spi(Arc::new(spi))
109 }
110
111 #[inline]
113 pub fn from_shared_spi(spi: Arc<dyn AsyncFileSystemSpi>) -> FsResult<Self> {
114 let core = FacadeCore::new(spi.properties())?;
115 Ok(Self {
116 spi,
117 core: Arc::new(core),
118 })
119 }
120
121 #[inline]
123 #[must_use]
124 pub fn properties(&self) -> &FileSystemProperties {
125 self.core.properties()
126 }
127
128 pub fn validate_write(&self, path: &Path, options: &WriteOptions) -> FsResult<()> {
130 self.core.validate_write_request(path, options)
131 }
132
133 pub fn assess_copy(&self, source: &Path, target: &Path, options: &CopyOptions) -> FsResult<CopyAssessment> {
139 self.core.assess_copy(source, target, options)
140 }
141
142 #[inline]
144 pub(crate) fn core(&self) -> &FacadeCore {
145 &self.core
146 }
147
148 #[inline]
150 pub(crate) fn spi(&self) -> &dyn AsyncFileSystemSpi {
151 self.spi.as_ref()
152 }
153
154 pub async fn stat(&self, path: &Path) -> FsResult<FileMetadata> {
164 self.validate_path(path, FsOperation::Stat)?;
165 let response = self
166 .spi
167 .stat(StatRequest::new(path, ()))
168 .await
169 .map_err(|error| self.enrich(error, path, FsOperation::Stat))?;
170 if response.path() != path {
171 return Err(self.contract_error(path, "provider returned metadata for a different path"));
172 }
173 Ok(response.into_metadata())
174 }
175
176 pub async fn exists(&self, path: &Path) -> FsResult<bool> {
182 match self.stat(path).await {
183 Ok(_) => Ok(true),
184 Err(error) if error.kind() == FsErrorKind::NotFound => Ok(false),
185 Err(error) => Err(error.with_operation(FsOperation::Exists)),
186 }
187 }
188
189 pub async fn list(&self, scope: &ListScope, options: ListOptions) -> FsResult<AsyncDirectoryStream> {
195 AsyncDirectoryOperation::new(self).list(scope, options).await
196 }
197
198 pub async fn open_reader(&self, path: &Path, options: ReadOptions) -> FsResult<AsyncFileReader> {
204 self.open_reader_resolved(path, ResolvedReadOptions::new(options)).await
205 }
206
207 pub(crate) async fn open_reader_resolved(
209 &self,
210 path: &Path,
211 resolved: ResolvedReadOptions,
212 ) -> FsResult<AsyncFileReader> {
213 let options = resolved.options().clone();
214 self.validate_path(path, FsOperation::OpenReader)?;
215 options
216 .validate_against(self.properties().capabilities())
217 .map_err(|error| self.enrich(error, path, FsOperation::OpenReader))?;
218 self.properties()
219 .limits()
220 .validate_read_range(path, options.length())
221 .map_err(|error| self.enrich(error, path, FsOperation::OpenReader))?;
222 self.require(FileSystemCapability::Read, FsOperation::OpenReader, path)?;
223 let opened = self
224 .spi
225 .open_reader(OpenReaderRequest::new(path, resolved))
226 .await
227 .map_err(|error| self.enrich(error, path, FsOperation::OpenReader))?;
228 self.validate_opened_info(opened.info(), path)?;
229 Ok(opened.into_reader())
230 }
231
232 pub async fn read_all(&self, path: &Path, options: ReadOptions, max_bytes: usize) -> FsResult<Vec<u8>> {
241 AsyncReadOperation::new(self).read_all(path, options, max_bytes).await
242 }
243
244 pub async fn read_prefix(
250 &self,
251 path: &Path,
252 options: ReadOptions,
253 max_bytes: usize,
254 ) -> FsResult<crate::read::PrefixReadOutcome> {
255 AsyncReadOperation::new(self)
256 .read_prefix(path, options, max_bytes)
257 .await
258 }
259
260 pub async fn open_writer(
262 &self,
263 path: &Path,
264 options: WriteOptions,
265 ) -> Result<AsyncFileWriter, OpenFailure<RejectedAsyncWriter>> {
266 self.core
267 .validate_write_request(path, &options)
268 .map_err(|error| OpenFailure::new(error, OpenFailureStage::Preflight, None))?;
269 let atomicity = options.atomicity();
270 let durability = options.durability();
271 let opened = self
272 .spi
273 .open_writer(OpenWriterRequest::new(path, ResolvedWriteOptions::new(options)))
274 .await
275 .map_err(|error| {
276 OpenFailure::new(
277 self.enrich(error, path, FsOperation::OpenWriter),
278 OpenFailureStage::ProviderOpen,
279 None,
280 )
281 })?;
282 let (info, session) = opened.into_parts();
283 if let Err(error) = self.validate_opened_info(&info, path) {
284 return Err(OpenFailure::new(
285 error,
286 OpenFailureStage::OutcomeValidation,
287 Some(RejectedAsyncWriter::new(
288 session,
289 self.properties().info().provider_id(),
290 Some(path.clone()),
291 )),
292 ));
293 }
294 Ok(AsyncFileWriter::new(
295 info,
296 session,
297 atomicity,
298 durability,
299 self.properties().info().provider_id(),
300 self.properties().limits().max_write_bytes().maximum(),
301 ))
302 }
303
304 pub fn begin_write_all(
323 &self,
324 path: Path,
325 bytes: Vec<u8>,
326 options: WriteOptions,
327 ) -> Result<AsyncWriteAllOperation, AsyncWriteAllOperationFailure> {
328 self.core.validate_write_request(&path, &options).map_err(|error| {
329 AsyncWriteAllOperationFailure::new(error, crate::write::WriteFailureState::NotPublished, 0)
330 })?;
331 self.properties()
332 .limits()
333 .validate_write_size(&path, bytes.len())
334 .map_err(|error| {
335 AsyncWriteAllOperationFailure::new(
336 self.core.enrich(error, Some(&path), FsOperation::Write),
337 crate::write::WriteFailureState::NotPublished,
338 0,
339 )
340 })?;
341 Ok(AsyncWriteAllOperation::new(self.clone(), path, bytes, options))
342 }
343
344 pub async fn create_directory(
350 &self,
351 path: &Path,
352 options: CreateDirectoryOptions,
353 ) -> FsResult<CreateDirectoryOutcome> {
354 self.validate_path(path, FsOperation::CreateDir)?;
355 self.require(FileSystemCapability::CreateDirectory, FsOperation::CreateDir, path)?;
356 let exists_ok = options.exists_ok();
357 let outcome = self
358 .spi
359 .create_directory(CreateDirectoryRequest::new(
360 path,
361 ResolvedCreateDirectoryOptions::new(options),
362 ))
363 .await
364 .map_err(|error| self.enrich(error, path, FsOperation::CreateDir))?;
365 if outcome.already_existed() && !exists_ok {
366 return Err(self.contract_error(path, "provider accepted an existing directory without exists_ok"));
367 }
368 Ok(outcome)
369 }
370
371 #[inline]
377 pub async fn delete_file(&self, path: &Path, options: DeleteOptions) -> FsResult<DeleteOutcome> {
378 self.delete(path, options, false).await
379 }
380
381 #[inline]
387 pub async fn delete_directory(&self, path: &Path, options: DeleteOptions) -> FsResult<DeleteOutcome> {
388 self.delete(path, options, true).await
389 }
390
391 pub async fn rename(
397 &self,
398 source: &Path,
399 target: &Path,
400 options: RenameOptions,
401 ) -> Result<RenameOutcome, RenameFailure> {
402 if let Err(error) = self.rename_preflight(source, target, &options) {
403 return Err(self.contextual_rename_failure(error, RenameFailureState::Unchanged, source, target));
404 }
405 match self
406 .spi
407 .rename(RenameRequest::new(
408 source,
409 target,
410 ResolvedRenameOptions::new(options.clone()),
411 ))
412 .await
413 {
414 Ok(outcome) => match validate_rename_outcome(&outcome, &options, source, target) {
415 Some(violation) => Err(self.contextual_rename_failure(
416 self.contract_error(source, violation.message),
417 violation.state,
418 source,
419 target,
420 )),
421 None => Ok(outcome),
422 },
423 Err(failure) => {
424 let (error, state) = failure.into_parts();
425 Err(self.contextual_rename_failure(error, state, source, target))
426 }
427 }
428 }
429
430 pub async fn create_temp_file(
437 &self,
438 options: TempOptions,
439 ) -> Result<AsyncTempFile, OpenFailure<RejectedAsyncTempResource>> {
440 let parent = options.parent().cloned();
441 self.core
442 .validate_temp_parent(parent.as_ref())
443 .map_err(|error| OpenFailure::new(error, OpenFailureStage::Preflight, None))?;
444 self.core
445 .require(FileSystemCapability::TempFile, FsOperation::CreateTemp, parent.as_ref())
446 .map_err(|error| OpenFailure::new(error, OpenFailureStage::Preflight, None))?;
447 let opened = self
448 .spi
449 .create_temp_file(crate::spi::CreateTempFileRequest::new(options))
450 .await
451 .map_err(|error| {
452 OpenFailure::new(
453 self.core.enrich(error, parent.as_ref(), FsOperation::CreateTemp),
454 OpenFailureStage::ProviderOpen,
455 None,
456 )
457 })?;
458 let (info, session) = opened.into_parts();
459 if let Err(cause) = self.validate_temp_info(&info, crate::metadata::FileKind::File) {
460 let error = FsError::with_source(
461 FsErrorKind::ProviderContractViolation,
462 FsOperation::ValidateProviderOutcome,
463 "provider returned an invalid temporary identity",
464 cause,
465 )
466 .with_provider(self.properties().info().provider_id());
467 let error = match parent.as_ref() {
468 Some(path) => error.with_path(path.clone()),
469 None => error,
470 };
471 return Err(OpenFailure::new(
472 error,
473 OpenFailureStage::OutcomeValidation,
474 Some(RejectedAsyncTempResource::new(
475 session,
476 self.properties().info().provider_id(),
477 parent,
478 )),
479 ));
480 }
481 Ok(AsyncTempFile::new(
482 self.clone(),
483 info.path().clone(),
484 session,
485 "temporary file",
486 ))
487 }
488
489 pub async fn create_temp_directory(
495 &self,
496 options: TempOptions,
497 ) -> Result<AsyncTempDirectory, OpenFailure<RejectedAsyncTempResource>> {
498 let parent = options.parent().cloned();
499 self.core
500 .validate_temp_parent(parent.as_ref())
501 .map_err(|error| OpenFailure::new(error, OpenFailureStage::Preflight, None))?;
502 self.core
503 .require(
504 FileSystemCapability::TempDirectory,
505 FsOperation::CreateTemp,
506 parent.as_ref(),
507 )
508 .map_err(|error| OpenFailure::new(error, OpenFailureStage::Preflight, None))?;
509 let opened = self
510 .spi
511 .create_temp_directory(crate::spi::CreateTempDirectoryRequest::new(options))
512 .await
513 .map_err(|error| {
514 OpenFailure::new(
515 self.core.enrich(error, parent.as_ref(), FsOperation::CreateTemp),
516 OpenFailureStage::ProviderOpen,
517 None,
518 )
519 })?;
520 let (info, session) = opened.into_parts();
521 if let Err(cause) = self.validate_temp_info(&info, crate::metadata::FileKind::Directory) {
522 let error = FsError::with_source(
523 FsErrorKind::ProviderContractViolation,
524 FsOperation::ValidateProviderOutcome,
525 "provider returned an invalid temporary identity",
526 cause,
527 )
528 .with_provider(self.properties().info().provider_id());
529 let error = match parent.as_ref() {
530 Some(path) => error.with_path(path.clone()),
531 None => error,
532 };
533 return Err(OpenFailure::new(
534 error,
535 OpenFailureStage::OutcomeValidation,
536 Some(RejectedAsyncTempResource::new(
537 session,
538 self.properties().info().provider_id(),
539 parent,
540 )),
541 ));
542 }
543 Ok(AsyncTempDirectory::new(self.clone(), info.path().clone(), session))
544 }
545
546 #[allow(clippy::result_large_err)]
554 pub fn begin_copy(
555 &self,
556 source: Path,
557 target: Path,
558 options: CopyOptions,
559 ) -> Result<AsyncCopyOperation, AsyncCopyFailure> {
560 self.copy_preflight(&source, &target, &options).map_err(|error| {
561 self.contextual_copy_failure(
562 error,
563 CopyFailureState::Unchanged,
564 CopyStats::default(),
565 &source,
566 &target,
567 )
568 })?;
569 let symlink_policy = options
570 .symlink_policy_override()
571 .unwrap_or(self.properties().symlink_policy());
572 Ok(AsyncCopyOperation::new(
573 self.clone(),
574 source,
575 target,
576 options,
577 symlink_policy,
578 ))
579 }
580
581 #[allow(clippy::result_large_err)]
583 pub(crate) fn verify_completed_copy(
584 &self,
585 outcome: CopyOutcome,
586 options: &CopyOptions,
587 source: &Path,
588 target: &Path,
589 ) -> Result<CopyOutcome, AsyncCopyFailure> {
590 if let Some(message) = outcome.contract_violation(options) {
591 return Err(self.contextual_copy_failure(
592 FsError::new(FsErrorKind::ProviderContractViolation, FsOperation::Copy, message),
593 CopyFailureState::Published,
594 *outcome.stats(),
595 source,
596 target,
597 ));
598 }
599 Ok(outcome)
600 }
601
602 fn copy_preflight(&self, source: &Path, target: &Path, options: &CopyOptions) -> FsResult<()> {
604 self.validate_path(source, FsOperation::Copy)?;
605 self.validate_path(target, FsOperation::Copy)?;
606 options
607 .validate_against(self.properties().capabilities())
608 .map_err(|error| {
609 self.enrich(error, source, FsOperation::Copy)
610 .with_target(target.clone())
611 })?;
612 if source == target {
613 return Err(FsError::new(
614 crate::error::FsErrorKind::InvalidOptions,
615 FsOperation::Copy,
616 "copy source and target must differ",
617 )
618 .with_path(source.clone())
619 .with_target(target.clone()));
620 }
621 Ok(())
622 }
623
624 fn rename_preflight(&self, source: &Path, target: &Path, options: &RenameOptions) -> FsResult<()> {
626 self.validate_path(source, FsOperation::Rename)?;
627 self.validate_path(target, FsOperation::Rename)?;
628 options
629 .validate_against(self.properties().capabilities())
630 .map_err(|error| {
631 self.enrich(error, source, FsOperation::Rename)
632 .with_target(target.clone())
633 })?;
634 self.require(FileSystemCapability::Rename, FsOperation::Rename, source)?;
635 if source == target {
636 return Err(FsError::new(
637 FsErrorKind::InvalidOptions,
638 FsOperation::Rename,
639 "rename source and target must differ",
640 )
641 .with_path(source.clone())
642 .with_target(target.clone()));
643 }
644 Ok(())
645 }
646
647 async fn delete(&self, path: &Path, options: DeleteOptions, directory: bool) -> FsResult<DeleteOutcome> {
649 self.validate_path(path, FsOperation::Delete)?;
650 options
651 .validate_against(self.properties().capabilities())
652 .map_err(|error| self.enrich(error, path, FsOperation::Delete))?;
653 self.require(FileSystemCapability::Delete, FsOperation::Delete, path)?;
654 let missing_ok = options.missing_ok();
655 let request_options = ResolvedDeleteOptions::new(options);
656 let outcome = if directory {
657 self.spi
658 .delete_directory(DeleteDirectoryRequest::new(path, request_options))
659 .await
660 } else {
661 self.spi
662 .delete_file(DeleteFileRequest::new(path, request_options))
663 .await
664 }
665 .map_err(|error| self.enrich(error, path, FsOperation::Delete))?;
666 if outcome.already_missing() && !missing_ok {
667 return Err(self.contract_error(path, "provider accepted a missing target without missing_ok"));
668 }
669 Ok(outcome)
670 }
671
672 fn validate_path(&self, path: &Path, operation: FsOperation) -> FsResult<()> {
674 self.core.validate_path(path, operation)
675 }
676
677 pub(crate) fn require(
679 &self,
680 capability: FileSystemCapability,
681 operation: FsOperation,
682 path: &Path,
683 ) -> FsResult<()> {
684 self.core.require(capability, operation, Some(path))
685 }
686
687 pub(crate) fn contextual_copy_failure(
690 &self,
691 error: FsError,
692 state: CopyFailureState,
693 stats: CopyStats,
694 source: &Path,
695 target: &Path,
696 ) -> AsyncCopyFailure {
697 AsyncCopyFailure::new(
698 error.with_operation(FsOperation::Copy).with_missing_context(
699 source,
700 Some(target),
701 self.properties().info().provider_id(),
702 ),
703 state,
704 stats,
705 )
706 }
707
708 fn contextual_rename_failure(
711 &self,
712 error: FsError,
713 state: RenameFailureState,
714 source: &Path,
715 target: &Path,
716 ) -> RenameFailure {
717 RenameFailure::new(
718 error.with_operation(FsOperation::Rename).with_missing_context(
719 source,
720 Some(target),
721 self.properties().info().provider_id(),
722 ),
723 state,
724 )
725 }
726
727 fn enrich(&self, error: FsError, path: &Path, operation: FsOperation) -> FsError {
729 self.core.enrich(error, Some(path), operation)
730 }
731
732 fn contract_error(&self, path: &Path, message: &'static str) -> FsError {
734 self.core
735 .contract_error(path, FsOperation::ValidateProviderOutcome, message)
736 }
737
738 fn validate_opened_info(&self, info: &crate::metadata::OpenedFileInfo, path: &Path) -> FsResult<()> {
741 if info.filesystem_id() != self.properties().info().id() || info.path() != path {
742 return Err(self.contract_error(path, "provider returned an opened handle with a different identity"));
743 }
744 Ok(())
745 }
746
747 fn validate_temp_info(
749 &self,
750 info: &crate::metadata::OpenedFileInfo,
751 expected_kind: crate::metadata::FileKind,
752 ) -> FsResult<()> {
753 if info.filesystem_id() != self.properties().info().id() {
754 return Err(self.contract_error(
755 info.path(),
756 "provider returned a temporary handle for a different filesystem",
757 ));
758 }
759 self.validate_path(info.path(), FsOperation::CreateTemp).map_err(|_| {
760 self.contract_error(
761 info.path(),
762 "provider returned a temporary handle with an invalid logical path",
763 )
764 })?;
765 if info.metadata().is_none_or(|metadata| metadata.kind() != &expected_kind) {
766 return Err(self.contract_error(
767 info.path(),
768 "provider returned a temporary handle with an inconsistent resource kind",
769 ));
770 }
771 Ok(())
772 }
773
774 pub(crate) fn preflight_temp_persist(
776 &self,
777 source: &Path,
778 target: &Path,
779 options: &PersistOptions,
780 ) -> FsResult<()> {
781 self.validate_path(source, FsOperation::PersistTemp)?;
782 self.validate_path(target, FsOperation::PersistTemp)?;
783 options
784 .validate_against(self.properties().capabilities())
785 .map_err(|error| {
786 self.enrich(error, source, FsOperation::PersistTemp)
787 .with_target(target.clone())
788 })
789 }
790
791 pub(crate) fn validate_temp_keep_target(&self, source: &Path, target: &Path) -> FsResult<()> {
793 self.validate_path(target, FsOperation::KeepTemp).map_err(|_| {
794 self.contract_error(
795 source,
796 "provider returned a temporary keep target with an invalid logical path",
797 )
798 .with_target(target.clone())
799 })
800 }
801}