use std::fmt::Debug;
use std::fmt::Formatter;
use std::fmt::Result as FmtResult;
use std::io::Error as IoError;
use std::io::ErrorKind as IoErrorKind;
use std::io::Result as IoResult;
use qubit_io::Output;
use crate::error::FsEffectState;
use crate::error::FsError;
use crate::error::FsErrorKind;
use crate::error::FsOperation;
use crate::error::FsResult;
use crate::facade::facade_core::FacadeCore;
use crate::facade::internal::ByteBudget;
use crate::facade::internal::FileSystemResource;
use crate::metadata::AchievedAtomicity;
use crate::metadata::AtomicityRequirement;
use crate::metadata::DurabilityRequirement;
use crate::metadata::OpenedFileInfo;
use crate::metadata::WriteOutcome;
use crate::spi::FileWriterSpi;
use crate::write::WriteAbortOutcome;
use crate::write::WriteFailure;
use crate::write::WriteFailureState;
use crate::write::WriterState;
pub struct FileWriter {
session: Box<dyn FileWriterSpi>,
info: OpenedFileInfo,
state: WriterState,
abort_completed: bool,
atomicity: AtomicityRequirement,
durability: DurabilityRequirement,
provider: Box<str>,
write_budget: Option<ByteBudget>,
written_bytes: u64,
}
impl FileWriter {
#[inline]
#[must_use]
pub(crate) fn new(
info: OpenedFileInfo,
session: Box<dyn FileWriterSpi>,
atomicity: AtomicityRequirement,
durability: DurabilityRequirement,
provider: &str,
max_write_bytes: Option<u64>,
) -> Self {
Self {
session,
info,
state: WriterState::Open,
abort_completed: false,
atomicity,
durability,
provider: provider.into(),
write_budget: max_write_bytes
.map(|maximum| FacadeCore::byte_budget(FileSystemResource::WriteBytes, maximum)),
written_bytes: 0,
}
}
#[inline]
#[must_use]
pub fn info(&self) -> &OpenedFileInfo {
&self.info
}
#[inline]
#[must_use]
pub const fn state(&self) -> WriterState {
self.state
}
#[inline]
#[must_use]
pub(crate) const fn written_bytes(&self) -> u64 {
self.written_bytes
}
pub fn commit(&mut self) -> Result<WriteOutcome, WriteFailure> {
if self.state != WriterState::Open {
let publication_state = self.state.publication_failure_state();
return Err(WriteFailure::new(
self.invalid_state(
FsOperation::CommitWriter,
"writer cannot be committed in its current state",
),
publication_state,
));
}
let outcome = self.session.commit();
match outcome {
Ok(outcome) => {
self.state = WriterState::Committed;
if self.atomicity == AtomicityRequirement::Required && outcome.atomicity() != AchievedAtomicity::Atomic
{
self.state = WriterState::Published;
return Err(WriteFailure::new(
FsError::new(
FsErrorKind::ProviderContractViolation,
FsOperation::CommitWriter,
"provider reported non-atomic success for an atomic-required write",
)
.with_path(self.info.path().clone())
.with_provider(&self.provider)
.with_effect_state(FsEffectState::Applied),
WriteFailureState::Published,
));
}
if self.durability == DurabilityRequirement::Required && !outcome.durable() {
self.state = WriterState::Published;
return Err(WriteFailure::new(
FsError::new(
FsErrorKind::ProviderContractViolation,
FsOperation::CommitWriter,
"provider reported non-durable success for a durability-required write",
)
.with_path(self.info.path().clone())
.with_provider(&self.provider)
.with_effect_state(FsEffectState::Applied),
WriteFailureState::Published,
));
}
if let Some(bytes_written) = outcome.bytes_written()
&& bytes_written != self.written_bytes
{
self.state = WriterState::Published;
return Err(WriteFailure::new(
FsError::new(
FsErrorKind::ProviderContractViolation,
FsOperation::CommitWriter,
"provider reported a byte count different from the bytes accepted by the writer",
)
.with_path(self.info.path().clone())
.with_provider(&self.provider)
.with_effect_state(FsEffectState::Applied),
WriteFailureState::Published,
));
}
Ok(outcome)
}
Err(failure) => {
self.state = match failure.state() {
WriteFailureState::RetryableNotPublished => WriterState::Open,
WriteFailureState::NotPublished => WriterState::NotPublished,
WriteFailureState::Published => WriterState::Published,
WriteFailureState::Indeterminate => WriterState::Indeterminate,
};
let (error, state) = failure.into_parts();
Err(WriteFailure::new(
self.contextual_error(error, FsOperation::CommitWriter),
state,
))
}
}
}
pub fn abort(&mut self) -> FsResult<WriteAbortOutcome> {
if self.abort_completed
|| !matches!(
self.state,
WriterState::Open | WriterState::NotPublished | WriterState::Published | WriterState::Indeterminate
)
{
return Err(self.invalid_state(
FsOperation::AbortWriter,
"writer cannot be aborted in its current state",
));
}
match self.session.abort() {
Ok(outcome) => {
self.abort_completed = true;
self.state = match outcome {
WriteAbortOutcome::NotPublished => WriterState::Aborted,
WriteAbortOutcome::Published => WriterState::Published,
WriteAbortOutcome::Indeterminate => WriterState::Indeterminate,
};
Ok(outcome)
}
Err(error) => {
if error.has_indeterminate_effect() {
self.state = WriterState::Indeterminate;
}
Err(self.contextual_error(error, FsOperation::AbortWriter))
}
}
}
fn invalid_state(&self, operation: FsOperation, message: &str) -> FsError {
FsError::new(FsErrorKind::InvalidState, operation, message)
.with_path(self.info.path().clone())
.with_provider(&self.provider)
}
fn closed_io_error(&self) -> IoError {
IoError::new(
IoErrorKind::BrokenPipe,
self.invalid_state(FsOperation::Write, "writer no longer accepts bytes"),
)
}
fn check_write_limit(&self, count: usize) -> IoResult<u64> {
let count = FacadeCore::quantity_from_usize(count, FsOperation::Write, self.info.path(), &self.provider)
.map_err(FsError::into_io_error)?;
if let Some(budget) = &self.write_budget
&& let Err(error) = budget.check_available(count)
{
return Err(FacadeCore::budget_error(
error,
FsOperation::Write,
self.info.path(),
&self.provider,
"write session exceeds the provider byte limit",
)
.into_io_error());
}
Ok(count)
}
fn record_written_bytes(&mut self, count: usize) -> IoResult<()> {
let count = FacadeCore::quantity_from_usize(count, FsOperation::Write, self.info.path(), &self.provider)
.map_err(FsError::into_io_error)?;
if let Some(error) = self
.write_budget
.as_mut()
.and_then(|budget| budget.try_consume(count).err())
{
return Err(FacadeCore::budget_error(
error,
FsOperation::Write,
self.info.path(),
&self.provider,
"write session exceeds the provider byte limit",
)
.into_io_error());
}
self.written_bytes = self
.written_bytes
.checked_add(count)
.ok_or_else(|| self.byte_count_error())?;
Ok(())
}
fn byte_count_error(&self) -> IoError {
FsError::new(
FsErrorKind::ResourceLimitExceeded,
FsOperation::Write,
"write byte count exceeds the filesystem API reporting range",
)
.with_path(self.info.path().clone())
.with_provider(&self.provider)
.into_io_error()
}
fn contextual_error(&self, error: FsError, operation: FsOperation) -> FsError {
error
.with_operation(operation)
.with_missing_context(self.info.path(), None, &self.provider)
}
}
impl Output for FileWriter {
type Item = u8;
#[inline]
fn is_buffered(&self) -> bool {
self.session.is_buffered()
}
unsafe fn write_unchecked(&mut self, input: &[u8], index: usize, count: usize) -> IoResult<usize> {
if self.state != WriterState::Open {
return Err(self.closed_io_error());
}
self.check_write_limit(count)?;
match unsafe { self.session.write_unchecked(input, index, count) } {
Ok(value) => {
if let Err(error) = self.record_written_bytes(value) {
self.state = WriterState::Indeterminate;
return Err(error);
}
Ok(value)
}
Err(error) => {
self.state = WriterState::Indeterminate;
Err(error)
}
}
}
fn flush(&mut self) -> IoResult<()> {
if self.state != WriterState::Open {
return Err(self.closed_io_error());
}
match self.session.flush() {
Ok(()) => Ok(()),
Err(error) => {
self.state = WriterState::Indeterminate;
Err(error)
}
}
}
}
impl Debug for FileWriter {
#[inline]
fn fmt(&self, formatter: &mut Formatter<'_>) -> FmtResult {
formatter
.debug_struct("FileWriter")
.field("info", &self.info)
.field("state", &self.state)
.finish_non_exhaustive()
}
}
impl Drop for FileWriter {
fn drop(&mut self) {
if !self.abort_completed
&& matches!(
self.state,
WriterState::Open | WriterState::NotPublished | WriterState::Published
)
{
let _ = self.session.abort();
}
}
}