Skip to main content

qubit_fs/copy/
async_copy_operation.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//! Owning asynchronous copy operation with cancellation-safe state tracking.
9
10use qubit_io::AsyncInput;
11use qubit_io::AsyncOutput;
12
13use super::fallback_failure_stats;
14use super::from_writer_state;
15use super::internal::CopyCancellationGuard;
16use super::internal::CopyDeadline;
17use super::internal::CopyRecoverySnapshot;
18use super::internal::StreamCopyPlan;
19use super::internal::from_completed_stats;
20use crate::AsyncFileSystem;
21use crate::copy::AsyncCopyFailure;
22use crate::copy::AsyncCopyOperationState;
23use crate::copy::CopyConflictPolicy;
24use crate::copy::CopyFailureState;
25use crate::copy::CopyOptions;
26use crate::copy::CopyOutcome;
27use crate::copy::CopyStats;
28use crate::copy::internal::from_write_failure_state;
29use crate::error::FsError;
30use crate::error::FsErrorKind;
31use crate::error::FsOperation;
32use crate::error::OpenFailureStage;
33use crate::metadata::FileSystemCapability;
34use crate::metadata::SymlinkPolicy;
35use crate::path::Path;
36use crate::read::ReadOptions;
37use crate::spi::CopyAttempt;
38use crate::spi::CopyRequest;
39use crate::spi::ProviderOperation;
40use crate::spi::ResolvedCopyOptions;
41use crate::spi::SpiFuture;
42use crate::write::AsyncWriterRecovery;
43use crate::write::internal::is_unchanged_open_failure;
44use crate::write::internal::open_failure_state;
45
46/// An owning copy request whose recovery writer remains accessible after
47/// failure.
48///
49/// # Examples
50///
51/// This example uses an isolated in-memory provider fixture.
52///
53/// ```rust
54/// # mod support { include!(concat!(env!("CARGO_MANIFEST_DIR"), "/tests/common/rustdoc_support.rs")); }
55/// # use support::*;
56/// # let (filesystem, _) = async_recording_spi::async_recording_file_system(Default::default());
57/// # poll_support::ready(async {
58/// use qubit_fs::Path;
59/// use qubit_fs::copy::CopyOptions;
60///
61/// let mut operation = filesystem.begin_copy(
62///     Path::parse("/source")?, Path::parse("/target")?, CopyOptions::default(),
63/// )?;
64/// let outcome = operation.execute().await?;
65/// assert_eq!(5, outcome.stats().bytes);
66/// # Ok::<(), Box<dyn std::error::Error>>(())
67/// # }).unwrap();
68/// ```
69pub struct AsyncCopyOperation {
70    /// Facade used to validate and execute the pending operation.
71    pub(crate) file_system: AsyncFileSystem,
72    /// Immutable validated source path.
73    source: Path,
74    /// Immutable validated destination path.
75    target: Path,
76    /// Facade-resolved copy policy.
77    options: ResolvedCopyOptions,
78    /// Observable lifecycle state.
79    state: AsyncCopyOperationState,
80    /// Destination writer retained when recovery remains possible.
81    writer: Option<AsyncWriterRecovery>,
82    /// Monotonic start used to enforce caller elapsed-time budgets.
83    deadline: CopyDeadline,
84    /// Publication state and progress retained for cancellation recovery.
85    recovery: CopyRecoverySnapshot,
86}
87
88impl AsyncCopyOperation {
89    /// Builds a preflight-validated operation without provider I/O.
90    pub(crate) fn new(
91        file_system: AsyncFileSystem,
92        source: Path,
93        target: Path,
94        options: CopyOptions,
95        symlink_policy: SymlinkPolicy,
96    ) -> Self {
97        Self {
98            file_system,
99            source,
100            target,
101            deadline: CopyDeadline::new(options.deadline()),
102            options: ResolvedCopyOptions::new(options, symlink_policy),
103            state: AsyncCopyOperationState::Ready,
104            writer: None,
105            recovery: CopyRecoverySnapshot::unchanged(),
106        }
107    }
108
109    /// Returns the immutable source path.
110    ///
111    /// # Returns
112    /// The source path captured when the operation was created.
113    #[inline]
114    #[must_use]
115    pub const fn source(&self) -> &Path {
116        &self.source
117    }
118
119    /// Returns the immutable destination path.
120    #[inline]
121    #[must_use]
122    pub const fn target(&self) -> &Path {
123        &self.target
124    }
125
126    /// Returns the current operation lifecycle state.
127    #[inline]
128    #[must_use]
129    pub const fn state(&self) -> AsyncCopyOperationState {
130        self.state
131    }
132
133    /// Returns whether a recovery writer is retained by this operation.
134    #[inline]
135    #[must_use]
136    pub const fn has_recovery(&self) -> bool {
137        self.writer.is_some()
138    }
139
140    /// Borrows the retained recovery writer, if one exists.
141    #[inline]
142    pub fn recovery(&mut self) -> Option<&mut AsyncWriterRecovery> {
143        self.writer.as_mut()
144    }
145
146    /// Takes ownership of the retained recovery writer, if one exists.
147    #[inline]
148    pub fn take_recovery(&mut self) -> Option<AsyncWriterRecovery> {
149        self.writer.take()
150    }
151
152    /// Executes the operation exactly once.
153    ///
154    /// The operation becomes running only when this future is first polled.
155    /// Dropping a polled pending future records indeterminate state and does
156    /// not invoke provider cleanup or start any additional I/O. The deadline
157    /// includes time spent waiting before execution and is cooperative: a
158    /// pending provider future is not forcibly woken by this operation, and
159    /// external cancellation does not roll back provider publication.
160    ///
161    /// # Returns
162    /// The provider-confirmed copy outcome.
163    ///
164    /// # Errors
165    /// Returns [`AsyncCopyFailure`] with the confirmed copy state and any
166    /// retained writer when preflight, provider execution, fallback, or
167    /// outcome validation fails.
168    pub async fn execute(&mut self) -> Result<CopyOutcome, AsyncCopyFailure> {
169        if self.state != AsyncCopyOperationState::Ready {
170            return Err(invalid_state_failure(
171                &self.source,
172                &self.target,
173                self.file_system.properties().info().provider_id(),
174                self.recovery,
175            ));
176        }
177        let Self {
178            file_system,
179            source,
180            target,
181            options,
182            state,
183            writer,
184            deadline,
185            recovery,
186        } = self;
187        let mut guard = CopyCancellationGuard::start(state, writer, recovery);
188        let result = execute_copy(file_system, source, target, options, *deadline, guard.writer_mut()).await;
189        guard.finish(&result);
190        result
191    }
192}
193
194/// Dispatches the asynchronous provider attempt and falls back only when the
195/// provider explicitly declines it.
196async fn execute_copy(
197    filesystem: &AsyncFileSystem,
198    source: &Path,
199    target: &Path,
200    options: &ResolvedCopyOptions,
201    deadline: CopyDeadline,
202    writer: &mut Option<AsyncWriterRecovery>,
203) -> Result<CopyOutcome, AsyncCopyFailure> {
204    let caller_options = options.options();
205    if caller_options.max_entries() == Some(0) {
206        return Err(filesystem.contextual_copy_failure(
207            budget_error(source, target, "copy entry limit was exceeded"),
208            CopyFailureState::Unchanged,
209            CopyStats::default(),
210            source,
211            target,
212        ));
213    }
214    if deadline.expired() {
215        return Err(filesystem.contextual_copy_failure(
216            budget_error(source, target, "copy deadline was exceeded"),
217            CopyFailureState::Unchanged,
218            CopyStats::default(),
219            source,
220            target,
221        ));
222    }
223    if !filesystem.core().provider_supports(ProviderOperation::TryCopy) {
224        return stream_copy_fallback(filesystem, source, target, options, deadline, writer).await;
225    }
226    match filesystem
227        .spi()
228        .try_copy(CopyRequest::new(source, target, options.clone()))
229        .await
230    {
231        Ok(CopyAttempt::Completed(outcome)) => {
232            let outcome = filesystem.verify_completed_copy(outcome, options.options(), source, target)?;
233            if deadline.expired() {
234                return Err(filesystem.contextual_copy_failure(
235                    budget_error(source, target, "copy deadline was exceeded"),
236                    from_completed_stats(outcome.stats()),
237                    *outcome.stats(),
238                    source,
239                    target,
240                ));
241            }
242            Ok(outcome)
243        }
244        Ok(CopyAttempt::Declined(_)) => {
245            stream_copy_fallback(filesystem, source, target, options, deadline, writer).await
246        }
247        Err(failure) => {
248            let (error, state, stats) = failure.into_parts();
249            Err(filesystem.contextual_copy_failure(error, state, stats, source, target))
250        }
251    }
252}
253
254/// Streams a declined asynchronous copy while retaining any recovery writer
255/// in `writer_slot` until publication or cleanup completes.
256#[inline]
257fn stream_copy_fallback<'a>(
258    filesystem: &'a AsyncFileSystem,
259    source: &'a Path,
260    target: &'a Path,
261    options: &'a ResolvedCopyOptions,
262    deadline: CopyDeadline,
263    writer_slot: &'a mut Option<AsyncWriterRecovery>,
264) -> SpiFuture<'a, Result<CopyOutcome, AsyncCopyFailure>> {
265    Box::pin(async move {
266        let options = options.options();
267        let plan = StreamCopyPlan::new(options, filesystem.properties().limits(), source, target);
268        if let Err(error) = plan.validate_options(filesystem.properties().symlink_policy()) {
269            return Err(filesystem.contextual_copy_failure(
270                error,
271                CopyFailureState::Unchanged,
272                CopyStats::default(),
273                source,
274                target,
275            ));
276        }
277        if deadline.expired() {
278            return Err(filesystem.contextual_copy_failure(
279                budget_error(source, target, "copy deadline was exceeded"),
280                CopyFailureState::Unchanged,
281                CopyStats::default(),
282                source,
283                target,
284            ));
285        }
286        filesystem
287            .require(FileSystemCapability::Read, FsOperation::Copy, source)
288            .and_then(|_| filesystem.require(FileSystemCapability::Write, FsOperation::Copy, target))
289            .map_err(|error| {
290                filesystem.contextual_copy_failure(
291                    error,
292                    CopyFailureState::Unchanged,
293                    CopyStats::default(),
294                    source,
295                    target,
296                )
297            })?;
298        let metadata = filesystem.stat(source).await.map_err(|error| {
299            filesystem.contextual_copy_failure(error, CopyFailureState::Unchanged, CopyStats::default(), source, target)
300        })?;
301        if deadline.expired() {
302            return Err(filesystem.contextual_copy_failure(
303                budget_error(source, target, "copy deadline was exceeded"),
304                CopyFailureState::Unchanged,
305                CopyStats::default(),
306                source,
307                target,
308            ));
309        }
310        if let Err(error) = plan.validate_metadata(&metadata) {
311            return Err(filesystem.contextual_copy_failure(
312                error,
313                CopyFailureState::Unchanged,
314                CopyStats::default(),
315                source,
316                target,
317            ));
318        }
319        let mut reader = filesystem
320            .open_reader(source, ReadOptions::default())
321            .await
322            .map_err(|error| {
323                filesystem.contextual_copy_failure(
324                    error,
325                    CopyFailureState::Unchanged,
326                    CopyStats::default(),
327                    source,
328                    target,
329                )
330            })?;
331        if deadline.expired() {
332            return Err(filesystem.contextual_copy_failure(
333                budget_error(source, target, "copy deadline was exceeded"),
334                CopyFailureState::Unchanged,
335                CopyStats::default(),
336                source,
337                target,
338            ));
339        }
340        let writer_options = plan.writer_options();
341        match filesystem.open_writer(target, writer_options).await {
342            Ok(writer) => *writer_slot = Some(AsyncWriterRecovery::Opened(Box::new(writer))),
343            Err(error)
344                if error.stage() == OpenFailureStage::ProviderOpen
345                    && error.recovery().is_none()
346                    && error.error().kind() == FsErrorKind::AlreadyExists
347                    && is_unchanged_open_failure(error.error())
348                    && options.conflict() == CopyConflictPolicy::Skip =>
349            {
350                return Ok(CopyOutcome::streamed_fallback(
351                    CopyStats {
352                        skipped: 1,
353                        ..CopyStats::default()
354                    },
355                    crate::metadata::AchievedAtomicity::NonAtomic,
356                    false,
357                ));
358            }
359            Err(error) => {
360                let (error, stage, recovery) = error.into_parts();
361                *writer_slot = recovery.map(AsyncWriterRecovery::Rejected);
362                let state = match stage {
363                    OpenFailureStage::Preflight => CopyFailureState::Unchanged,
364                    OpenFailureStage::ProviderOpen => from_write_failure_state(open_failure_state(&error)),
365                    OpenFailureStage::OutcomeValidation => CopyFailureState::Indeterminate,
366                };
367                return Err(filesystem.contextual_copy_failure(error, state, CopyStats::default(), source, target));
368            }
369        }
370        if deadline.expired() {
371            let writer = writer_slot
372                .as_ref()
373                .and_then(AsyncWriterRecovery::opened)
374                .expect("writer is retained before transfer");
375            return Err(filesystem.contextual_copy_failure(
376                budget_error(source, target, "copy deadline was exceeded"),
377                from_writer_state(writer.state()),
378                fallback_failure_stats(writer.written_bytes()),
379                source,
380                target,
381            ));
382        }
383        let mut bytes = 0_u64;
384        let mut buffer = [0_u8; 8192];
385        loop {
386            if deadline.expired() {
387                let writer = writer_slot
388                    .as_ref()
389                    .and_then(AsyncWriterRecovery::opened)
390                    .expect("writer is retained before transfer");
391                return Err(filesystem.contextual_copy_failure(
392                    budget_error(source, target, "copy deadline was exceeded"),
393                    from_writer_state(writer.state()),
394                    fallback_failure_stats(writer.written_bytes()),
395                    source,
396                    target,
397                ));
398            }
399            let read = reader.read_async(&mut buffer).await.map_err(|error| {
400                filesystem.contextual_copy_failure(
401                    FsError::from_stream_io(error, FsOperation::Read, source),
402                    from_writer_state(
403                        writer_slot
404                            .as_ref()
405                            .and_then(AsyncWriterRecovery::opened)
406                            .expect("writer is retained before transfer")
407                            .state(),
408                    ),
409                    fallback_failure_stats(
410                        writer_slot
411                            .as_ref()
412                            .and_then(AsyncWriterRecovery::opened)
413                            .expect("writer is retained before transfer")
414                            .written_bytes(),
415                    ),
416                    source,
417                    target,
418                )
419            })?;
420            if deadline.expired() {
421                let writer = writer_slot
422                    .as_ref()
423                    .and_then(AsyncWriterRecovery::opened)
424                    .expect("writer is retained before transfer");
425                return Err(filesystem.contextual_copy_failure(
426                    budget_error(source, target, "copy deadline was exceeded"),
427                    from_writer_state(writer.state()),
428                    fallback_failure_stats(writer.written_bytes()),
429                    source,
430                    target,
431                ));
432            }
433            if read == 0 {
434                break;
435            }
436            let writer = writer_slot
437                .as_mut()
438                .and_then(AsyncWriterRecovery::opened_mut)
439                .expect("writer is retained before transfer");
440            let next_bytes = plan.next_bytes(bytes, read).map_err(|error| {
441                filesystem.contextual_copy_failure(
442                    error,
443                    from_writer_state(writer.state()),
444                    fallback_failure_stats(writer.written_bytes()),
445                    source,
446                    target,
447                )
448            })?;
449            writer.write_fully_async(&buffer[..read]).await.map_err(|error| {
450                filesystem.contextual_copy_failure(
451                    FsError::from_stream_io(error, FsOperation::Write, target),
452                    from_writer_state(writer.state()),
453                    fallback_failure_stats(writer.written_bytes()),
454                    source,
455                    target,
456                )
457            })?;
458            if deadline.expired() {
459                return Err(filesystem.contextual_copy_failure(
460                    budget_error(source, target, "copy deadline was exceeded"),
461                    from_writer_state(writer.state()),
462                    fallback_failure_stats(writer.written_bytes()),
463                    source,
464                    target,
465                ));
466            }
467            bytes = next_bytes;
468        }
469        let writer = writer_slot
470            .as_mut()
471            .and_then(AsyncWriterRecovery::opened_mut)
472            .expect("writer is retained before flush");
473        if deadline.expired() {
474            return Err(filesystem.contextual_copy_failure(
475                budget_error(source, target, "copy deadline was exceeded"),
476                from_writer_state(writer.state()),
477                fallback_failure_stats(writer.written_bytes()),
478                source,
479                target,
480            ));
481        }
482        writer.flush_async().await.map_err(|error| {
483            filesystem.contextual_copy_failure(
484                FsError::from_stream_io(error, FsOperation::Write, target),
485                from_writer_state(writer.state()),
486                fallback_failure_stats(writer.written_bytes()),
487                source,
488                target,
489            )
490        })?;
491        if deadline.expired() {
492            return Err(filesystem.contextual_copy_failure(
493                budget_error(source, target, "copy deadline was exceeded"),
494                from_writer_state(writer.state()),
495                fallback_failure_stats(writer.written_bytes()),
496                source,
497                target,
498            ));
499        }
500        let writer = writer_slot
501            .as_mut()
502            .and_then(AsyncWriterRecovery::opened_mut)
503            .expect("writer is retained before commit");
504        let write_outcome = match writer.commit_async().await {
505            Ok(outcome) => outcome,
506            Err(failure)
507                if failure.error().kind() == FsErrorKind::AlreadyExists
508                    && plan.may_skip_conflict(from_writer_state(writer.state())) =>
509            {
510                if let Err(cleanup_error) = writer.abort_async().await {
511                    return Err(filesystem.contextual_copy_failure(
512                        cleanup_error,
513                        from_writer_state(writer.state()),
514                        fallback_failure_stats(writer.written_bytes()),
515                        source,
516                        target,
517                    ));
518                }
519                let _ = writer_slot.take();
520                return Ok(CopyOutcome::streamed_fallback(
521                    CopyStats {
522                        skipped: 1,
523                        ..CopyStats::default()
524                    },
525                    crate::metadata::AchievedAtomicity::NonAtomic,
526                    false,
527                ));
528            }
529            Err(failure) => {
530                return Err(filesystem.contextual_copy_failure(
531                    failure.into_error(),
532                    from_writer_state(writer.state()),
533                    fallback_failure_stats(writer.written_bytes()),
534                    source,
535                    target,
536                ));
537            }
538        };
539        if deadline.expired() {
540            let _ = writer_slot.take();
541            return Err(filesystem.contextual_copy_failure(
542                budget_error(source, target, "copy deadline was exceeded"),
543                CopyFailureState::Published,
544                StreamCopyPlan::completed_stats(bytes),
545                source,
546                target,
547            ));
548        }
549        let _ = writer_slot.take();
550        Ok(CopyOutcome::streamed_fallback(
551            StreamCopyPlan::completed_stats(bytes),
552            write_outcome.atomicity(),
553            write_outcome.durable(),
554        ))
555    })
556}
557
558/// Builds a caller-budget error for an asynchronous copy.
559fn budget_error(source: &Path, target: &Path, message: &str) -> FsError {
560    FsError::new(FsErrorKind::ResourceLimitExceeded, FsOperation::Copy, message)
561        .with_path(source.clone())
562        .with_target(target.clone())
563}
564
565/// Builds the stable failure used for an invalid execute retry.
566fn invalid_state_failure(
567    source: &Path,
568    target: &Path,
569    provider: &str,
570    snapshot: CopyRecoverySnapshot,
571) -> AsyncCopyFailure {
572    AsyncCopyFailure::new(
573        FsError::new(
574            FsErrorKind::InvalidState,
575            FsOperation::Copy,
576            "copy operation cannot execute in its current state",
577        )
578        .with_path(source.clone())
579        .with_target(target.clone())
580        .with_provider(provider),
581        snapshot.state,
582        snapshot.stats,
583    )
584}