use std::sync::Arc;
use crate::copy::AsyncCopyFailure;
use crate::copy::AsyncCopyOperation;
use crate::copy::CopyAssessment;
use crate::copy::CopyFailureState;
use crate::copy::CopyOptions;
use crate::copy::CopyOutcome;
use crate::copy::CopyStats;
use crate::directory::AsyncDirectoryOperation;
use crate::directory::AsyncDirectoryStream;
use crate::directory::CreateDirectoryOptions;
use crate::directory::CreateDirectoryOutcome;
use crate::directory::DeleteOptions;
use crate::directory::DeleteOutcome;
use crate::directory::ListOptions;
use crate::directory::ListScope;
use crate::error::FsError;
use crate::error::FsErrorKind;
use crate::error::FsOperation;
use crate::error::FsResult;
use crate::error::OpenFailure;
use crate::error::OpenFailureStage;
use crate::facade::facade_core::FacadeCore;
use crate::metadata::FileMetadata;
use crate::metadata::FileSystemCapability;
use crate::metadata::FileSystemProperties;
use crate::path::Path;
use crate::read::AsyncFileReader;
use crate::read::AsyncReadOperation;
use crate::read::ReadOptions;
use crate::rename::RenameFailure;
use crate::rename::RenameFailureState;
use crate::rename::RenameOptions;
use crate::rename::RenameOutcome;
use crate::rename::validate_rename_outcome;
use crate::spi::AsyncFileSystemSpi;
use crate::spi::CreateDirectoryRequest;
use crate::spi::DeleteDirectoryRequest;
use crate::spi::DeleteFileRequest;
use crate::spi::OpenReaderRequest;
use crate::spi::OpenWriterRequest;
use crate::spi::RenameRequest;
use crate::spi::ResolvedCreateDirectoryOptions;
use crate::spi::ResolvedDeleteOptions;
use crate::spi::ResolvedReadOptions;
use crate::spi::ResolvedRenameOptions;
use crate::spi::ResolvedWriteOptions;
use crate::spi::StatRequest;
use crate::temp::AsyncTempDirectory;
use crate::temp::AsyncTempFile;
use crate::temp::PersistOptions;
use crate::temp::RejectedAsyncTempResource;
use crate::temp::TempOptions;
use crate::write::AsyncFileWriter;
use crate::write::AsyncWriteAllOperation;
use crate::write::AsyncWriteAllOperationFailure;
use crate::write::RejectedAsyncWriter;
use crate::write::WriteOptions;
#[derive(Clone)]
pub struct AsyncFileSystem {
spi: Arc<dyn AsyncFileSystemSpi>,
core: Arc<FacadeCore>,
}
impl AsyncFileSystem {
#[inline]
pub fn from_spi<S>(spi: S) -> FsResult<Self>
where
S: AsyncFileSystemSpi + 'static,
{
Self::from_shared_spi(Arc::new(spi))
}
#[inline]
pub fn from_shared_spi(spi: Arc<dyn AsyncFileSystemSpi>) -> FsResult<Self> {
let core = FacadeCore::new(spi.properties())?;
Ok(Self {
spi,
core: Arc::new(core),
})
}
#[inline]
#[must_use]
pub fn properties(&self) -> &FileSystemProperties {
self.core.properties()
}
pub fn validate_write(&self, path: &Path, options: &WriteOptions) -> FsResult<()> {
self.core.validate_write_request(path, options)
}
pub fn assess_copy(&self, source: &Path, target: &Path, options: &CopyOptions) -> FsResult<CopyAssessment> {
self.core.assess_copy(source, target, options)
}
#[inline]
pub(crate) fn core(&self) -> &FacadeCore {
&self.core
}
#[inline]
pub(crate) fn spi(&self) -> &dyn AsyncFileSystemSpi {
self.spi.as_ref()
}
pub async fn stat(&self, path: &Path) -> FsResult<FileMetadata> {
self.validate_path(path, FsOperation::Stat)?;
let response = self
.spi
.stat(StatRequest::new(path, ()))
.await
.map_err(|error| self.enrich(error, path, FsOperation::Stat))?;
if response.path() != path {
return Err(self.contract_error(path, "provider returned metadata for a different path"));
}
Ok(response.into_metadata())
}
pub async fn exists(&self, path: &Path) -> FsResult<bool> {
match self.stat(path).await {
Ok(_) => Ok(true),
Err(error) if error.kind() == FsErrorKind::NotFound => Ok(false),
Err(error) => Err(error.with_operation(FsOperation::Exists)),
}
}
pub async fn list(&self, scope: &ListScope, options: ListOptions) -> FsResult<AsyncDirectoryStream> {
AsyncDirectoryOperation::new(self).list(scope, options).await
}
pub async fn open_reader(&self, path: &Path, options: ReadOptions) -> FsResult<AsyncFileReader> {
self.open_reader_resolved(path, ResolvedReadOptions::new(options)).await
}
pub(crate) async fn open_reader_resolved(
&self,
path: &Path,
resolved: ResolvedReadOptions,
) -> FsResult<AsyncFileReader> {
let options = resolved.options().clone();
self.validate_path(path, FsOperation::OpenReader)?;
options
.validate_against(self.properties().capabilities())
.map_err(|error| self.enrich(error, path, FsOperation::OpenReader))?;
self.properties()
.limits()
.validate_read_range(path, options.length())
.map_err(|error| self.enrich(error, path, FsOperation::OpenReader))?;
self.require(FileSystemCapability::Read, FsOperation::OpenReader, path)?;
let opened = self
.spi
.open_reader(OpenReaderRequest::new(path, resolved))
.await
.map_err(|error| self.enrich(error, path, FsOperation::OpenReader))?;
self.validate_opened_info(opened.info(), path)?;
Ok(opened.into_reader())
}
pub async fn read_all(&self, path: &Path, options: ReadOptions, max_bytes: usize) -> FsResult<Vec<u8>> {
AsyncReadOperation::new(self).read_all(path, options, max_bytes).await
}
pub async fn read_prefix(
&self,
path: &Path,
options: ReadOptions,
max_bytes: usize,
) -> FsResult<crate::read::PrefixReadOutcome> {
AsyncReadOperation::new(self)
.read_prefix(path, options, max_bytes)
.await
}
pub async fn open_writer(
&self,
path: &Path,
options: WriteOptions,
) -> Result<AsyncFileWriter, OpenFailure<RejectedAsyncWriter>> {
self.core
.validate_write_request(path, &options)
.map_err(|error| OpenFailure::new(error, OpenFailureStage::Preflight, None))?;
let atomicity = options.atomicity();
let durability = options.durability();
let opened = self
.spi
.open_writer(OpenWriterRequest::new(path, ResolvedWriteOptions::new(options)))
.await
.map_err(|error| {
OpenFailure::new(
self.enrich(error, path, FsOperation::OpenWriter),
OpenFailureStage::ProviderOpen,
None,
)
})?;
let (info, session) = opened.into_parts();
if let Err(error) = self.validate_opened_info(&info, path) {
return Err(OpenFailure::new(
error,
OpenFailureStage::OutcomeValidation,
Some(RejectedAsyncWriter::new(
session,
self.properties().info().provider_id(),
Some(path.clone()),
)),
));
}
Ok(AsyncFileWriter::new(
info,
session,
atomicity,
durability,
self.properties().info().provider_id(),
self.properties().limits().max_write_bytes().maximum(),
))
}
pub fn begin_write_all(
&self,
path: Path,
bytes: Vec<u8>,
options: WriteOptions,
) -> Result<AsyncWriteAllOperation, AsyncWriteAllOperationFailure> {
self.core.validate_write_request(&path, &options).map_err(|error| {
AsyncWriteAllOperationFailure::new(error, crate::write::WriteFailureState::NotPublished, 0)
})?;
self.properties()
.limits()
.validate_write_size(&path, bytes.len())
.map_err(|error| {
AsyncWriteAllOperationFailure::new(
self.core.enrich(error, Some(&path), FsOperation::Write),
crate::write::WriteFailureState::NotPublished,
0,
)
})?;
Ok(AsyncWriteAllOperation::new(self.clone(), path, bytes, options))
}
pub async fn create_directory(
&self,
path: &Path,
options: CreateDirectoryOptions,
) -> FsResult<CreateDirectoryOutcome> {
self.validate_path(path, FsOperation::CreateDir)?;
self.require(FileSystemCapability::CreateDirectory, FsOperation::CreateDir, path)?;
let exists_ok = options.exists_ok();
let outcome = self
.spi
.create_directory(CreateDirectoryRequest::new(
path,
ResolvedCreateDirectoryOptions::new(options),
))
.await
.map_err(|error| self.enrich(error, path, FsOperation::CreateDir))?;
if outcome.already_existed() && !exists_ok {
return Err(self.contract_error(path, "provider accepted an existing directory without exists_ok"));
}
Ok(outcome)
}
#[inline]
pub async fn delete_file(&self, path: &Path, options: DeleteOptions) -> FsResult<DeleteOutcome> {
self.delete(path, options, false).await
}
#[inline]
pub async fn delete_directory(&self, path: &Path, options: DeleteOptions) -> FsResult<DeleteOutcome> {
self.delete(path, options, true).await
}
pub async fn rename(
&self,
source: &Path,
target: &Path,
options: RenameOptions,
) -> Result<RenameOutcome, RenameFailure> {
if let Err(error) = self.rename_preflight(source, target, &options) {
return Err(self.contextual_rename_failure(error, RenameFailureState::Unchanged, source, target));
}
match self
.spi
.rename(RenameRequest::new(
source,
target,
ResolvedRenameOptions::new(options.clone()),
))
.await
{
Ok(outcome) => match validate_rename_outcome(&outcome, &options, source, target) {
Some(violation) => Err(self.contextual_rename_failure(
self.contract_error(source, violation.message),
violation.state,
source,
target,
)),
None => Ok(outcome),
},
Err(failure) => {
let (error, state) = failure.into_parts();
Err(self.contextual_rename_failure(error, state, source, target))
}
}
}
pub async fn create_temp_file(
&self,
options: TempOptions,
) -> Result<AsyncTempFile, OpenFailure<RejectedAsyncTempResource>> {
let parent = options.parent().cloned();
self.core
.validate_temp_parent(parent.as_ref())
.map_err(|error| OpenFailure::new(error, OpenFailureStage::Preflight, None))?;
self.core
.require(FileSystemCapability::TempFile, FsOperation::CreateTemp, parent.as_ref())
.map_err(|error| OpenFailure::new(error, OpenFailureStage::Preflight, None))?;
let opened = self
.spi
.create_temp_file(crate::spi::CreateTempFileRequest::new(options))
.await
.map_err(|error| {
OpenFailure::new(
self.core.enrich(error, parent.as_ref(), FsOperation::CreateTemp),
OpenFailureStage::ProviderOpen,
None,
)
})?;
let (info, session) = opened.into_parts();
if let Err(cause) = self.validate_temp_info(&info, crate::metadata::FileKind::File) {
let error = FsError::with_source(
FsErrorKind::ProviderContractViolation,
FsOperation::ValidateProviderOutcome,
"provider returned an invalid temporary identity",
cause,
)
.with_provider(self.properties().info().provider_id());
let error = match parent.as_ref() {
Some(path) => error.with_path(path.clone()),
None => error,
};
return Err(OpenFailure::new(
error,
OpenFailureStage::OutcomeValidation,
Some(RejectedAsyncTempResource::new(
session,
self.properties().info().provider_id(),
parent,
)),
));
}
Ok(AsyncTempFile::new(
self.clone(),
info.path().clone(),
session,
"temporary file",
))
}
pub async fn create_temp_directory(
&self,
options: TempOptions,
) -> Result<AsyncTempDirectory, OpenFailure<RejectedAsyncTempResource>> {
let parent = options.parent().cloned();
self.core
.validate_temp_parent(parent.as_ref())
.map_err(|error| OpenFailure::new(error, OpenFailureStage::Preflight, None))?;
self.core
.require(
FileSystemCapability::TempDirectory,
FsOperation::CreateTemp,
parent.as_ref(),
)
.map_err(|error| OpenFailure::new(error, OpenFailureStage::Preflight, None))?;
let opened = self
.spi
.create_temp_directory(crate::spi::CreateTempDirectoryRequest::new(options))
.await
.map_err(|error| {
OpenFailure::new(
self.core.enrich(error, parent.as_ref(), FsOperation::CreateTemp),
OpenFailureStage::ProviderOpen,
None,
)
})?;
let (info, session) = opened.into_parts();
if let Err(cause) = self.validate_temp_info(&info, crate::metadata::FileKind::Directory) {
let error = FsError::with_source(
FsErrorKind::ProviderContractViolation,
FsOperation::ValidateProviderOutcome,
"provider returned an invalid temporary identity",
cause,
)
.with_provider(self.properties().info().provider_id());
let error = match parent.as_ref() {
Some(path) => error.with_path(path.clone()),
None => error,
};
return Err(OpenFailure::new(
error,
OpenFailureStage::OutcomeValidation,
Some(RejectedAsyncTempResource::new(
session,
self.properties().info().provider_id(),
parent,
)),
));
}
Ok(AsyncTempDirectory::new(self.clone(), info.path().clone(), session))
}
#[allow(clippy::result_large_err)]
pub fn begin_copy(
&self,
source: Path,
target: Path,
options: CopyOptions,
) -> Result<AsyncCopyOperation, AsyncCopyFailure> {
self.copy_preflight(&source, &target, &options).map_err(|error| {
self.contextual_copy_failure(
error,
CopyFailureState::Unchanged,
CopyStats::default(),
&source,
&target,
)
})?;
let symlink_policy = options
.symlink_policy_override()
.unwrap_or(self.properties().symlink_policy());
Ok(AsyncCopyOperation::new(
self.clone(),
source,
target,
options,
symlink_policy,
))
}
#[allow(clippy::result_large_err)]
pub(crate) fn verify_completed_copy(
&self,
outcome: CopyOutcome,
options: &CopyOptions,
source: &Path,
target: &Path,
) -> Result<CopyOutcome, AsyncCopyFailure> {
if let Some(message) = outcome.contract_violation(options) {
return Err(self.contextual_copy_failure(
FsError::new(FsErrorKind::ProviderContractViolation, FsOperation::Copy, message),
CopyFailureState::Published,
*outcome.stats(),
source,
target,
));
}
Ok(outcome)
}
fn copy_preflight(&self, source: &Path, target: &Path, options: &CopyOptions) -> FsResult<()> {
self.validate_path(source, FsOperation::Copy)?;
self.validate_path(target, FsOperation::Copy)?;
options
.validate_against(self.properties().capabilities())
.map_err(|error| {
self.enrich(error, source, FsOperation::Copy)
.with_target(target.clone())
})?;
if source == target {
return Err(FsError::new(
crate::error::FsErrorKind::InvalidOptions,
FsOperation::Copy,
"copy source and target must differ",
)
.with_path(source.clone())
.with_target(target.clone()));
}
Ok(())
}
fn rename_preflight(&self, source: &Path, target: &Path, options: &RenameOptions) -> FsResult<()> {
self.validate_path(source, FsOperation::Rename)?;
self.validate_path(target, FsOperation::Rename)?;
options
.validate_against(self.properties().capabilities())
.map_err(|error| {
self.enrich(error, source, FsOperation::Rename)
.with_target(target.clone())
})?;
self.require(FileSystemCapability::Rename, FsOperation::Rename, source)?;
if source == target {
return Err(FsError::new(
FsErrorKind::InvalidOptions,
FsOperation::Rename,
"rename source and target must differ",
)
.with_path(source.clone())
.with_target(target.clone()));
}
Ok(())
}
async fn delete(&self, path: &Path, options: DeleteOptions, directory: bool) -> FsResult<DeleteOutcome> {
self.validate_path(path, FsOperation::Delete)?;
options
.validate_against(self.properties().capabilities())
.map_err(|error| self.enrich(error, path, FsOperation::Delete))?;
self.require(FileSystemCapability::Delete, FsOperation::Delete, path)?;
let missing_ok = options.missing_ok();
let request_options = ResolvedDeleteOptions::new(options);
let outcome = if directory {
self.spi
.delete_directory(DeleteDirectoryRequest::new(path, request_options))
.await
} else {
self.spi
.delete_file(DeleteFileRequest::new(path, request_options))
.await
}
.map_err(|error| self.enrich(error, path, FsOperation::Delete))?;
if outcome.already_missing() && !missing_ok {
return Err(self.contract_error(path, "provider accepted a missing target without missing_ok"));
}
Ok(outcome)
}
fn validate_path(&self, path: &Path, operation: FsOperation) -> FsResult<()> {
self.core.validate_path(path, operation)
}
pub(crate) fn require(
&self,
capability: FileSystemCapability,
operation: FsOperation,
path: &Path,
) -> FsResult<()> {
self.core.require(capability, operation, Some(path))
}
pub(crate) fn contextual_copy_failure(
&self,
error: FsError,
state: CopyFailureState,
stats: CopyStats,
source: &Path,
target: &Path,
) -> AsyncCopyFailure {
AsyncCopyFailure::new(
error.with_operation(FsOperation::Copy).with_missing_context(
source,
Some(target),
self.properties().info().provider_id(),
),
state,
stats,
)
}
fn contextual_rename_failure(
&self,
error: FsError,
state: RenameFailureState,
source: &Path,
target: &Path,
) -> RenameFailure {
RenameFailure::new(
error.with_operation(FsOperation::Rename).with_missing_context(
source,
Some(target),
self.properties().info().provider_id(),
),
state,
)
}
fn enrich(&self, error: FsError, path: &Path, operation: FsOperation) -> FsError {
self.core.enrich(error, Some(path), operation)
}
fn contract_error(&self, path: &Path, message: &'static str) -> FsError {
self.core
.contract_error(path, FsOperation::ValidateProviderOutcome, message)
}
fn validate_opened_info(&self, info: &crate::metadata::OpenedFileInfo, path: &Path) -> FsResult<()> {
if info.filesystem_id() != self.properties().info().id() || info.path() != path {
return Err(self.contract_error(path, "provider returned an opened handle with a different identity"));
}
Ok(())
}
fn validate_temp_info(
&self,
info: &crate::metadata::OpenedFileInfo,
expected_kind: crate::metadata::FileKind,
) -> FsResult<()> {
if info.filesystem_id() != self.properties().info().id() {
return Err(self.contract_error(
info.path(),
"provider returned a temporary handle for a different filesystem",
));
}
self.validate_path(info.path(), FsOperation::CreateTemp).map_err(|_| {
self.contract_error(
info.path(),
"provider returned a temporary handle with an invalid logical path",
)
})?;
if info.metadata().is_none_or(|metadata| metadata.kind() != &expected_kind) {
return Err(self.contract_error(
info.path(),
"provider returned a temporary handle with an inconsistent resource kind",
));
}
Ok(())
}
pub(crate) fn preflight_temp_persist(
&self,
source: &Path,
target: &Path,
options: &PersistOptions,
) -> FsResult<()> {
self.validate_path(source, FsOperation::PersistTemp)?;
self.validate_path(target, FsOperation::PersistTemp)?;
options
.validate_against(self.properties().capabilities())
.map_err(|error| {
self.enrich(error, source, FsOperation::PersistTemp)
.with_target(target.clone())
})
}
pub(crate) fn validate_temp_keep_target(&self, source: &Path, target: &Path) -> FsResult<()> {
self.validate_path(target, FsOperation::KeepTemp).map_err(|_| {
self.contract_error(
source,
"provider returned a temporary keep target with an invalid logical path",
)
.with_target(target.clone())
})
}
}