use std::fmt::Debug;
use std::fmt::Formatter;
use std::fmt::Result as FmtResult;
use std::io::Error as IoError;
use qubit_io::AsyncOutput;
use crate::AsyncFileSystem;
use crate::error::FsError;
use crate::error::FsErrorKind;
use crate::error::FsOperation;
use crate::error::OpenFailureStage;
use crate::metadata::WriteOutcome;
use crate::path::Path;
use crate::write::AsyncWriteAllOperationFailure;
use crate::write::AsyncWriteAllOperationState;
use crate::write::AsyncWriterRecovery;
use crate::write::WriteFailureState;
use crate::write::WriteOptions;
use crate::write::WriterState;
use crate::write::internal::WriteAllCancellationGuard;
use crate::write::internal::WriteAllRecoverySnapshot;
use crate::write::internal::open_failure_state;
#[must_use]
pub struct AsyncWriteAllOperation {
filesystem: AsyncFileSystem,
path: Path,
bytes: Vec<u8>,
options: WriteOptions,
state: AsyncWriteAllOperationState,
writer: Option<AsyncWriterRecovery>,
recovery: WriteAllRecoverySnapshot,
}
impl AsyncWriteAllOperation {
pub(crate) fn new(filesystem: AsyncFileSystem, path: Path, bytes: Vec<u8>, options: WriteOptions) -> Self {
Self {
filesystem,
path,
bytes,
options,
state: AsyncWriteAllOperationState::Ready,
writer: None,
recovery: WriteAllRecoverySnapshot::new(),
}
}
#[inline]
#[must_use]
pub const fn filesystem(&self) -> &AsyncFileSystem {
&self.filesystem
}
#[inline]
#[must_use]
pub const fn path(&self) -> &Path {
&self.path
}
#[inline]
#[must_use = "inspect publication state before choosing a recovery action"]
pub const fn state(&self) -> AsyncWriteAllOperationState {
self.state
}
#[inline]
#[must_use]
pub const fn has_recovery(&self) -> bool {
self.writer.is_some()
}
#[inline]
#[must_use]
pub fn recovery(&mut self) -> Option<&mut AsyncWriterRecovery> {
self.writer.as_mut()
}
#[inline]
#[must_use]
pub fn take_recovery(&mut self) -> Option<AsyncWriterRecovery> {
self.writer.take()
}
#[inline]
#[must_use]
pub const fn written_bytes(&self) -> u64 {
self.recovery.written_bytes
}
pub async fn execute(&mut self) -> Result<WriteOutcome, AsyncWriteAllOperationFailure> {
if self.state != AsyncWriteAllOperationState::Ready {
return Err(AsyncWriteAllOperationFailure::new(
invalid_state(&self.path, &self.filesystem),
self.recovery.state,
self.recovery.written_bytes,
));
}
let Self {
filesystem,
path,
bytes,
options,
state,
writer,
recovery,
} = self;
let bytes = std::mem::take(bytes);
let mut guard = WriteAllCancellationGuard::start(state, writer, recovery);
let result = execute_write(filesystem, path, &bytes, options, guard.writer_mut()).await;
guard.finish(&result);
result
}
}
async fn execute_write(
filesystem: &AsyncFileSystem,
path: &Path,
bytes: &[u8],
options: &WriteOptions,
slot: &mut Option<AsyncWriterRecovery>,
) -> Result<WriteOutcome, AsyncWriteAllOperationFailure> {
if slot.is_none() {
match filesystem.open_writer(path, options.clone()).await {
Ok(writer) => *slot = Some(AsyncWriterRecovery::Opened(Box::new(writer))),
Err(failure) => {
let (error, stage, recovery) = failure.into_parts();
*slot = recovery.map(AsyncWriterRecovery::Rejected);
let state = match stage {
OpenFailureStage::Preflight => WriteFailureState::NotPublished,
OpenFailureStage::ProviderOpen => open_failure_state(&error),
OpenFailureStage::OutcomeValidation => WriteFailureState::Indeterminate,
};
return Err(AsyncWriteAllOperationFailure::new(error, state, 0));
}
}
}
let writer = slot
.as_mut()
.and_then(AsyncWriterRecovery::opened_mut)
.expect("writer is retained after open");
if let Err(error) = writer.write_fully_async(bytes).await {
let error = contextual(filesystem, error, path);
let state = state_for(error.has_indeterminate_effect(), writer.state());
return Err(AsyncWriteAllOperationFailure::new(error, state, writer.written_bytes()));
}
if let Err(error) = writer.flush_async().await {
let error = contextual(filesystem, error, path);
let state = state_for(error.has_indeterminate_effect(), writer.state());
return Err(AsyncWriteAllOperationFailure::new(error, state, writer.written_bytes()));
}
match writer.commit_async().await {
Ok(outcome) => Ok(outcome),
Err(failure) => {
let (error, state) = failure.into_parts();
Err(AsyncWriteAllOperationFailure::new(error, state, writer.written_bytes()))
}
}
}
fn contextual(filesystem: &AsyncFileSystem, error: IoError, path: &Path) -> FsError {
filesystem.core().enrich(
FsError::from_stream_io(error, FsOperation::Write, path),
Some(path),
FsOperation::Write,
)
}
fn state_for(indeterminate: bool, state: WriterState) -> WriteFailureState {
if indeterminate {
WriteFailureState::Indeterminate
} else {
state.publication_failure_state()
}
}
fn invalid_state(path: &Path, filesystem: &AsyncFileSystem) -> FsError {
FsError::new(
FsErrorKind::InvalidState,
FsOperation::Write,
"async whole-file write cannot execute in its current state",
)
.with_path(path.clone())
.with_provider(filesystem.properties().info().provider_id())
}
impl Debug for AsyncWriteAllOperation {
fn fmt(&self, formatter: &mut Formatter<'_>) -> FmtResult {
formatter
.debug_struct("AsyncWriteAllOperation")
.field("path", &self.path)
.field("state", &self.state)
.field("written_bytes", &self.recovery.written_bytes)
.field("has_recovery", &self.writer.is_some())
.finish()
}
}