use std::pin::Pin;
use crate::AsyncFileSystem;
use crate::error::FsError;
use crate::error::FsErrorKind;
use crate::error::FsOperation;
use crate::error::FsResult;
use crate::metadata::AchievedAtomicity;
use crate::metadata::AtomicityRequirement;
use crate::path::Path;
use crate::path::PathComponent;
use crate::spi::AsyncTempResourceSpi;
use crate::spi::PersistRequest;
use crate::spi::SpiFuture;
use crate::temp::PersistFailure;
use crate::temp::PersistFailureState;
use crate::temp::PersistOptions;
use crate::temp::PersistOutcome;
use crate::temp::TempResourceState;
use crate::temp::internal::TempLifecycle;
pub struct AsyncTempFile {
file_system: AsyncFileSystem,
path: Path,
session: Pin<Box<dyn AsyncTempResourceSpi>>,
lifecycle: TempLifecycle,
resource_name: &'static str,
}
impl AsyncTempFile {
pub(crate) fn new(
file_system: AsyncFileSystem,
path: Path,
session: Box<dyn AsyncTempResourceSpi>,
resource_name: &'static str,
) -> Self {
Self {
file_system,
path,
session: Box::into_pin(session),
lifecycle: TempLifecycle::new(),
resource_name,
}
}
#[inline]
#[must_use]
pub const fn path(&self) -> &Path {
&self.path
}
#[inline]
#[must_use]
pub const fn state(&self) -> TempResourceState {
self.lifecycle.state()
}
#[inline]
#[must_use]
pub fn child(&self, component: &PathComponent) -> Path {
self.path.child(component)
}
#[inline]
#[must_use]
pub fn descendant(&self, relative: &crate::path::RelativePath) -> Path {
self.path.join(relative)
}
#[inline]
pub fn cleanup(&mut self) -> SpiFuture<'_, FsResult<()>> {
self.lifecycle("cannot be cleaned now", FsOperation::CleanupTemp, |session| {
session.cleanup()
})
}
#[inline]
pub fn keep(&mut self) -> SpiFuture<'_, Result<PersistOutcome, PersistFailure>> {
if self.lifecycle.state() != TempResourceState::Owned {
let error = self.invalid_state(FsOperation::KeepTemp, "cannot be kept now");
return Box::pin(async move {
Err(PersistFailure::new(error, self.lifecycle.failure_state())
.with_publication_target(self.lifecycle.publication_target()))
});
}
Box::pin(async move {
self.lifecycle.begin_pending();
match self.session.as_mut().keep().await {
Ok(outcome) => {
if let Err(error) = self.file_system.validate_temp_keep_target(&self.path, outcome.target()) {
return Err(PersistFailure::new(error, PersistFailureState::Indeterminate));
}
self.path = outcome.target().clone();
self.lifecycle.record_success(true, outcome.target().clone());
Ok(outcome)
}
Err(failure) => {
let (error, state) = failure.into_parts();
self.lifecycle.record_failure(state, error.target().cloned(), true);
let target = self.path.clone();
Err(PersistFailure::new(
error.with_operation(FsOperation::KeepTemp).with_missing_context(
&self.path,
Some(&target),
self.file_system.properties().info().provider_id(),
),
state,
)
.with_publication_target(self.lifecycle.publication_target()))
}
}
})
}
pub fn persist<'a>(
&'a mut self,
target: &'a Path,
options: PersistOptions,
) -> SpiFuture<'a, Result<PersistOutcome, PersistFailure>> {
if self.lifecycle.state() != TempResourceState::Owned {
let error = self.invalid_state(FsOperation::PersistTemp, "cannot be persisted now");
return Box::pin(async move {
Err(PersistFailure::new(error, self.lifecycle.failure_state())
.with_publication_target(self.lifecycle.publication_target()))
});
}
if let Err(error) = self.file_system.preflight_temp_persist(&self.path, target, &options) {
return Box::pin(async move { Err(PersistFailure::new(error, PersistFailureState::NotPublished)) });
}
Box::pin(async move {
self.lifecycle.begin_pending();
let atomicity = options.atomicity();
let result = self
.session
.as_mut()
.persist(PersistRequest::new(target, options))
.await;
match &result {
Ok(outcome)
if outcome.target() == target
&& !(atomicity == AtomicityRequirement::Required
&& outcome.atomicity() != AchievedAtomicity::Atomic) =>
{
self.lifecycle.record_success(false, outcome.target().clone());
}
Ok(outcome) if outcome.target() != target => {
self.lifecycle
.record_failure(PersistFailureState::Indeterminate, Some(target.clone()), false);
}
Ok(_) => self.lifecycle.record_failure(
PersistFailureState::PublishedSourceRetained,
Some(target.clone()),
false,
),
Err(failure) => self
.lifecycle
.record_failure(failure.state(), Some(target.clone()), false),
}
match result {
Ok(outcome) if outcome.target() != target => Err(PersistFailure::new(
FsError::new(
FsErrorKind::ProviderContractViolation,
FsOperation::PersistTemp,
"provider reported a persistence target different from the request",
)
.with_path(self.path.clone())
.with_target(target.clone()),
PersistFailureState::Indeterminate,
)),
Ok(outcome)
if atomicity == AtomicityRequirement::Required
&& outcome.atomicity() != AchievedAtomicity::Atomic =>
{
Err(PersistFailure::new(
FsError::new(
FsErrorKind::ProviderContractViolation,
FsOperation::PersistTemp,
"provider reported non-atomic success for atomic-required persist",
)
.with_path(self.path.clone())
.with_target(target.clone()),
PersistFailureState::PublishedSourceRetained,
)
.with_publication_target(self.lifecycle.publication_target()))
}
Err(failure) => {
let (error, state) = failure.into_parts();
Err(PersistFailure::new(self.contextual_persist_error(error, target), state)
.with_publication_target(self.lifecycle.publication_target()))
}
Ok(outcome) => Ok(outcome),
}
})
}
fn lifecycle<'a, F>(
&'a mut self,
action: &'static str,
operation: FsOperation,
call: F,
) -> SpiFuture<'a, FsResult<()>>
where
F: FnOnce(Pin<&'a mut dyn AsyncTempResourceSpi>) -> SpiFuture<'a, FsResult<()>> + Send + 'a,
{
if !matches!(
self.lifecycle.state(),
TempResourceState::Owned | TempResourceState::CleanupRequired
) {
let error = self.invalid_state(operation, action);
return Box::pin(async move { Err(error) });
}
Box::pin(async move {
let previous_lifecycle = self.lifecycle.clone();
self.lifecycle.begin_pending();
let result = call(self.session.as_mut()).await;
self.lifecycle = previous_lifecycle;
match &result {
Ok(()) => self.lifecycle.record_cleanup_success(),
Err(error) => self.lifecycle.record_cleanup_error(error),
}
result.map_err(|error| {
error.with_operation(operation).with_missing_context(
&self.path,
None,
self.file_system.properties().info().provider_id(),
)
})
})
}
fn invalid_state(&self, operation: FsOperation, action: &str) -> FsError {
let message = format!("{} {}", self.resource_name, action);
FsError::new(FsErrorKind::InvalidState, operation, &message).with_path(self.path.clone())
}
fn contextual_persist_error(&self, error: FsError, target: &Path) -> FsError {
error.with_operation(FsOperation::PersistTemp).with_missing_context(
&self.path,
Some(target),
self.file_system.properties().info().provider_id(),
)
}
}
impl Drop for AsyncTempFile {
fn drop(&mut self) {
if matches!(
self.lifecycle.state(),
TempResourceState::Owned | TempResourceState::CleanupRequired
) {
self.session.as_mut().cancel_on_drop();
}
}
}