use qubit_io::AsyncInput;
use qubit_io::AsyncOutput;
use super::fallback_failure_stats;
use super::from_writer_state;
use super::internal::CopyCancellationGuard;
use super::internal::CopyDeadline;
use super::internal::CopyRecoverySnapshot;
use super::internal::StreamCopyPlan;
use super::internal::from_completed_stats;
use crate::AsyncFileSystem;
use crate::copy::AsyncCopyFailure;
use crate::copy::AsyncCopyOperationState;
use crate::copy::CopyConflictPolicy;
use crate::copy::CopyFailureState;
use crate::copy::CopyOptions;
use crate::copy::CopyOutcome;
use crate::copy::CopyStats;
use crate::copy::internal::from_write_failure_state;
use crate::error::FsError;
use crate::error::FsErrorKind;
use crate::error::FsOperation;
use crate::error::OpenFailureStage;
use crate::metadata::FileSystemCapability;
use crate::metadata::SymlinkPolicy;
use crate::path::Path;
use crate::read::ReadOptions;
use crate::spi::CopyAttempt;
use crate::spi::CopyRequest;
use crate::spi::ProviderOperation;
use crate::spi::ResolvedCopyOptions;
use crate::spi::SpiFuture;
use crate::write::AsyncWriterRecovery;
use crate::write::internal::is_unchanged_open_failure;
use crate::write::internal::open_failure_state;
pub struct AsyncCopyOperation {
pub(crate) file_system: AsyncFileSystem,
source: Path,
target: Path,
options: ResolvedCopyOptions,
state: AsyncCopyOperationState,
writer: Option<AsyncWriterRecovery>,
deadline: CopyDeadline,
recovery: CopyRecoverySnapshot,
}
impl AsyncCopyOperation {
pub(crate) fn new(
file_system: AsyncFileSystem,
source: Path,
target: Path,
options: CopyOptions,
symlink_policy: SymlinkPolicy,
) -> Self {
Self {
file_system,
source,
target,
deadline: CopyDeadline::new(options.deadline()),
options: ResolvedCopyOptions::new(options, symlink_policy),
state: AsyncCopyOperationState::Ready,
writer: None,
recovery: CopyRecoverySnapshot::unchanged(),
}
}
#[inline]
#[must_use]
pub const fn source(&self) -> &Path {
&self.source
}
#[inline]
#[must_use]
pub const fn target(&self) -> &Path {
&self.target
}
#[inline]
#[must_use]
pub const fn state(&self) -> AsyncCopyOperationState {
self.state
}
#[inline]
#[must_use]
pub const fn has_recovery(&self) -> bool {
self.writer.is_some()
}
#[inline]
pub fn recovery(&mut self) -> Option<&mut AsyncWriterRecovery> {
self.writer.as_mut()
}
#[inline]
pub fn take_recovery(&mut self) -> Option<AsyncWriterRecovery> {
self.writer.take()
}
pub async fn execute(&mut self) -> Result<CopyOutcome, AsyncCopyFailure> {
if self.state != AsyncCopyOperationState::Ready {
return Err(invalid_state_failure(
&self.source,
&self.target,
self.file_system.properties().info().provider_id(),
self.recovery,
));
}
let Self {
file_system,
source,
target,
options,
state,
writer,
deadline,
recovery,
} = self;
let mut guard = CopyCancellationGuard::start(state, writer, recovery);
let result = execute_copy(file_system, source, target, options, *deadline, guard.writer_mut()).await;
guard.finish(&result);
result
}
}
async fn execute_copy(
filesystem: &AsyncFileSystem,
source: &Path,
target: &Path,
options: &ResolvedCopyOptions,
deadline: CopyDeadline,
writer: &mut Option<AsyncWriterRecovery>,
) -> Result<CopyOutcome, AsyncCopyFailure> {
let caller_options = options.options();
if caller_options.max_entries() == Some(0) {
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy entry limit was exceeded"),
CopyFailureState::Unchanged,
CopyStats::default(),
source,
target,
));
}
if deadline.expired() {
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
CopyFailureState::Unchanged,
CopyStats::default(),
source,
target,
));
}
if !filesystem.core().provider_supports(ProviderOperation::TryCopy) {
return stream_copy_fallback(filesystem, source, target, options, deadline, writer).await;
}
match filesystem
.spi()
.try_copy(CopyRequest::new(source, target, options.clone()))
.await
{
Ok(CopyAttempt::Completed(outcome)) => {
let outcome = filesystem.verify_completed_copy(outcome, options.options(), source, target)?;
if deadline.expired() {
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
from_completed_stats(outcome.stats()),
*outcome.stats(),
source,
target,
));
}
Ok(outcome)
}
Ok(CopyAttempt::Declined(_)) => {
stream_copy_fallback(filesystem, source, target, options, deadline, writer).await
}
Err(failure) => {
let (error, state, stats) = failure.into_parts();
Err(filesystem.contextual_copy_failure(error, state, stats, source, target))
}
}
}
#[inline]
fn stream_copy_fallback<'a>(
filesystem: &'a AsyncFileSystem,
source: &'a Path,
target: &'a Path,
options: &'a ResolvedCopyOptions,
deadline: CopyDeadline,
writer_slot: &'a mut Option<AsyncWriterRecovery>,
) -> SpiFuture<'a, Result<CopyOutcome, AsyncCopyFailure>> {
Box::pin(async move {
let options = options.options();
let plan = StreamCopyPlan::new(options, filesystem.properties().limits(), source, target);
if let Err(error) = plan.validate_options(filesystem.properties().symlink_policy()) {
return Err(filesystem.contextual_copy_failure(
error,
CopyFailureState::Unchanged,
CopyStats::default(),
source,
target,
));
}
if deadline.expired() {
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
CopyFailureState::Unchanged,
CopyStats::default(),
source,
target,
));
}
filesystem
.require(FileSystemCapability::Read, FsOperation::Copy, source)
.and_then(|_| filesystem.require(FileSystemCapability::Write, FsOperation::Copy, target))
.map_err(|error| {
filesystem.contextual_copy_failure(
error,
CopyFailureState::Unchanged,
CopyStats::default(),
source,
target,
)
})?;
let metadata = filesystem.stat(source).await.map_err(|error| {
filesystem.contextual_copy_failure(error, CopyFailureState::Unchanged, CopyStats::default(), source, target)
})?;
if deadline.expired() {
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
CopyFailureState::Unchanged,
CopyStats::default(),
source,
target,
));
}
if let Err(error) = plan.validate_metadata(&metadata) {
return Err(filesystem.contextual_copy_failure(
error,
CopyFailureState::Unchanged,
CopyStats::default(),
source,
target,
));
}
let mut reader = filesystem
.open_reader(source, ReadOptions::default())
.await
.map_err(|error| {
filesystem.contextual_copy_failure(
error,
CopyFailureState::Unchanged,
CopyStats::default(),
source,
target,
)
})?;
if deadline.expired() {
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
CopyFailureState::Unchanged,
CopyStats::default(),
source,
target,
));
}
let writer_options = plan.writer_options();
match filesystem.open_writer(target, writer_options).await {
Ok(writer) => *writer_slot = Some(AsyncWriterRecovery::Opened(Box::new(writer))),
Err(error)
if error.stage() == OpenFailureStage::ProviderOpen
&& error.recovery().is_none()
&& error.error().kind() == FsErrorKind::AlreadyExists
&& is_unchanged_open_failure(error.error())
&& options.conflict() == CopyConflictPolicy::Skip =>
{
return Ok(CopyOutcome::streamed_fallback(
CopyStats {
skipped: 1,
..CopyStats::default()
},
crate::metadata::AchievedAtomicity::NonAtomic,
false,
));
}
Err(error) => {
let (error, stage, recovery) = error.into_parts();
*writer_slot = recovery.map(AsyncWriterRecovery::Rejected);
let state = match stage {
OpenFailureStage::Preflight => CopyFailureState::Unchanged,
OpenFailureStage::ProviderOpen => from_write_failure_state(open_failure_state(&error)),
OpenFailureStage::OutcomeValidation => CopyFailureState::Indeterminate,
};
return Err(filesystem.contextual_copy_failure(error, state, CopyStats::default(), source, target));
}
}
if deadline.expired() {
let writer = writer_slot
.as_ref()
.and_then(AsyncWriterRecovery::opened)
.expect("writer is retained before transfer");
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
from_writer_state(writer.state()),
fallback_failure_stats(writer.written_bytes()),
source,
target,
));
}
let mut bytes = 0_u64;
let mut buffer = [0_u8; 8192];
loop {
if deadline.expired() {
let writer = writer_slot
.as_ref()
.and_then(AsyncWriterRecovery::opened)
.expect("writer is retained before transfer");
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
from_writer_state(writer.state()),
fallback_failure_stats(writer.written_bytes()),
source,
target,
));
}
let read = reader.read_async(&mut buffer).await.map_err(|error| {
filesystem.contextual_copy_failure(
FsError::from_stream_io(error, FsOperation::Read, source),
from_writer_state(
writer_slot
.as_ref()
.and_then(AsyncWriterRecovery::opened)
.expect("writer is retained before transfer")
.state(),
),
fallback_failure_stats(
writer_slot
.as_ref()
.and_then(AsyncWriterRecovery::opened)
.expect("writer is retained before transfer")
.written_bytes(),
),
source,
target,
)
})?;
if deadline.expired() {
let writer = writer_slot
.as_ref()
.and_then(AsyncWriterRecovery::opened)
.expect("writer is retained before transfer");
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
from_writer_state(writer.state()),
fallback_failure_stats(writer.written_bytes()),
source,
target,
));
}
if read == 0 {
break;
}
let writer = writer_slot
.as_mut()
.and_then(AsyncWriterRecovery::opened_mut)
.expect("writer is retained before transfer");
let next_bytes = plan.next_bytes(bytes, read).map_err(|error| {
filesystem.contextual_copy_failure(
error,
from_writer_state(writer.state()),
fallback_failure_stats(writer.written_bytes()),
source,
target,
)
})?;
writer.write_fully_async(&buffer[..read]).await.map_err(|error| {
filesystem.contextual_copy_failure(
FsError::from_stream_io(error, FsOperation::Write, target),
from_writer_state(writer.state()),
fallback_failure_stats(writer.written_bytes()),
source,
target,
)
})?;
if deadline.expired() {
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
from_writer_state(writer.state()),
fallback_failure_stats(writer.written_bytes()),
source,
target,
));
}
bytes = next_bytes;
}
let writer = writer_slot
.as_mut()
.and_then(AsyncWriterRecovery::opened_mut)
.expect("writer is retained before flush");
if deadline.expired() {
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
from_writer_state(writer.state()),
fallback_failure_stats(writer.written_bytes()),
source,
target,
));
}
writer.flush_async().await.map_err(|error| {
filesystem.contextual_copy_failure(
FsError::from_stream_io(error, FsOperation::Write, target),
from_writer_state(writer.state()),
fallback_failure_stats(writer.written_bytes()),
source,
target,
)
})?;
if deadline.expired() {
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
from_writer_state(writer.state()),
fallback_failure_stats(writer.written_bytes()),
source,
target,
));
}
let writer = writer_slot
.as_mut()
.and_then(AsyncWriterRecovery::opened_mut)
.expect("writer is retained before commit");
let write_outcome = match writer.commit_async().await {
Ok(outcome) => outcome,
Err(failure)
if failure.error().kind() == FsErrorKind::AlreadyExists
&& plan.may_skip_conflict(from_writer_state(writer.state())) =>
{
if let Err(cleanup_error) = writer.abort_async().await {
return Err(filesystem.contextual_copy_failure(
cleanup_error,
from_writer_state(writer.state()),
fallback_failure_stats(writer.written_bytes()),
source,
target,
));
}
let _ = writer_slot.take();
return Ok(CopyOutcome::streamed_fallback(
CopyStats {
skipped: 1,
..CopyStats::default()
},
crate::metadata::AchievedAtomicity::NonAtomic,
false,
));
}
Err(failure) => {
return Err(filesystem.contextual_copy_failure(
failure.into_error(),
from_writer_state(writer.state()),
fallback_failure_stats(writer.written_bytes()),
source,
target,
));
}
};
if deadline.expired() {
let _ = writer_slot.take();
return Err(filesystem.contextual_copy_failure(
budget_error(source, target, "copy deadline was exceeded"),
CopyFailureState::Published,
StreamCopyPlan::completed_stats(bytes),
source,
target,
));
}
let _ = writer_slot.take();
Ok(CopyOutcome::streamed_fallback(
StreamCopyPlan::completed_stats(bytes),
write_outcome.atomicity(),
write_outcome.durable(),
))
})
}
fn budget_error(source: &Path, target: &Path, message: &str) -> FsError {
FsError::new(FsErrorKind::ResourceLimitExceeded, FsOperation::Copy, message)
.with_path(source.clone())
.with_target(target.clone())
}
fn invalid_state_failure(
source: &Path,
target: &Path,
provider: &str,
snapshot: CopyRecoverySnapshot,
) -> AsyncCopyFailure {
AsyncCopyFailure::new(
FsError::new(
FsErrorKind::InvalidState,
FsOperation::Copy,
"copy operation cannot execute in its current state",
)
.with_path(source.clone())
.with_target(target.clone())
.with_provider(provider),
snapshot.state,
snapshot.stats,
)
}