Skip to main content

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