Skip to main content

qubit_fs/write/
file_writer.rs

1// =============================================================================
2//    Copyright (c) 2026 Haixing Hu.
3//
4//    SPDX-License-Identifier: Apache-2.0
5//
6//    Licensed under the Apache License, Version 2.0.
7// =============================================================================
8//! Concrete synchronous file writer handle.
9
10use std::fmt::Debug;
11use std::fmt::Formatter;
12use std::fmt::Result as FmtResult;
13use std::io::Error as IoError;
14use std::io::ErrorKind as IoErrorKind;
15use std::io::Result as IoResult;
16
17use qubit_io::Output;
18
19use crate::error::FsEffectState;
20use crate::error::FsError;
21use crate::error::FsErrorKind;
22use crate::error::FsOperation;
23use crate::error::FsResult;
24use crate::facade::facade_core::FacadeCore;
25use crate::facade::internal::ByteBudget;
26use crate::facade::internal::FileSystemResource;
27use crate::metadata::AchievedAtomicity;
28use crate::metadata::AtomicityRequirement;
29use crate::metadata::DurabilityRequirement;
30use crate::metadata::OpenedFileInfo;
31use crate::metadata::WriteOutcome;
32use crate::spi::FileWriterSpi;
33use crate::write::WriteAbortOutcome;
34use crate::write::WriteFailure;
35use crate::write::WriteFailureState;
36use crate::write::WriterState;
37
38/// Type-erased provider write session explicitly associated with a file.
39///
40/// # Examples
41///
42/// This example uses an isolated in-memory provider fixture.
43///
44/// ```rust
45/// # mod support { include!(concat!(env!("CARGO_MANIFEST_DIR"), "/tests/common/rustdoc_support.rs")); }
46/// # let filesystem = support::rustdoc_provider::filesystem();
47/// use qubit_fs::Path;
48/// use qubit_fs::write::WriteOptions;
49/// use qubit_fs::write::WriterState;
50/// use qubit_io::Output;
51///
52/// let mut writer = filesystem.open_writer(&Path::parse("/new-report")?, WriteOptions::default())?;
53/// writer.write_fully(b"report")?;
54/// let outcome = writer.commit()?;
55/// assert_eq!(Some(6), outcome.bytes_written());
56/// assert_eq!(WriterState::Committed, writer.state());
57/// # Ok::<(), Box<dyn std::error::Error>>(())
58/// ```
59pub struct FileWriter {
60    /// Provider write session.
61    session: Box<dyn FileWriterSpi>,
62    /// Stable identity and metadata captured at open time.
63    info: OpenedFileInfo,
64    /// Current publication lifecycle state.
65    state: WriterState,
66    /// Whether explicit provider cleanup has completed.
67    abort_completed: bool,
68    /// Atomicity required by the caller.
69    atomicity: AtomicityRequirement,
70    /// Durability required by the caller.
71    durability: DurabilityRequirement,
72    /// Provider identifier attached to facade-generated errors.
73    provider: Box<str>,
74    /// Optional inclusive byte limit for this write session.
75    write_budget: Option<ByteBudget>,
76    /// Bytes accepted by the provider session so far.
77    written_bytes: u64,
78}
79
80impl FileWriter {
81    /// Wraps an already-open provider write session.
82    ///
83    /// # Parameters
84    /// - `session`: Provider session accepting file bytes.
85    /// - `info`: File identity and optional open-time metadata snapshot.
86    ///
87    /// # Returns
88    /// A concrete writer in [`WriterState::Open`].
89    #[inline]
90    #[must_use]
91    pub(crate) fn new(
92        info: OpenedFileInfo,
93        session: Box<dyn FileWriterSpi>,
94        atomicity: AtomicityRequirement,
95        durability: DurabilityRequirement,
96        provider: &str,
97        max_write_bytes: Option<u64>,
98    ) -> Self {
99        Self {
100            session,
101            info,
102            state: WriterState::Open,
103            abort_completed: false,
104            atomicity,
105            durability,
106            provider: provider.into(),
107            write_budget: max_write_bytes
108                .map(|maximum| FacadeCore::byte_budget(FileSystemResource::WriteBytes, maximum)),
109            written_bytes: 0,
110        }
111    }
112
113    /// Returns the fixed identity and open-time metadata snapshot.
114    ///
115    /// # Returns
116    /// Information captured when the writer was opened.
117    #[inline]
118    #[must_use]
119    pub fn info(&self) -> &OpenedFileInfo {
120        &self.info
121    }
122
123    /// Returns the current lifecycle state.
124    ///
125    /// # Returns
126    /// Current writer state.
127    #[inline]
128    #[must_use]
129    pub const fn state(&self) -> WriterState {
130        self.state
131    }
132
133    /// Returns the bytes accepted by the underlying write session.
134    #[inline]
135    #[must_use]
136    pub(crate) const fn written_bytes(&self) -> u64 {
137        self.written_bytes
138    }
139
140    /// Publishes bytes accepted by this session.
141    ///
142    /// This method borrows rather than consumes the writer. The provider's
143    /// typed failure determines whether retry remains safe or only explicit
144    /// cleanup is available.
145    ///
146    /// # Returns
147    /// Actual publication method and atomicity on success.
148    ///
149    /// # Errors
150    /// Returns [`FsErrorKind::InvalidState`] after commit or abort, or the
151    /// provider publication error.
152    pub fn commit(&mut self) -> Result<WriteOutcome, WriteFailure> {
153        if self.state != WriterState::Open {
154            let publication_state = self.state.publication_failure_state();
155            return Err(WriteFailure::new(
156                self.invalid_state(
157                    FsOperation::CommitWriter,
158                    "writer cannot be committed in its current state",
159                ),
160                publication_state,
161            ));
162        }
163        let outcome = self.session.commit();
164        match outcome {
165            Ok(outcome) => {
166                self.state = WriterState::Committed;
167                if self.atomicity == AtomicityRequirement::Required && outcome.atomicity() != AchievedAtomicity::Atomic
168                {
169                    self.state = WriterState::Published;
170                    return Err(WriteFailure::new(
171                        FsError::new(
172                            FsErrorKind::ProviderContractViolation,
173                            FsOperation::CommitWriter,
174                            "provider reported non-atomic success for an atomic-required write",
175                        )
176                        .with_path(self.info.path().clone())
177                        .with_provider(&self.provider)
178                        .with_effect_state(FsEffectState::Applied),
179                        WriteFailureState::Published,
180                    ));
181                }
182                if self.durability == DurabilityRequirement::Required && !outcome.durable() {
183                    self.state = WriterState::Published;
184                    return Err(WriteFailure::new(
185                        FsError::new(
186                            FsErrorKind::ProviderContractViolation,
187                            FsOperation::CommitWriter,
188                            "provider reported non-durable success for a durability-required write",
189                        )
190                        .with_path(self.info.path().clone())
191                        .with_provider(&self.provider)
192                        .with_effect_state(FsEffectState::Applied),
193                        WriteFailureState::Published,
194                    ));
195                }
196                if let Some(bytes_written) = outcome.bytes_written()
197                    && bytes_written != self.written_bytes
198                {
199                    self.state = WriterState::Published;
200                    return Err(WriteFailure::new(
201                        FsError::new(
202                            FsErrorKind::ProviderContractViolation,
203                            FsOperation::CommitWriter,
204                            "provider reported a byte count different from the bytes accepted by the writer",
205                        )
206                        .with_path(self.info.path().clone())
207                        .with_provider(&self.provider)
208                        .with_effect_state(FsEffectState::Applied),
209                        WriteFailureState::Published,
210                    ));
211                }
212                Ok(outcome)
213            }
214            Err(failure) => {
215                self.state = match failure.state() {
216                    WriteFailureState::RetryableNotPublished => WriterState::Open,
217                    WriteFailureState::NotPublished => WriterState::NotPublished,
218                    WriteFailureState::Published => WriterState::Published,
219                    WriteFailureState::Indeterminate => WriterState::Indeterminate,
220                };
221                let (error, state) = failure.into_parts();
222                Err(WriteFailure::new(
223                    self.contextual_error(error, FsOperation::CommitWriter),
224                    state,
225                ))
226            }
227        }
228    }
229
230    /// Aborts this writer and releases provider staging resources.
231    ///
232    /// Abort is allowed while open and after every failed commit. Successful
233    /// cleanup does not imply that an already-published target was rolled back.
234    /// An indeterminate abort disables automatic drop cleanup so the caller can
235    /// inspect provider state before choosing a recovery action.
236    ///
237    /// # Returns
238    /// Provider-confirmed destination publication state after cleanup.
239    ///
240    /// # Errors
241    /// Returns [`FsErrorKind::InvalidState`] after commit or a previous abort,
242    /// or returns the provider cleanup error while retaining the session.
243    pub fn abort(&mut self) -> FsResult<WriteAbortOutcome> {
244        if self.abort_completed
245            || !matches!(
246                self.state,
247                WriterState::Open | WriterState::NotPublished | WriterState::Published | WriterState::Indeterminate
248            )
249        {
250            return Err(self.invalid_state(
251                FsOperation::AbortWriter,
252                "writer cannot be aborted in its current state",
253            ));
254        }
255        match self.session.abort() {
256            Ok(outcome) => {
257                self.abort_completed = true;
258                self.state = match outcome {
259                    WriteAbortOutcome::NotPublished => WriterState::Aborted,
260                    WriteAbortOutcome::Published => WriterState::Published,
261                    WriteAbortOutcome::Indeterminate => WriterState::Indeterminate,
262                };
263                Ok(outcome)
264            }
265            Err(error) => {
266                if error.has_indeterminate_effect() {
267                    self.state = WriterState::Indeterminate;
268                }
269                Err(self.contextual_error(error, FsOperation::AbortWriter))
270            }
271        }
272    }
273
274    /// Builds a stable invalid-state error for this writer.
275    fn invalid_state(&self, operation: FsOperation, message: &str) -> FsError {
276        FsError::new(FsErrorKind::InvalidState, operation, message)
277            .with_path(self.info.path().clone())
278            .with_provider(&self.provider)
279    }
280
281    /// Builds a stream error for byte transfer after lifecycle completion.
282    fn closed_io_error(&self) -> IoError {
283        IoError::new(
284            IoErrorKind::BrokenPipe,
285            self.invalid_state(FsOperation::Write, "writer no longer accepts bytes"),
286        )
287    }
288
289    /// Checks whether a provider write can fit in the session budget.
290    fn check_write_limit(&self, count: usize) -> IoResult<u64> {
291        let count = FacadeCore::quantity_from_usize(count, FsOperation::Write, self.info.path(), &self.provider)
292            .map_err(FsError::into_io_error)?;
293        if let Some(budget) = &self.write_budget
294            && let Err(error) = budget.check_available(count)
295        {
296            return Err(FacadeCore::budget_error(
297                error,
298                FsOperation::Write,
299                self.info.path(),
300                &self.provider,
301                "write session exceeds the provider byte limit",
302            )
303            .into_io_error());
304        }
305        Ok(count)
306    }
307
308    /// Records bytes accepted by the provider in the public `u64` accounting
309    /// domain.
310    ///
311    /// Returns an I/O error when the native byte count or accumulated total
312    /// cannot be represented by the filesystem API's `u64` byte counters.
313    fn record_written_bytes(&mut self, count: usize) -> IoResult<()> {
314        let count = FacadeCore::quantity_from_usize(count, FsOperation::Write, self.info.path(), &self.provider)
315            .map_err(FsError::into_io_error)?;
316        if let Some(error) = self
317            .write_budget
318            .as_mut()
319            .and_then(|budget| budget.try_consume(count).err())
320        {
321            return Err(FacadeCore::budget_error(
322                error,
323                FsOperation::Write,
324                self.info.path(),
325                &self.provider,
326                "write session exceeds the provider byte limit",
327            )
328            .into_io_error());
329        }
330        self.written_bytes = self
331            .written_bytes
332            .checked_add(count)
333            .ok_or_else(|| self.byte_count_error())?;
334        Ok(())
335    }
336
337    /// Builds the error used when native byte accounting exceeds the public
338    /// filesystem API's `u64` reporting range.
339    fn byte_count_error(&self) -> IoError {
340        FsError::new(
341            FsErrorKind::ResourceLimitExceeded,
342            FsOperation::Write,
343            "write byte count exceeds the filesystem API reporting range",
344        )
345        .with_path(self.info.path().clone())
346        .with_provider(&self.provider)
347        .into_io_error()
348    }
349
350    /// Adds only missing facade context to a provider lifecycle error.
351    fn contextual_error(&self, error: FsError, operation: FsOperation) -> FsError {
352        error
353            .with_operation(operation)
354            .with_missing_context(self.info.path(), None, &self.provider)
355    }
356}
357
358impl Output for FileWriter {
359    type Item = u8;
360
361    #[inline]
362    fn is_buffered(&self) -> bool {
363        self.session.is_buffered()
364    }
365
366    unsafe fn write_unchecked(&mut self, input: &[u8], index: usize, count: usize) -> IoResult<usize> {
367        if self.state != WriterState::Open {
368            return Err(self.closed_io_error());
369        }
370        self.check_write_limit(count)?;
371        // SAFETY: The caller guarantees the same range contract required by
372        // the wrapped output session.
373        match unsafe { self.session.write_unchecked(input, index, count) } {
374            Ok(value) => {
375                if let Err(error) = self.record_written_bytes(value) {
376                    self.state = WriterState::Indeterminate;
377                    return Err(error);
378                }
379                Ok(value)
380            }
381            Err(error) => {
382                self.state = WriterState::Indeterminate;
383                Err(error)
384            }
385        }
386    }
387
388    fn flush(&mut self) -> IoResult<()> {
389        if self.state != WriterState::Open {
390            return Err(self.closed_io_error());
391        }
392        match self.session.flush() {
393            Ok(()) => Ok(()),
394            Err(error) => {
395                self.state = WriterState::Indeterminate;
396                Err(error)
397            }
398        }
399    }
400}
401
402impl Debug for FileWriter {
403    #[inline]
404    fn fmt(&self, formatter: &mut Formatter<'_>) -> FmtResult {
405        formatter
406            .debug_struct("FileWriter")
407            .field("info", &self.info)
408            .field("state", &self.state)
409            .finish_non_exhaustive()
410    }
411}
412
413impl Drop for FileWriter {
414    fn drop(&mut self) {
415        if !self.abort_completed
416            && matches!(
417                self.state,
418                WriterState::Open | WriterState::NotPublished | WriterState::Published
419            )
420        {
421            let _ = self.session.abort();
422        }
423    }
424}