1use 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
46pub struct AsyncCopyOperation {
70 pub(crate) file_system: AsyncFileSystem,
72 source: Path,
74 target: Path,
76 options: ResolvedCopyOptions,
78 state: AsyncCopyOperationState,
80 writer: Option<AsyncWriterRecovery>,
82 deadline: CopyDeadline,
84 recovery: CopyRecoverySnapshot,
86}
87
88impl AsyncCopyOperation {
89 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 #[inline]
114 #[must_use]
115 pub const fn source(&self) -> &Path {
116 &self.source
117 }
118
119 #[inline]
121 #[must_use]
122 pub const fn target(&self) -> &Path {
123 &self.target
124 }
125
126 #[inline]
128 #[must_use]
129 pub const fn state(&self) -> AsyncCopyOperationState {
130 self.state
131 }
132
133 #[inline]
135 #[must_use]
136 pub const fn has_recovery(&self) -> bool {
137 self.writer.is_some()
138 }
139
140 #[inline]
142 pub fn recovery(&mut self) -> Option<&mut AsyncWriterRecovery> {
143 self.writer.as_mut()
144 }
145
146 #[inline]
148 pub fn take_recovery(&mut self) -> Option<AsyncWriterRecovery> {
149 self.writer.take()
150 }
151
152 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
194async 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#[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
558fn 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
565fn 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}