1use clap::Args;
6use jiff::Timestamp;
7use rc_core::alias::RetryConfig;
8use rc_core::{
9 AliasManager, Error, MetadataDirective, MultipartCopyCancellation, MultipartCopyOptions,
10 ObjectAttributes, ObjectEncryptionRequest, ObjectInfo, ObjectKeyPolicy, ObjectStore as _,
11 ObjectWriteOptions, ParsedPath, RemotePath, SseCustomerKey, TransferCancellation,
12 TransferCandidate, TransferControls, TransferCopyOptions, TransferExecutor,
13 TransferOutcomeState, TransferPlan, TransferReadOptions, TransferSelection,
14 normalize_relative_key, parse_path, relative_local_path_from_key,
15};
16use rc_s3::S3Client;
17use serde::Serialize;
18use std::collections::{BTreeSet, HashMap, HashSet};
19use std::fmt;
20use std::path::{Path, PathBuf};
21use std::sync::{Arc, Mutex as StdMutex};
22use tokio::sync::Mutex as AsyncMutex;
23
24use crate::exit_code::ExitCode;
25use crate::output::{Formatter, OutputConfig, ProgressBar, V3SuccessEnvelope};
26use crate::secret_input::{SecretLocator, resolve_secret_locator};
27
28use super::object_identity::set_source_identity;
29use super::transfer_fidelity::{MetadataDirectiveArg, TaggingDirectiveArg, TransferFidelityArgs};
30
31const CP_AFTER_HELP: &str = "\
32Examples:
33 rc object copy ./report.json local/my-bucket/reports/
34 rc cp ./report.json local/my-bucket/reports/
35 rc object copy local/source-bucket/archive.tar.gz ./downloads/archive.tar.gz";
36
37pub(crate) const GET_AFTER_HELP: &str = "\
38Examples:
39 rc get local/my-bucket/report.json ./report.json
40 rc get local/my-bucket/archive.tar.gz ./downloads/archive.tar.gz";
41
42pub(crate) const PUT_AFTER_HELP: &str = "\
43Examples:
44 rc put ./report.json local/my-bucket/reports/
45 rc put ./january.csv ./february.csv local/my-bucket/reports/";
46
47const REMOTE_PATH_SUGGESTION: &str =
48 "Use a local filesystem path or a remote path in the form alias/bucket[/key].";
49const DEFAULT_TRANSFER_CONCURRENCY: usize = 4;
50const DEFAULT_RETRY_ATTEMPTS: u32 = 3;
51const DEFAULT_RETRY_INITIAL_BACKOFF_MS: u64 = 100;
52const DEFAULT_RETRY_MAX_BACKOFF_MS: u64 = 10_000;
53#[cfg(test)]
54const MAX_SINGLE_COPY_SIZE: u64 = rc_core::S3_SINGLE_COPY_MAX_SIZE;
55
56#[derive(Args, Clone)]
58#[command(after_help = CP_AFTER_HELP)]
59pub struct CpArgs {
60 #[arg(required = true, num_args = 1.., value_name = "SOURCE")]
62 pub sources: Vec<String>,
63
64 pub target: String,
66
67 #[arg(short, long)]
69 pub recursive: bool,
70
71 #[arg(short, long, conflicts_with = "metadata_directive")]
73 pub preserve: bool,
74
75 #[arg(long, value_enum)]
77 pub(crate) metadata_directive: Option<MetadataDirectiveArg>,
78
79 #[arg(long, value_enum)]
81 pub(crate) tagging_directive: Option<TaggingDirectiveArg>,
82
83 #[arg(long)]
85 pub continue_on_error: bool,
86
87 #[arg(
89 long,
90 default_value_t = true,
91 action = clap::ArgAction::Set,
92 num_args = 0..=1,
93 default_missing_value = "true"
94 )]
95 pub overwrite: bool,
96
97 #[arg(long)]
99 pub skip_existing: bool,
100
101 #[arg(long)]
103 pub dry_run: bool,
104
105 #[arg(long)]
107 pub storage_class: Option<String>,
108
109 #[arg(long)]
111 pub content_type: Option<String>,
112
113 #[command(flatten)]
114 pub(crate) fidelity: TransferFidelityArgs,
115
116 #[arg(long = "enc-s3")]
118 pub enc_s3: Vec<String>,
119
120 #[arg(long = "enc-kms")]
122 pub enc_kms: Vec<String>,
123
124 #[arg(long = "enc-c-source-key-file")]
126 pub enc_c_source_key_file: Option<PathBuf>,
127
128 #[arg(long = "enc-c-source-key-env")]
130 pub enc_c_source_key_env: Option<String>,
131
132 #[arg(long = "enc-c-destination-key-file")]
134 pub enc_c_destination_key_file: Option<PathBuf>,
135
136 #[arg(long = "enc-c-destination-key-env")]
138 pub enc_c_destination_key_env: Option<String>,
139
140 #[arg(skip)]
141 pub(crate) source_customer_key: Option<SseCustomerKey>,
142
143 #[arg(skip)]
144 pub(crate) destination_customer_key: Option<SseCustomerKey>,
145
146 #[arg(long)]
148 pub include: Vec<String>,
149
150 #[arg(long)]
152 pub exclude: Vec<String>,
153
154 #[arg(long)]
156 pub newer_than: Option<String>,
157
158 #[arg(long)]
160 pub older_than: Option<String>,
161
162 #[arg(long)]
164 pub rewind: Option<String>,
165
166 #[arg(long)]
168 pub concurrency: Option<usize>,
169
170 #[arg(long)]
172 pub rate_limit: Option<String>,
173
174 #[arg(long)]
176 pub retry_attempts: Option<u32>,
177
178 #[arg(long)]
180 pub retry_initial_backoff_ms: Option<u64>,
181
182 #[arg(long)]
184 pub retry_max_backoff_ms: Option<u64>,
185
186 #[arg(long)]
188 pub fail_empty: bool,
189
190 #[arg(long)]
192 pub summary: bool,
193
194 #[arg(long)]
196 pub portable_names: bool,
197}
198
199impl fmt::Debug for CpArgs {
200 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
201 formatter.write_str("CpArgs { .. }")
202 }
203}
204
205impl CpArgs {
206 pub(crate) fn single(source: impl Into<String>, target: impl Into<String>) -> Self {
207 Self {
208 sources: vec![source.into()],
209 target: target.into(),
210 recursive: false,
211 preserve: false,
212 metadata_directive: None,
213 tagging_directive: None,
214 continue_on_error: false,
215 overwrite: true,
216 skip_existing: false,
217 dry_run: false,
218 storage_class: None,
219 content_type: None,
220 fidelity: TransferFidelityArgs::default(),
221 enc_s3: Vec::new(),
222 enc_kms: Vec::new(),
223 enc_c_source_key_file: None,
224 enc_c_source_key_env: None,
225 enc_c_destination_key_file: None,
226 enc_c_destination_key_env: None,
227 source_customer_key: None,
228 destination_customer_key: None,
229 include: Vec::new(),
230 exclude: Vec::new(),
231 newer_than: None,
232 older_than: None,
233 rewind: None,
234 concurrency: None,
235 rate_limit: None,
236 retry_attempts: None,
237 retry_initial_backoff_ms: None,
238 retry_max_backoff_ms: None,
239 fail_empty: false,
240 summary: false,
241 portable_names: false,
242 }
243 }
244
245 fn local_key_policy(&self) -> ObjectKeyPolicy {
246 ObjectKeyPolicy::for_local_destination(self.portable_names)
247 }
248}
249
250#[derive(Args, Debug)]
252#[command(
253 override_usage = "rc get [OPTIONS] <SOURCE> <TARGET>",
254 after_help = GET_AFTER_HELP
255)]
256pub struct GetArgs {
257 #[command(flatten)]
258 pub transfer: CpArgs,
259}
260
261#[derive(Args, Debug)]
263#[command(after_help = PUT_AFTER_HELP)]
264pub struct PutArgs {
265 #[command(flatten)]
266 pub transfer: CpArgs,
267}
268
269#[derive(Clone, Copy, Debug, Eq, PartialEq)]
270enum TransferAlias {
271 Copy,
272 Get,
273 Put,
274}
275
276#[derive(Debug, Serialize)]
277struct CpOutput {
278 status: &'static str,
279 source: String,
280 target: String,
281 #[serde(skip_serializing_if = "Option::is_none")]
282 size_bytes: Option<i64>,
283 #[serde(skip_serializing_if = "Option::is_none")]
284 size_human: Option<String>,
285 #[serde(skip_serializing_if = "Option::is_none")]
286 version_id: Option<String>,
287 #[serde(skip_serializing_if = "Option::is_none")]
288 source_version_id: Option<String>,
289}
290
291#[derive(Debug, Serialize)]
292struct VersionCopyData {
293 operation: &'static str,
294 source: String,
295 target: String,
296 source_version_id: Option<String>,
297 version_id: Option<String>,
298 size_bytes: Option<i64>,
299 size_human: Option<String>,
300}
301
302pub async fn execute(args: CpArgs, output_config: OutputConfig) -> ExitCode {
304 execute_with_alias(args, output_config, TransferAlias::Copy).await
305}
306
307pub async fn execute_get(args: GetArgs, output_config: OutputConfig) -> ExitCode {
309 execute_with_alias(args.transfer, output_config, TransferAlias::Get).await
310}
311
312pub async fn execute_put(args: PutArgs, output_config: OutputConfig) -> ExitCode {
314 execute_with_alias(args.transfer, output_config, TransferAlias::Put).await
315}
316
317async fn execute_with_alias(
318 mut args: CpArgs,
319 output_config: OutputConfig,
320 alias: TransferAlias,
321) -> ExitCode {
322 let formatter = Formatter::new(output_config);
323
324 if alias == TransferAlias::Get && args.sources.len() != 1 {
325 return formatter.fail(
326 ExitCode::UsageError,
327 "get requires exactly one remote source and one local target",
328 );
329 }
330
331 if let Err(error) = validate_destination_storage_class(args.storage_class.as_deref()) {
332 return formatter.fail(exit_code_for_core_error(&error), &error.to_string());
333 }
334 if let Err(error) = object_write_options(
335 &args.fidelity,
336 args.content_type.as_deref(),
337 None,
338 None,
339 args.storage_class.clone(),
340 ) {
341 return formatter.fail(exit_code_for_core_error(&error), &error.to_string());
342 }
343
344 if formatter.is_json() && uses_transfer_planner(&args) {
345 return formatter.fail(
346 ExitCode::UnsupportedFeature,
347 "Bulk transfer planning currently supports human output only; JSON batch output requires the versioned output contract",
348 );
349 }
350
351 let selection = match build_transfer_selection(&args, Timestamp::now()) {
352 Ok(selection) => selection,
353 Err(error) => return formatter.fail(ExitCode::UsageError, &error),
354 };
355 let controls = match build_transfer_controls(&args) {
356 Ok(controls) => controls,
357 Err(error) => return formatter.fail(ExitCode::UsageError, &error),
358 };
359
360 let alias_manager = AliasManager::new();
361
362 let mut sources = Vec::with_capacity(args.sources.len());
364 for source in &args.sources {
365 let parsed = match alias {
366 TransferAlias::Put => Ok(ParsedPath::Local(PathBuf::from(source))),
369 TransferAlias::Copy | TransferAlias::Get => {
370 parse_cp_path(source, alias_manager.as_ref().ok())
371 }
372 };
373 match parsed {
374 Ok(path) => sources.push(path),
375 Err(error) => {
376 return formatter.fail_with_suggestion(
377 ExitCode::UsageError,
378 &format!("Invalid source path '{source}': {error}"),
379 REMOTE_PATH_SUGGESTION,
380 );
381 }
382 }
383 }
384
385 let parsed_target = match alias {
386 TransferAlias::Get => Ok(ParsedPath::Local(PathBuf::from(&args.target))),
389 TransferAlias::Copy | TransferAlias::Put => {
390 parse_cp_path(&args.target, alias_manager.as_ref().ok())
391 }
392 };
393 let target = match parsed_target {
394 Ok(p) => p,
395 Err(e) => {
396 return formatter.fail_with_suggestion(
397 ExitCode::UsageError,
398 &format!("Invalid target path: {e}"),
399 REMOTE_PATH_SUGGESTION,
400 );
401 }
402 };
403 if let Err(error) = validate_alias_direction(alias, &sources, &target) {
404 return formatter.fail(ExitCode::UsageError, error);
405 }
406 if let Err(error) = validate_fidelity_directions(&args, &sources, &target) {
407 return formatter.fail(exit_code_for_core_error(&error), &error.to_string());
408 }
409 let source_key_locator = match resolve_secret_locator(
410 args.enc_c_source_key_file.clone(),
411 args.enc_c_source_key_env.clone(),
412 ) {
413 Ok(locator) => locator,
414 Err(error) => {
415 return formatter.fail(exit_code_for_core_error(&error), &error.to_string());
416 }
417 };
418 let destination_key_locator = match resolve_secret_locator(
419 args.enc_c_destination_key_file.clone(),
420 args.enc_c_destination_key_env.clone(),
421 ) {
422 Ok(locator) => locator,
423 Err(error) => {
424 return formatter.fail(exit_code_for_core_error(&error), &error.to_string());
425 }
426 };
427 if source_key_locator.is_some()
428 && sources
429 .iter()
430 .any(|path| matches!(path, ParsedPath::Local(_)))
431 {
432 return formatter.fail(
433 ExitCode::UsageError,
434 "SSE-C source keys require every source to be remote",
435 );
436 }
437 if destination_key_locator.is_some() && matches!(target, ParsedPath::Local(_)) {
438 return formatter.fail(
439 ExitCode::UsageError,
440 "SSE-C destination keys require a remote destination",
441 );
442 }
443 if args.storage_class.is_some() && matches!(target, ParsedPath::Local(_)) {
444 return formatter.fail(
445 ExitCode::UsageError,
446 "--storage-class requires a remote destination",
447 );
448 }
449
450 let target_is_container = is_container_target(&args.target, &target);
451 if args.sources.len() > 1 && !target_is_container {
452 return formatter.fail(
453 ExitCode::UsageError,
454 "Multiple copy sources require a directory or remote prefix destination ending in '/'",
455 );
456 }
457
458 let target_encryption = match parse_destination_encryption(&args.enc_s3, &args.enc_kms, &target)
459 {
460 Ok(encryption) => encryption,
461 Err(error) => return formatter.fail(ExitCode::UsageError, &error),
462 };
463 if destination_key_locator.is_some() && target_encryption.is_some() {
464 return formatter.fail(
465 ExitCode::UsageError,
466 "SSE-C destination keys cannot be combined with --enc-s3 or --enc-kms",
467 );
468 }
469 if !args.dry_run
470 && (source_key_locator.is_some() || destination_key_locator.is_some())
471 && sources
472 .iter()
473 .all(|path| matches!(path, ParsedPath::Remote(_)))
474 && matches!(target, ParsedPath::Remote(_))
475 {
476 return formatter.fail(
477 ExitCode::UnsupportedFeature,
478 "RustFS beta.10 server-side SSE-C copy is not compatibility-proven; tracked by rustfs/backlog#1467",
479 );
480 }
481 if !args.dry_run {
482 args.source_customer_key = match load_customer_key(source_key_locator.as_ref()) {
483 Ok(key) => key,
484 Err(error) => {
485 return formatter.fail(exit_code_for_core_error(&error), &error.to_string());
486 }
487 };
488 args.destination_customer_key = match load_customer_key(destination_key_locator.as_ref()) {
489 Ok(key) => key,
490 Err(error) => {
491 return formatter.fail(exit_code_for_core_error(&error), &error.to_string());
492 }
493 };
494 }
495
496 let single_local_to_local = matches!(
497 (sources.as_slice(), &target),
498 ([ParsedPath::Local(_)], ParsedPath::Local(_))
499 );
500 if uses_transfer_planner(&args) && !single_local_to_local {
501 let alias_manager = match alias_manager.as_ref() {
502 Ok(manager) => manager,
503 Err(error) => {
504 return formatter.fail(
505 ExitCode::UsageError,
506 &format!("Failed to load aliases: {error}"),
507 );
508 }
509 };
510 return execute_transfer_plan(
511 &args,
512 sources,
513 target,
514 target_is_container,
515 target_encryption,
516 selection,
517 controls,
518 &formatter,
519 alias_manager,
520 )
521 .await;
522 }
523
524 let Some(source) = sources.first() else {
525 return formatter.fail(ExitCode::UsageError, "At least one copy source is required");
526 };
527
528 execute_single_copy(
529 source,
530 &target,
531 &args,
532 &formatter,
533 target_encryption.as_ref(),
534 )
535 .await
536}
537
538fn validate_alias_direction(
539 alias: TransferAlias,
540 sources: &[ParsedPath],
541 target: &ParsedPath,
542) -> Result<(), &'static str> {
543 match alias {
544 TransferAlias::Copy => Ok(()),
545 TransferAlias::Get if sources.len() == 1 && sources[0].is_remote() && target.is_local() => {
546 Ok(())
547 }
548 TransferAlias::Get => Err("get requires exactly one remote source and one local target"),
549 TransferAlias::Put if sources.iter().all(ParsedPath::is_local) && target.is_remote() => {
550 Ok(())
551 }
552 TransferAlias::Put => Err("put requires one or more local sources and one remote target"),
553 }
554}
555
556fn uses_transfer_planner(args: &CpArgs) -> bool {
557 args.sources.len() > 1
558 || args.recursive
559 || args.skip_existing
560 || !args.overwrite
561 || !args.include.is_empty()
562 || !args.exclude.is_empty()
563 || args.newer_than.is_some()
564 || args.older_than.is_some()
565 || args.rewind.is_some()
566 || args.concurrency.is_some()
567 || args.rate_limit.is_some()
568 || args.retry_attempts.is_some()
569 || args.retry_initial_backoff_ms.is_some()
570 || args.retry_max_backoff_ms.is_some()
571 || args.fail_empty
572 || args.summary
573}
574
575pub(super) fn validate_destination_storage_class(value: Option<&str>) -> rc_core::Result<()> {
576 let Some(value) = value else {
577 return Ok(());
578 };
579 match value {
580 "STANDARD" | "REDUCED_REDUNDANCY" => Ok(()),
581 "DEEP_ARCHIVE"
582 | "EXPRESS_ONEZONE"
583 | "FSX_ONTAP"
584 | "FSX_OPENZFS"
585 | "GLACIER"
586 | "GLACIER_IR"
587 | "INTELLIGENT_TIERING"
588 | "ONEZONE_IA"
589 | "OUTPOSTS"
590 | "SNOW"
591 | "STANDARD_IA" => Err(Error::UnsupportedFeature(format!(
592 "RustFS beta.10 does not provide meaningful storage policy '{value}'"
593 ))),
594 value => Err(Error::InvalidPath(format!(
595 "Unknown destination storage class '{value}'"
596 ))),
597 }
598}
599
600fn object_write_options(
601 fidelity: &TransferFidelityArgs,
602 content_type: Option<&str>,
603 encryption: Option<&ObjectEncryptionRequest>,
604 customer_key: Option<&SseCustomerKey>,
605 storage_class: Option<String>,
606) -> rc_core::Result<ObjectWriteOptions> {
607 fidelity.build_write_options(content_type, encryption, customer_key, storage_class)
608}
609
610fn requested_metadata_directive(args: &CpArgs) -> Option<MetadataDirective> {
611 if args.preserve {
612 Some(MetadataDirective::Copy)
613 } else {
614 args.metadata_directive.map(Into::into)
615 }
616}
617
618fn transfer_copy_options(
619 args: &CpArgs,
620 source_version_id: Option<String>,
621 source_etag: Option<String>,
622 encryption: Option<&ObjectEncryptionRequest>,
623) -> rc_core::Result<TransferCopyOptions> {
624 let metadata_directive = requested_metadata_directive(args);
625 let tagging_directive = args.tagging_directive.map(Into::into);
626 let mut destination = object_write_options(
627 &args.fidelity,
628 args.content_type.as_deref(),
629 encryption,
630 args.destination_customer_key.as_ref(),
631 args.storage_class.clone(),
632 )?;
633 if matches!(metadata_directive, Some(MetadataDirective::Replace))
634 && destination.attributes.is_none()
635 {
636 destination.attributes = Some(Default::default());
637 }
638 if matches!(tagging_directive, Some(rc_core::TaggingDirective::Replace))
639 && destination.tags.is_none()
640 {
641 destination.tags = Some(Default::default());
642 }
643 let options = TransferCopyOptions {
644 source: TransferReadOptions {
645 version_id: source_version_id,
646 customer_key: args.source_customer_key.clone(),
647 ..TransferReadOptions::default()
648 },
649 source_etag,
650 metadata_directive,
651 tagging_directive,
652 destination,
653 };
654 options.validate()?;
655 Ok(options)
656}
657
658fn validate_fidelity_directions(
659 args: &CpArgs,
660 sources: &[ParsedPath],
661 target: &ParsedPath,
662) -> rc_core::Result<()> {
663 let all_remote = sources
664 .iter()
665 .all(|source| matches!(source, ParsedPath::Remote(_)));
666 let any_remote = sources
667 .iter()
668 .any(|source| matches!(source, ParsedPath::Remote(_)));
669 let target_remote = matches!(target, ParsedPath::Remote(_));
670 let has_copy_directive =
671 args.preserve || args.metadata_directive.is_some() || args.tagging_directive.is_some();
672 if has_copy_directive && !(all_remote && target_remote) {
673 return Err(Error::InvalidPath(
674 "Metadata and tagging copy directives require remote sources and a remote destination"
675 .to_string(),
676 ));
677 }
678 if !target_remote
679 && (args.storage_class.is_some()
680 || args.content_type.is_some()
681 || args.fidelity.has_write_policy()
682 || !args.enc_s3.is_empty()
683 || !args.enc_kms.is_empty()
684 || args.enc_c_destination_key_file.is_some()
685 || args.enc_c_destination_key_env.is_some())
686 {
687 return Err(Error::InvalidPath(
688 "Destination transfer policies require a remote destination".to_string(),
689 ));
690 }
691 let same_alias_remote_copy = target_remote
692 && sources.iter().any(|source| {
693 matches!(
694 source,
695 ParsedPath::Remote(source)
696 if target
697 .as_remote()
698 .is_some_and(|target| source.alias == target.alias)
699 )
700 });
701 if any_remote && target_remote {
702 let copy_options = transfer_copy_options(args, None, None, None)?;
703 if same_alias_remote_copy
704 && matches!(
705 copy_options.metadata_directive,
706 Some(MetadataDirective::Replace)
707 )
708 {
709 return Err(Error::UnsupportedFeature(
710 "RustFS beta.10 does not preserve complete metadata REPLACE semantics; tracked by rustfs/backlog#1463"
711 .to_string(),
712 ));
713 }
714 if copy_options.tagging_directive.is_some() || copy_options.destination.tags.is_some() {
715 return Err(Error::UnsupportedFeature(
716 "RustFS beta.10 does not preserve CopyObject tagging directives; tracked by rustfs/backlog#1462"
717 .to_string(),
718 ));
719 }
720 if copy_options.destination.checksum.is_some() {
721 return Err(Error::UnsupportedFeature(
722 "RustFS beta.10 does not preserve CopyObject checksum selection; tracked by rustfs/backlog#1466"
723 .to_string(),
724 ));
725 }
726 }
727 Ok(())
728}
729
730fn load_customer_key(locator: Option<&SecretLocator>) -> rc_core::Result<Option<SseCustomerKey>> {
731 locator.map(SecretLocator::load_customer_key).transpose()
732}
733
734async fn execute_single_copy(
735 source: &ParsedPath,
736 target: &ParsedPath,
737 args: &CpArgs,
738 formatter: &Formatter,
739 encryption: Option<&ObjectEncryptionRequest>,
740) -> ExitCode {
741 match (source, target) {
743 (ParsedPath::Local(src), ParsedPath::Remote(dst)) => {
744 copy_local_to_s3_prepared(src, dst, args, formatter, encryption).await
745 }
746 (ParsedPath::Remote(src), ParsedPath::Local(dst)) => {
747 copy_s3_to_local(src, dst, args, formatter).await
748 }
749 (ParsedPath::Remote(src), ParsedPath::Remote(dst)) => {
750 if args.recursive {
751 return formatter.fail(
752 ExitCode::UnsupportedFeature,
753 "Recursive S3-to-S3 copy is not implemented",
754 );
755 }
756 copy_s3_to_s3_prepared(src, dst, args, formatter, encryption).await
757 }
758 (ParsedPath::Local(_), ParsedPath::Local(_)) => formatter.fail_with_suggestion(
759 ExitCode::UsageError,
760 "Cannot copy between two local paths. Use system cp command.",
761 "Use your local shell cp command when both paths are on the filesystem.",
762 ),
763 }
764}
765
766#[derive(Debug, Clone)]
767enum CpOperation {
768 LocalToRemote {
769 source: PathBuf,
770 target: RemotePath,
771 encryption: Option<ObjectEncryptionRequest>,
772 },
773 RemoteToLocal {
774 source: RemotePath,
775 target: PathBuf,
776 },
777 RemoteToRemote {
778 source: RemotePath,
779 target: RemotePath,
780 source_info: Box<ObjectInfo>,
781 encryption: Option<ObjectEncryptionRequest>,
782 },
783}
784
785#[derive(Debug, Clone)]
786struct PlannedCopyDetail {
787 source_version_id: Option<String>,
788 destination_version_id: Option<String>,
789 upload_id: Option<String>,
790}
791
792type PlannedCopyDetails = Arc<AsyncMutex<HashMap<(String, String), PlannedCopyDetail>>>;
793
794#[derive(Debug)]
795struct PlannedCopyProgress {
796 positions: StdMutex<HashMap<(String, String), u64>>,
797 bar: Option<ProgressBar>,
798}
799
800impl PlannedCopyProgress {
801 fn new(output_config: OutputConfig, total: u64) -> Self {
802 Self {
803 positions: StdMutex::new(HashMap::new()),
804 bar: (total > 0).then(|| ProgressBar::new(output_config, total)),
805 }
806 }
807
808 fn reset(&self, key: &(String, String)) {
809 self.set(key, 0);
810 }
811
812 fn set(&self, key: &(String, String), bytes: u64) {
813 let mut positions = self
814 .positions
815 .lock()
816 .expect("planned copy progress lock should not be poisoned");
817 positions.insert(key.clone(), bytes);
818 let aggregate = positions.values().copied().fold(0_u64, u64::saturating_add);
819 if let Some(bar) = &self.bar {
820 bar.set_position(aggregate);
821 }
822 }
823
824 fn finish(&self) {
825 if let Some(bar) = &self.bar {
826 bar.finish_and_clear();
827 }
828 }
829}
830
831type SharedPlannedCopyProgress = Arc<PlannedCopyProgress>;
832
833#[allow(clippy::too_many_arguments)]
834async fn execute_transfer_plan(
835 args: &CpArgs,
836 sources: Vec<ParsedPath>,
837 target: ParsedPath,
838 target_is_container: bool,
839 encryption: Option<ObjectEncryptionRequest>,
840 selection: TransferSelection,
841 controls: TransferControls,
842 formatter: &Formatter,
843 alias_manager: &AliasManager,
844) -> ExitCode {
845 let candidates = match build_transfer_candidates(
846 &sources,
847 &target,
848 target_is_container,
849 args.recursive,
850 encryption,
851 args.source_customer_key.as_ref(),
852 alias_manager,
853 args.local_key_policy(),
854 )
855 .await
856 {
857 Ok(candidates) => candidates,
858 Err(error) => {
859 return formatter.fail(
860 exit_code_for_core_error(&error),
861 &format!("Failed to plan copy: {error}"),
862 );
863 }
864 };
865
866 let mut plan = TransferPlan::build(candidates, &selection);
867 if let Err(error) = validate_storage_class_plan(&plan, args.storage_class.as_deref()) {
868 return formatter.fail(exit_code_for_core_error(&error), &error.to_string());
869 }
870 if let Err(error) = validate_plan_targets(&plan) {
871 return formatter.fail(ExitCode::UsageError, &error.to_string());
872 }
873 if plan.items.is_empty() {
874 if args.fail_empty {
875 return formatter.fail(
876 ExitCode::NotFound,
877 "No copy sources matched the requested selection",
878 );
879 }
880 if args.summary || args.recursive || args.sources.len() > 1 {
881 print_transfer_summary(formatter, &plan.summary);
882 }
883 return ExitCode::Success;
884 }
885
886 let clients = match create_planned_client_cache(&plan.items, alias_manager).await {
887 Ok(clients) => Arc::new(clients),
888 Err(error) => {
889 return formatter.fail(
890 exit_code_for_core_error(&error),
891 &format!("Failed to prepare copy clients: {error}"),
892 );
893 }
894 };
895 let skipped_existing = if args.skip_existing || !args.overwrite {
896 match skip_existing_remote_targets(
897 &mut plan,
898 &clients,
899 args.destination_customer_key.as_ref(),
900 )
901 .await
902 {
903 Ok(skipped) => skipped,
904 Err(error) => {
905 return formatter.fail(
906 exit_code_for_core_error(&error),
907 &format!("Failed to inspect copy destinations: {error}"),
908 );
909 }
910 }
911 } else {
912 Vec::new()
913 };
914
915 if args.dry_run {
916 for item in &plan.items {
917 formatter.println(&format!(
918 "Would copy: {} -> {}{}",
919 formatter.style_file(&item.source),
920 formatter.style_file(&item.target),
921 transfer_policy_suffix(args)
922 ));
923 }
924 for item in &skipped_existing {
925 print_skipped_existing(formatter, item, true);
926 }
927 if args.summary || args.recursive || args.sources.len() > 1 {
928 print_transfer_summary(formatter, &plan.summary);
929 }
930 return ExitCode::Success;
931 }
932
933 for item in &skipped_existing {
934 print_skipped_existing(formatter, item, false);
935 }
936 if plan.items.is_empty() {
937 if args.summary || args.recursive || args.sources.len() > 1 {
938 print_transfer_summary(formatter, &plan.summary);
939 }
940 return ExitCode::Success;
941 }
942
943 let total_bytes = plan
944 .items
945 .iter()
946 .filter_map(|item| {
947 if matches!(item.payload, CpOperation::RemoteToRemote { .. }) {
948 item.size_bytes
949 } else {
950 None
951 }
952 })
953 .sum::<u64>();
954 let executor = match TransferExecutor::new(controls) {
955 Ok(executor) => executor,
956 Err(error) => {
957 return formatter.fail(ExitCode::UsageError, &error.to_string());
958 }
959 };
960 let operation_args = Arc::new(args.clone());
961 let transfer_cancellation = TransferCancellation::new();
962 let multipart_cancellation = MultipartCopyCancellation::new();
963 let signal_task = tokio::spawn({
964 let transfer_cancellation = transfer_cancellation.clone();
965 let multipart_cancellation = multipart_cancellation.clone();
966 async move {
967 if tokio::signal::ctrl_c().await.is_ok() {
968 multipart_cancellation.cancel();
969 transfer_cancellation.cancel();
970 }
971 }
972 });
973 let copy_details = PlannedCopyDetails::default();
974 let copy_progress = Arc::new(PlannedCopyProgress::new(
975 formatter.output_config(),
976 total_bytes,
977 ));
978 let report = executor
979 .execute_with_cancellation(plan, transfer_cancellation, {
980 let operation_args = Arc::clone(&operation_args);
981 let clients = Arc::clone(&clients);
982 let multipart_cancellation = multipart_cancellation.clone();
983 let copy_details = Arc::clone(©_details);
984 let copy_progress = Arc::clone(©_progress);
985 move |item| {
986 let operation_args = Arc::clone(&operation_args);
987 let clients = Arc::clone(&clients);
988 let multipart_cancellation = multipart_cancellation.clone();
989 let copy_details = Arc::clone(©_details);
990 let copy_progress = Arc::clone(©_progress);
991 async move {
992 execute_planned_operation(
993 item,
994 &operation_args,
995 &clients,
996 &multipart_cancellation,
997 ©_details,
998 ©_progress,
999 )
1000 .await
1001 }
1002 }
1003 })
1004 .await;
1005 signal_task.abort();
1006 let _ = signal_task.await;
1007 copy_progress.finish();
1008
1009 let copy_details = copy_details.lock().await;
1010 for outcome in &report.outcomes {
1011 match &outcome.state {
1012 TransferOutcomeState::Success { bytes_transferred } => {
1013 let key = (outcome.item.source.clone(), outcome.item.target.clone());
1014 print_planned_success(
1015 formatter,
1016 &outcome.item,
1017 *bytes_transferred,
1018 copy_details.get(&key),
1019 );
1020 }
1021 TransferOutcomeState::Failed { error } => {
1022 formatter.error_with_code(exit_code_for_core_error(error), &error.to_string());
1023 }
1024 TransferOutcomeState::Cancelled { error } => {
1025 if let Some(error) = error {
1026 formatter.warning(&format!(
1027 "Cancelled transfer: {} ({error})",
1028 outcome.item.source
1029 ));
1030 } else {
1031 formatter.warning(&format!(
1032 "Cancelled before transfer: {}",
1033 outcome.item.source
1034 ));
1035 }
1036 }
1037 }
1038 }
1039
1040 if args.summary || args.recursive || args.sources.len() > 1 {
1041 print_transfer_summary(formatter, &report.summary);
1042 }
1043
1044 if report.was_cancelled {
1045 ExitCode::Interrupted
1046 } else {
1047 report
1048 .first_failure()
1049 .map_or(ExitCode::Success, exit_code_for_core_error)
1050 }
1051}
1052
1053fn transfer_policy_suffix(args: &CpArgs) -> String {
1054 let mut policies = Vec::new();
1055 if let Some(value) = args.storage_class.as_deref() {
1056 policies.push(format!("storage-class={value}"));
1057 }
1058 if let Some(directive) = requested_metadata_directive(args) {
1059 policies.push(format!(
1060 "metadata={}",
1061 match directive {
1062 MetadataDirective::Copy => "copy",
1063 MetadataDirective::Replace => "replace",
1064 }
1065 ));
1066 } else if args.fidelity.has_attribute_policy() || args.content_type.is_some() {
1067 policies.push(format!("metadata=write({})", args.fidelity.metadata.len()));
1068 }
1069 if let Some(directive) = args.tagging_directive {
1070 policies.push(format!(
1071 "tags={}({})",
1072 match directive {
1073 TaggingDirectiveArg::Copy => "copy",
1074 TaggingDirectiveArg::Replace => "replace",
1075 },
1076 args.fidelity.tags.len()
1077 ));
1078 } else if args.fidelity.has_tag_policy() {
1079 policies.push(format!("tags=write({})", args.fidelity.tags.len()));
1080 }
1081 if args.fidelity.checksum.is_some() {
1082 policies.push("checksum=sha256".to_string());
1083 }
1084 if let Some(mode) = args.fidelity.retention_mode.as_deref() {
1085 policies.push(format!("retention={}", mode.to_ascii_lowercase()));
1086 }
1087 if let Some(state) = args.fidelity.legal_hold.as_deref() {
1088 policies.push(format!("legal-hold={}", state.to_ascii_lowercase()));
1089 }
1090 if args.enc_c_destination_key_file.is_some() || args.enc_c_destination_key_env.is_some() {
1091 policies.push("encryption=sse-c".to_string());
1092 } else if !args.enc_kms.is_empty() {
1093 policies.push("encryption=sse-kms".to_string());
1094 } else if !args.enc_s3.is_empty() {
1095 policies.push("encryption=sse-s3".to_string());
1096 }
1097 if policies.is_empty() {
1098 String::new()
1099 } else {
1100 format!(" [policy:{}]", policies.join(","))
1101 }
1102}
1103
1104fn validate_storage_class_plan(
1105 plan: &TransferPlan<CpOperation>,
1106 storage_class: Option<&str>,
1107) -> rc_core::Result<()> {
1108 if storage_class.is_none() {
1109 return Ok(());
1110 }
1111 for item in &plan.items {
1112 let size = item.size_bytes.ok_or_else(|| {
1113 Error::UnsupportedFeature(format!(
1114 "Cannot guarantee storage class for a transfer with unknown size: {}",
1115 item.source
1116 ))
1117 })?;
1118 let multipart = match &item.payload {
1119 CpOperation::LocalToRemote { .. } => size > MULTIPART_THRESHOLD,
1120 CpOperation::RemoteToRemote { source, target, .. } if source.alias != target.alias => {
1121 size > MULTIPART_THRESHOLD
1122 }
1123 CpOperation::RemoteToRemote { .. } => rc_core::requires_multipart_copy(size),
1124 CpOperation::RemoteToLocal { .. } => false,
1125 };
1126 if multipart {
1127 return Err(Error::UnsupportedFeature(format!(
1128 "RustFS beta.10 does not persist storage class for multipart transfer: {}",
1129 item.source
1130 )));
1131 }
1132 }
1133 Ok(())
1134}
1135
1136async fn execute_planned_operation(
1137 item: TransferCandidate<CpOperation>,
1138 args: &CpArgs,
1139 clients: &HashMap<String, Arc<S3Client>>,
1140 multipart_cancellation: &MultipartCopyCancellation,
1141 copy_details: &PlannedCopyDetails,
1142 copy_progress: &SharedPlannedCopyProgress,
1143) -> rc_core::Result<u64> {
1144 let client = planned_client(clients, operation_alias(&item.payload))?;
1145 match &item.payload {
1146 CpOperation::LocalToRemote {
1147 source,
1148 target,
1149 encryption,
1150 } => perform_planned_upload(client, source, target, encryption.as_ref(), args).await,
1151 CpOperation::RemoteToLocal { source, target } => {
1152 perform_planned_download(client, source, target, args).await
1153 }
1154 CpOperation::RemoteToRemote {
1155 source,
1156 target,
1157 source_info,
1158 encryption,
1159 } => {
1160 let source_client = planned_client(clients, &source.alias)?;
1161 let target_client = planned_client(clients, &target.alias)?;
1162 let progress_key = (item.source.clone(), item.target.clone());
1163 copy_progress.reset(&progress_key);
1164 let result = perform_planned_remote_copy(
1165 source_client,
1166 target_client,
1167 source,
1168 target,
1169 source_info,
1170 encryption.as_ref(),
1171 multipart_cancellation,
1172 &|bytes| copy_progress.set(&progress_key, bytes),
1173 args,
1174 )
1175 .await?;
1176 copy_progress.set(&progress_key, result.bytes_copied);
1177 copy_details.lock().await.insert(
1178 (item.source, item.target),
1179 PlannedCopyDetail {
1180 source_version_id: result.source_version_id,
1181 destination_version_id: result.destination_version_id,
1182 upload_id: result.upload_id,
1183 },
1184 );
1185 Ok(result.bytes_copied)
1186 }
1187 }
1188}
1189
1190async fn perform_planned_upload(
1191 client: &S3Client,
1192 source: &Path,
1193 target: &RemotePath,
1194 encryption: Option<&ObjectEncryptionRequest>,
1195 args: &CpArgs,
1196) -> rc_core::Result<u64> {
1197 let metadata = tokio::fs::metadata(source).await?;
1198 let file_size = metadata.len();
1199 let guessed_type = mime_guess::from_path(source)
1200 .first()
1201 .map(|mime| mime.essence_str().to_string());
1202 let content_type = select_upload_content_type(
1203 args.content_type.as_deref(),
1204 guessed_type.as_deref(),
1205 file_size,
1206 );
1207 let options = object_write_options(
1208 &args.fidelity,
1209 content_type,
1210 encryption,
1211 args.destination_customer_key.as_ref(),
1212 args.storage_class.clone(),
1213 )?;
1214 let info = client
1215 .put_object_from_path_with_options(target, source, &options, |_| {})
1216 .await?;
1217 Ok(info
1218 .size_bytes
1219 .and_then(|size| u64::try_from(size).ok())
1220 .unwrap_or(file_size))
1221}
1222
1223async fn perform_planned_download(
1224 client: &S3Client,
1225 source: &RemotePath,
1226 target: &Path,
1227 args: &CpArgs,
1228) -> rc_core::Result<u64> {
1229 if target.exists() && !args.overwrite {
1230 return Err(Error::Conflict(format!(
1231 "Destination exists: {}. Use --overwrite to replace.",
1232 target.display()
1233 )));
1234 }
1235 if let Some(parent) = target.parent() {
1236 tokio::fs::create_dir_all(parent).await?;
1237 }
1238 client
1239 .download_object_to_path_with_transfer_options(
1240 source,
1241 target,
1242 &TransferReadOptions {
1243 customer_key: args.source_customer_key.clone(),
1244 ..TransferReadOptions::default()
1245 },
1246 |_, _| {},
1247 )
1248 .await
1249}
1250
1251struct PlannedRemoteCopyResult {
1252 bytes_copied: u64,
1253 source_version_id: Option<String>,
1254 source_etag: Option<String>,
1255 destination_version_id: Option<String>,
1256 upload_id: Option<String>,
1257 object: ObjectInfo,
1258}
1259
1260#[allow(clippy::too_many_arguments)]
1261async fn perform_planned_remote_copy(
1262 source_client: &S3Client,
1263 target_client: &S3Client,
1264 source: &RemotePath,
1265 target: &RemotePath,
1266 source_info: &ObjectInfo,
1267 encryption: Option<&ObjectEncryptionRequest>,
1268 cancellation: &MultipartCopyCancellation,
1269 on_progress: &(dyn Fn(u64) + Send + Sync),
1270 args: &CpArgs,
1271) -> rc_core::Result<PlannedRemoteCopyResult> {
1272 if source.alias != target.alias {
1273 return perform_cross_alias_remote_copy(
1274 source_client,
1275 target_client,
1276 source,
1277 target,
1278 source_info,
1279 encryption,
1280 on_progress,
1281 args,
1282 )
1283 .await;
1284 }
1285 if args.source_customer_key.is_some() || args.destination_customer_key.is_some() {
1286 return Err(Error::UnsupportedFeature(
1287 "RustFS beta.10 server-side SSE-C copy is not compatibility-proven; tracked by rustfs/backlog#1467"
1288 .to_string(),
1289 ));
1290 }
1291 let planned_size = source_info
1292 .size_bytes
1293 .and_then(|size| u64::try_from(size).ok())
1294 .ok_or_else(|| Error::InvalidPath(format!("Source size is unavailable: {source}")))?;
1295 if rc_core::requires_multipart_copy(planned_size) {
1296 if args.storage_class.is_some() {
1297 return Err(Error::UnsupportedFeature(
1298 "RustFS beta.10 does not persist storage class for multipart copies".to_string(),
1299 ));
1300 }
1301 let current = source_client.head_object(source).await?;
1302 if !source_identity_matches(source_info, ¤t) {
1303 return Err(Error::Conflict(format!(
1304 "Source changed after copy planning: {source}"
1305 )));
1306 }
1307 let options = multipart_options_from_source(¤t)?;
1308 let transfer = transfer_copy_options(args, current.version_id.clone(), None, encryption)?;
1309 let copied = source_client
1310 .multipart_copy_with_transfer_options(
1311 source,
1312 target,
1313 &options,
1314 &transfer,
1315 cancellation,
1316 on_progress,
1317 )
1318 .await?;
1319 return Ok(PlannedRemoteCopyResult {
1320 bytes_copied: copied.bytes_copied,
1321 source_version_id: options.source_version_id,
1322 source_etag: current.etag.clone(),
1323 destination_version_id: copied.object.version_id.clone(),
1324 upload_id: Some(copied.upload_id),
1325 object: copied.object,
1326 });
1327 }
1328 let options = transfer_copy_options(
1329 args,
1330 source_info.version_id.clone(),
1331 source_info.etag.clone(),
1332 encryption,
1333 )?;
1334 let copied = source_client
1335 .copy_object_with_transfer_options(source, target, &options)
1336 .await?;
1337 let bytes_copied = copied
1338 .size_bytes
1339 .and_then(|size| u64::try_from(size).ok())
1340 .unwrap_or(planned_size);
1341 on_progress(bytes_copied);
1342 Ok(PlannedRemoteCopyResult {
1343 bytes_copied,
1344 source_version_id: copied
1345 .source_version_id
1346 .clone()
1347 .or_else(|| source_info.version_id.clone()),
1348 source_etag: source_info.etag.clone(),
1349 destination_version_id: copied.version_id.clone(),
1350 upload_id: None,
1351 object: copied,
1352 })
1353}
1354
1355#[allow(clippy::too_many_arguments)]
1356async fn perform_cross_alias_remote_copy(
1357 source_client: &S3Client,
1358 target_client: &S3Client,
1359 source: &RemotePath,
1360 target: &RemotePath,
1361 source_info: &ObjectInfo,
1362 encryption: Option<&ObjectEncryptionRequest>,
1363 on_progress: &(dyn Fn(u64) + Send + Sync),
1364 args: &CpArgs,
1365) -> rc_core::Result<PlannedRemoteCopyResult> {
1366 if args.source_customer_key.is_some() || args.destination_customer_key.is_some() {
1370 return Err(Error::UnsupportedFeature(
1371 "RustFS beta.10 server-side SSE-C copy is not compatibility-proven; tracked by rustfs/backlog#1467"
1372 .to_string(),
1373 ));
1374 }
1375 let planned_size = source_info
1376 .size_bytes
1377 .and_then(|size| u64::try_from(size).ok())
1378 .ok_or_else(|| Error::InvalidPath(format!("Source size is unavailable: {source}")))?;
1379 if args.storage_class.is_some() && planned_size > MULTIPART_THRESHOLD {
1380 return Err(Error::UnsupportedFeature(
1381 "RustFS beta.10 does not persist storage class for multipart uploads".to_string(),
1382 ));
1383 }
1384
1385 let current = source_client.head_object(source).await?;
1386 if !source_identity_matches(source_info, ¤t) {
1387 return Err(Error::Conflict(format!(
1388 "Source changed after copy planning: {source}"
1389 )));
1390 }
1391
1392 let staging = tempfile::Builder::new()
1393 .prefix("rc-cp-cross-alias-")
1394 .suffix(".part")
1395 .tempfile()?
1396 .into_temp_path();
1397 async {
1398 let downloaded = source_client
1399 .download_object_to_path_with_transfer_options(
1400 source,
1401 &staging,
1402 &TransferReadOptions {
1403 version_id: current.version_id.clone(),
1404 customer_key: args.source_customer_key.clone(),
1405 ..TransferReadOptions::default()
1406 },
1407 |copied, _total| on_progress(copied),
1408 )
1409 .await?;
1410 if downloaded != planned_size {
1411 return Err(Error::Conflict(format!(
1412 "Source changed after copy planning: {source}"
1413 )));
1414 }
1415 let after = source_client
1416 .head_object_with_transfer_options(
1417 source,
1418 &TransferReadOptions {
1419 version_id: current.version_id.clone(),
1420 customer_key: args.source_customer_key.clone(),
1421 ..TransferReadOptions::default()
1422 },
1423 )
1424 .await?;
1425 if !source_identity_matches(¤t, &after) {
1426 return Err(Error::Conflict(format!(
1427 "Source changed after copy planning: {source}"
1428 )));
1429 }
1430 let options = piped_copy_write_options(args, ¤t, encryption)?;
1431 let object = target_client
1432 .put_object_from_path_with_options(target, &staging, &options, |copied| {
1433 on_progress(copied)
1434 })
1435 .await?;
1436 let bytes_copied = object
1437 .size_bytes
1438 .and_then(|size| u64::try_from(size).ok())
1439 .unwrap_or(downloaded);
1440 on_progress(bytes_copied);
1441 Ok(PlannedRemoteCopyResult {
1442 bytes_copied,
1443 source_version_id: current.version_id.clone(),
1444 source_etag: current.etag.clone(),
1445 destination_version_id: object.version_id.clone(),
1446 upload_id: None,
1447 object,
1448 })
1449 }
1450 .await
1451}
1452
1453fn source_identity_matches(planned: &ObjectInfo, current: &ObjectInfo) -> bool {
1454 if planned.size_bytes != current.size_bytes {
1455 return false;
1456 }
1457 match (&planned.etag, ¤t.etag) {
1458 (Some(planned), Some(current)) => planned == current,
1459 (None, None) => {
1462 planned.version_id.is_some()
1463 && planned.version_id.as_ref() == current.version_id.as_ref()
1464 }
1465 _ => false,
1466 }
1467}
1468
1469#[derive(Debug, Clone)]
1475pub(super) struct CrossAliasCopyResult {
1476 pub(super) object: ObjectInfo,
1477 pub(super) source_version_id: Option<String>,
1478 pub(super) source_etag: Option<String>,
1479}
1480
1481pub(super) async fn copy_object_across_aliases(
1482 source_client: &S3Client,
1483 target_client: &S3Client,
1484 source: &RemotePath,
1485 target: &RemotePath,
1486 encryption: Option<&ObjectEncryptionRequest>,
1487) -> rc_core::Result<CrossAliasCopyResult> {
1488 let source_info = source_client.head_object(source).await?;
1489 let args = CpArgs::single(source.to_string(), target.to_string());
1490 let ignore_progress = |_: u64| {};
1491 let result = perform_cross_alias_remote_copy(
1492 source_client,
1493 target_client,
1494 source,
1495 target,
1496 &source_info,
1497 encryption,
1498 &ignore_progress,
1499 &args,
1500 )
1501 .await?;
1502 Ok(CrossAliasCopyResult {
1503 object: result.object,
1504 source_version_id: result.source_version_id,
1505 source_etag: result.source_etag,
1506 })
1507}
1508
1509fn piped_copy_write_options(
1510 args: &CpArgs,
1511 source: &ObjectInfo,
1512 encryption: Option<&ObjectEncryptionRequest>,
1513) -> rc_core::Result<ObjectWriteOptions> {
1514 let mut options = object_write_options(
1515 &args.fidelity,
1516 args.content_type.as_deref(),
1517 encryption,
1518 args.destination_customer_key.as_ref(),
1519 args.storage_class.clone(),
1520 )?;
1521 let replace_metadata = matches!(
1522 requested_metadata_directive(args),
1523 Some(MetadataDirective::Replace)
1524 );
1525 let mut attributes = options.attributes.take().unwrap_or_default();
1526 if !replace_metadata {
1527 if attributes.content_type.is_none() {
1528 attributes.content_type = source.content_type.clone();
1529 }
1530 if attributes.user_metadata.is_empty()
1531 && let Some(metadata) = &source.metadata
1532 {
1533 attributes.user_metadata.clone_from(metadata);
1534 }
1535 }
1536 if let Some(source_etag) = source.etag.as_deref() {
1541 set_source_identity(&mut attributes, source_etag);
1542 }
1543 if attributes != ObjectAttributes::default() {
1544 options.attributes = Some(attributes);
1545 }
1546 Ok(options)
1547}
1548
1549#[cfg(test)]
1550fn requires_multipart_copy(planned_size: Option<u64>) -> bool {
1551 planned_size.is_some_and(rc_core::requires_multipart_copy)
1552}
1553
1554fn multipart_options_from_source(source: &ObjectInfo) -> rc_core::Result<MultipartCopyOptions> {
1555 let source_size = source
1556 .size_bytes
1557 .and_then(|size| u64::try_from(size).ok())
1558 .ok_or_else(|| {
1559 Error::InvalidPath("Multipart copy source size is unavailable".to_string())
1560 })?;
1561 let source_etag = source.etag.clone().ok_or_else(|| {
1562 Error::InvalidPath("Multipart copy source ETag is unavailable".to_string())
1563 })?;
1564 let mut options = MultipartCopyOptions::new(source_size, source_etag)?;
1565 options.source_version_id = source.version_id.clone();
1566 options.content_type = source.content_type.clone();
1567 let mut attributes = ObjectAttributes {
1568 user_metadata: source.metadata.clone().unwrap_or_default(),
1569 ..ObjectAttributes::default()
1570 };
1571 if let Some(source_etag) = source.etag.as_deref() {
1572 set_source_identity(&mut attributes, source_etag);
1573 }
1574 options.metadata = attributes.user_metadata;
1575 Ok(options)
1576}
1577
1578fn operation_alias(operation: &CpOperation) -> &str {
1579 match operation {
1580 CpOperation::LocalToRemote { target, .. } => &target.alias,
1581 CpOperation::RemoteToLocal { source, .. } | CpOperation::RemoteToRemote { source, .. } => {
1582 &source.alias
1583 }
1584 }
1585}
1586
1587fn operation_client_aliases(operation: &CpOperation) -> Vec<&str> {
1588 match operation {
1589 CpOperation::LocalToRemote { target, .. } => vec![target.alias.as_str()],
1590 CpOperation::RemoteToLocal { source, .. } => vec![source.alias.as_str()],
1591 CpOperation::RemoteToRemote { source, target, .. } => {
1592 if source.alias == target.alias {
1593 vec![source.alias.as_str()]
1594 } else {
1595 vec![source.alias.as_str(), target.alias.as_str()]
1596 }
1597 }
1598 }
1599}
1600
1601fn planned_client_aliases(items: &[TransferCandidate<CpOperation>]) -> BTreeSet<String> {
1602 items
1603 .iter()
1604 .flat_map(|item| {
1605 operation_client_aliases(&item.payload)
1606 .into_iter()
1607 .map(ToOwned::to_owned)
1608 })
1609 .collect()
1610}
1611
1612async fn create_planned_client_cache(
1613 items: &[TransferCandidate<CpOperation>],
1614 alias_manager: &AliasManager,
1615) -> rc_core::Result<HashMap<String, Arc<S3Client>>> {
1616 let mut clients = HashMap::new();
1617 for alias_name in planned_client_aliases(items) {
1618 match create_leaf_s3_client(alias_manager, &alias_name).await {
1619 Ok(client) => {
1620 clients.insert(alias_name, Arc::new(client));
1621 }
1622 Err(Error::AliasNotFound(_)) => {}
1625 Err(error) => return Err(error),
1626 }
1627 }
1628 Ok(clients)
1629}
1630
1631fn planned_client<'a>(
1632 clients: &'a HashMap<String, Arc<S3Client>>,
1633 alias_name: &str,
1634) -> rc_core::Result<&'a S3Client> {
1635 clients
1636 .get(alias_name)
1637 .map(Arc::as_ref)
1638 .ok_or_else(|| Error::AliasNotFound(alias_name.to_string()))
1639}
1640
1641async fn create_leaf_s3_client(
1642 alias_manager: &AliasManager,
1643 alias_name: &str,
1644) -> rc_core::Result<S3Client> {
1645 let mut alias = alias_manager
1646 .get(alias_name)
1647 .map_err(|_| Error::AliasNotFound(alias_name.to_string()))?;
1648 alias.retry = Some(RetryConfig {
1651 max_attempts: 1,
1652 initial_backoff_ms: 1,
1653 max_backoff_ms: 1,
1654 });
1655 S3Client::new(alias).await
1656}
1657
1658fn print_planned_success(
1659 formatter: &Formatter,
1660 item: &TransferCandidate<CpOperation>,
1661 bytes_transferred: u64,
1662 detail: Option<&PlannedCopyDetail>,
1663) {
1664 let size_bytes = i64::try_from(bytes_transferred).ok();
1665 let size_human = size_bytes.map(|size| humansize::format_size(size as u64, humansize::BINARY));
1666 if formatter.is_json() {
1667 formatter.json(&CpOutput {
1668 status: "success",
1669 source: item.source.clone(),
1670 target: item.target.clone(),
1671 size_bytes,
1672 size_human,
1673 version_id: None,
1674 source_version_id: None,
1675 });
1676 } else {
1677 let mut line = format!(
1678 "{} -> {} ({})",
1679 formatter.style_file(&item.source),
1680 formatter.style_file(&item.target),
1681 formatter.style_size(size_human.as_deref().unwrap_or_default())
1682 );
1683 if let Some(detail) = detail {
1684 if let Some(version_id) = &detail.source_version_id {
1685 line.push_str(&format!(" source-version={version_id}"));
1686 }
1687 if let Some(version_id) = &detail.destination_version_id {
1688 line.push_str(&format!(" destination-version={version_id}"));
1689 }
1690 if let Some(upload_id) = &detail.upload_id {
1691 line.push_str(&format!(" upload-id={upload_id}"));
1692 }
1693 }
1694 formatter.println(&line);
1695 }
1696}
1697
1698fn print_skipped_existing(
1699 formatter: &Formatter,
1700 item: &TransferCandidate<CpOperation>,
1701 dry_run: bool,
1702) {
1703 let action = if dry_run {
1704 "Would skip existing"
1705 } else {
1706 "Skipped existing"
1707 };
1708 formatter.println(&format!(
1709 "{action}: {} -> {}",
1710 formatter.style_file(&item.source),
1711 formatter.style_file(&item.target)
1712 ));
1713}
1714
1715fn print_transfer_summary(formatter: &Formatter, summary: &rc_core::TransferSummary) {
1716 formatter.println(&format!(
1717 "Summary: {} planned, {} skipped, {} succeeded, {} failed, {} cancelled, {} transferred",
1718 summary.planned,
1719 summary.skipped,
1720 summary.successful,
1721 summary.failed,
1722 summary.cancelled,
1723 humansize::format_size(summary.transferred_bytes, humansize::BINARY)
1724 ));
1725}
1726
1727pub(super) fn exit_code_for_core_error(error: &Error) -> ExitCode {
1728 ExitCode::from_i32(error.exit_code()).unwrap_or(ExitCode::GeneralError)
1729}
1730
1731fn build_transfer_selection(args: &CpArgs, now: Timestamp) -> Result<TransferSelection, String> {
1732 let newer_than = args
1733 .newer_than
1734 .as_deref()
1735 .map(|value| parse_age_cutoff(value, now))
1736 .transpose()?;
1737 let older_than = args
1738 .older_than
1739 .as_deref()
1740 .map(|value| parse_age_cutoff(value, now))
1741 .transpose()?;
1742 let rewind = args
1743 .rewind
1744 .as_deref()
1745 .map(|value| parse_rewind_cutoff(value, now))
1746 .transpose()?;
1747
1748 TransferSelection::new(&args.include, &args.exclude, newer_than, older_than, rewind)
1749 .map_err(|error| error.to_string())
1750}
1751
1752fn build_transfer_controls(args: &CpArgs) -> Result<TransferControls, String> {
1753 let bytes_per_second = args
1754 .rate_limit
1755 .as_deref()
1756 .map(parse_byte_rate)
1757 .transpose()?;
1758 let controls = TransferControls {
1759 concurrency: args.concurrency.unwrap_or(DEFAULT_TRANSFER_CONCURRENCY),
1760 bytes_per_second,
1761 retry: RetryConfig {
1762 max_attempts: args.retry_attempts.unwrap_or(DEFAULT_RETRY_ATTEMPTS),
1763 initial_backoff_ms: args
1764 .retry_initial_backoff_ms
1765 .unwrap_or(DEFAULT_RETRY_INITIAL_BACKOFF_MS),
1766 max_backoff_ms: args
1767 .retry_max_backoff_ms
1768 .unwrap_or(DEFAULT_RETRY_MAX_BACKOFF_MS),
1769 },
1770 continue_on_error: args.continue_on_error,
1771 };
1772 controls.validate().map_err(|error| error.to_string())?;
1773 Ok(controls)
1774}
1775
1776pub(super) fn parse_age_cutoff(value: &str, now: Timestamp) -> Result<Timestamp, String> {
1777 let value = value.trim();
1778 if value.is_empty() {
1779 return Err("Transfer age must not be empty".to_string());
1780 }
1781 let suffix_index = value
1782 .find(|character: char| character.is_ascii_alphabetic())
1783 .unwrap_or(value.len());
1784 let number = &value[..suffix_index];
1785 let suffix = &value[suffix_index..];
1786 let amount: i64 = number
1787 .parse()
1788 .map_err(|_| format!("Invalid transfer age number: {number}"))?;
1789 if amount < 0 {
1790 return Err("Transfer age must not be negative".to_string());
1791 }
1792 let multiplier = match suffix.to_ascii_lowercase().as_str() {
1793 "" | "s" => 1,
1794 "m" => 60,
1795 "h" => 3_600,
1796 "d" => 86_400,
1797 "w" => 604_800,
1798 _ => return Err(format!("Unknown transfer age suffix: {suffix}")),
1799 };
1800 let seconds = amount
1801 .checked_mul(multiplier)
1802 .ok_or_else(|| "Transfer age is too large".to_string())?;
1803 now.checked_sub(jiff::Span::new().seconds(seconds))
1804 .map_err(|error| format!("Transfer age overflow: {error}"))
1805}
1806
1807fn parse_rewind_cutoff(value: &str, now: Timestamp) -> Result<Timestamp, String> {
1808 value
1809 .parse::<Timestamp>()
1810 .or_else(|_| parse_age_cutoff(value, now))
1811 .map_err(|error| format!("Invalid rewind value '{value}': {error}"))
1812}
1813
1814pub(super) fn parse_byte_rate(value: &str) -> Result<u64, String> {
1815 let normalized = value.trim().to_ascii_lowercase();
1816 let normalized = normalized
1817 .strip_suffix("/s")
1818 .or_else(|| normalized.strip_suffix("ps"))
1819 .unwrap_or(&normalized);
1820 let suffix_index = normalized
1821 .find(|character: char| character.is_ascii_alphabetic())
1822 .unwrap_or(normalized.len());
1823 let number = &normalized[..suffix_index];
1824 let suffix = &normalized[suffix_index..];
1825 let amount: u64 = number
1826 .parse()
1827 .map_err(|_| format!("Invalid transfer rate number: {number}"))?;
1828 let multiplier = match suffix {
1829 "" | "b" => 1u64,
1830 "k" | "kb" => 1_000,
1831 "ki" | "kib" => 1_024,
1832 "m" | "mb" => 1_000_000,
1833 "mi" | "mib" => 1_048_576,
1834 "g" | "gb" => 1_000_000_000,
1835 "gi" | "gib" => 1_073_741_824,
1836 _ => return Err(format!("Unknown transfer rate suffix: {suffix}")),
1837 };
1838 let rate = amount
1839 .checked_mul(multiplier)
1840 .ok_or_else(|| "Transfer rate is too large".to_string())?;
1841 if rate == 0 {
1842 return Err("Transfer rate must be greater than zero".to_string());
1843 }
1844 Ok(rate)
1845}
1846
1847fn is_container_target(raw: &str, target: &ParsedPath) -> bool {
1848 match target {
1849 ParsedPath::Local(path) => path.is_dir() || raw.ends_with(['/', '\\']),
1850 ParsedPath::Remote(path) => path.key.is_empty() || raw.ends_with('/'),
1851 }
1852}
1853
1854#[allow(clippy::too_many_arguments)]
1855async fn build_transfer_candidates(
1856 sources: &[ParsedPath],
1857 target: &ParsedPath,
1858 target_is_container: bool,
1859 recursive: bool,
1860 encryption: Option<ObjectEncryptionRequest>,
1861 source_customer_key: Option<&SseCustomerKey>,
1862 alias_manager: &AliasManager,
1863 key_policy: ObjectKeyPolicy,
1864) -> rc_core::Result<Vec<TransferCandidate<CpOperation>>> {
1865 let mut candidates = Vec::new();
1866 let mut planning_clients = HashMap::new();
1867 let multiple_sources = sources.len() > 1;
1868
1869 for source in sources {
1870 match source {
1871 ParsedPath::Local(source) => build_local_candidates(
1872 source,
1873 target,
1874 target_is_container,
1875 recursive,
1876 multiple_sources,
1877 encryption.clone(),
1878 &mut candidates,
1879 )?,
1880 ParsedPath::Remote(source) => {
1881 build_remote_candidates(
1882 source,
1883 target,
1884 target_is_container,
1885 recursive,
1886 multiple_sources,
1887 encryption.clone(),
1888 source_customer_key,
1889 alias_manager,
1890 &mut planning_clients,
1891 &mut candidates,
1892 key_policy,
1893 )
1894 .await?;
1895 }
1896 }
1897 }
1898
1899 Ok(candidates)
1900}
1901
1902fn validate_plan_targets(plan: &TransferPlan<CpOperation>) -> rc_core::Result<()> {
1903 let mut targets = HashSet::with_capacity(plan.items.len());
1904 for candidate in &plan.items {
1905 if let CpOperation::RemoteToRemote { source, target, .. } = &candidate.payload
1906 && source == target
1907 {
1908 return Err(Error::InvalidPath(format!(
1909 "Source and destination resolve to the same object '{}'",
1910 candidate.source
1911 )));
1912 }
1913 if !targets.insert(candidate.target.clone()) {
1914 return Err(Error::InvalidPath(format!(
1915 "Multiple sources resolve to the same destination '{}'",
1916 candidate.target
1917 )));
1918 }
1919 }
1920 Ok(())
1921}
1922
1923async fn skip_existing_remote_targets(
1924 plan: &mut TransferPlan<CpOperation>,
1925 clients: &HashMap<String, Arc<S3Client>>,
1926 destination_customer_key: Option<&SseCustomerKey>,
1927) -> rc_core::Result<Vec<TransferCandidate<CpOperation>>> {
1928 let mut retained = Vec::with_capacity(plan.items.len());
1929 let mut skipped = Vec::new();
1930 for candidate in plan.items.drain(..) {
1931 let destination = match &candidate.payload {
1932 CpOperation::LocalToRemote { target, .. }
1933 | CpOperation::RemoteToRemote { target, .. } => Some(target),
1934 CpOperation::RemoteToLocal { .. } => None,
1935 };
1936 let Some(destination) = destination else {
1937 retained.push(candidate);
1938 continue;
1939 };
1940 let client = planned_client(clients, &destination.alias)?;
1941 match client
1942 .head_object_with_transfer_options(
1943 destination,
1944 &TransferReadOptions {
1945 customer_key: destination_customer_key.cloned(),
1946 ..TransferReadOptions::default()
1947 },
1948 )
1949 .await
1950 {
1951 Ok(_) => skipped.push(candidate),
1952 Err(Error::NotFound(_)) | Err(Error::VersionNotFound { .. }) => {
1953 retained.push(candidate);
1954 }
1955 Err(error) => return Err(error),
1956 }
1957 }
1958 plan.items = retained;
1959 plan.summary.planned = plan.items.len();
1960 plan.summary.skipped = plan.summary.skipped.saturating_add(skipped.len());
1961 Ok(skipped)
1962}
1963
1964#[allow(clippy::too_many_arguments)]
1965fn build_local_candidates(
1966 source: &Path,
1967 target: &ParsedPath,
1968 target_is_container: bool,
1969 recursive: bool,
1970 multiple_sources: bool,
1971 encryption: Option<ObjectEncryptionRequest>,
1972 candidates: &mut Vec<TransferCandidate<CpOperation>>,
1973) -> rc_core::Result<()> {
1974 let ParsedPath::Remote(target) = target else {
1975 return Err(Error::InvalidPath(
1976 "Cannot copy between two local paths. Use the system cp command.".to_string(),
1977 ));
1978 };
1979 let metadata = std::fs::metadata(source).map_err(|error| {
1980 if error.kind() == std::io::ErrorKind::NotFound {
1981 Error::NotFound(format!("Source not found: {}", source.display()))
1982 } else {
1983 Error::Io(error)
1984 }
1985 })?;
1986
1987 if metadata.is_file() {
1988 let name = local_file_name(source)?;
1989 let destination = if target_is_container {
1990 remote_child(target, &name, ObjectKeyPolicy::for_remote_destination())?
1991 } else {
1992 normalize_remote_target(target, ObjectKeyPolicy::for_remote_destination())?
1993 };
1994 candidates.push(local_transfer_candidate(
1995 source.to_path_buf(),
1996 destination,
1997 name,
1998 metadata,
1999 encryption,
2000 ));
2001 return Ok(());
2002 }
2003
2004 if !recursive {
2005 return Err(Error::InvalidPath(
2006 "Source is a directory. Use -r/--recursive to copy directories.".to_string(),
2007 ));
2008 }
2009 if !target_is_container {
2010 return Err(Error::InvalidPath(
2011 "Recursive copy requires a directory or remote prefix destination ending in '/'"
2012 .to_string(),
2013 ));
2014 }
2015
2016 let mut files = Vec::new();
2017 collect_local_files(source, source, &mut files)?;
2018 files.sort_by(|left, right| left.1.cmp(&right.1));
2019 let source_root = multiple_sources
2020 .then(|| local_file_name(source))
2021 .transpose()?;
2022 for (path, relative, metadata) in files {
2023 let relative = relative.replace('\\', "/");
2024 let target_relative = source_root
2025 .as_deref()
2026 .map_or_else(|| relative.clone(), |root| format!("{root}/{relative}"));
2027 candidates.push(local_transfer_candidate(
2028 path,
2029 remote_child(
2030 target,
2031 &target_relative,
2032 ObjectKeyPolicy::for_remote_destination(),
2033 )?,
2034 relative,
2035 metadata,
2036 encryption.clone(),
2037 ));
2038 }
2039 Ok(())
2040}
2041
2042fn local_file_name(path: &Path) -> rc_core::Result<String> {
2043 path.file_name()
2044 .map(|name| name.to_string_lossy().to_string())
2045 .filter(|name| !name.is_empty())
2046 .ok_or_else(|| Error::InvalidPath(format!("Path has no file name: {}", path.display())))
2047}
2048
2049fn local_transfer_candidate(
2050 source: PathBuf,
2051 target: RemotePath,
2052 relative_path: String,
2053 metadata: std::fs::Metadata,
2054 encryption: Option<ObjectEncryptionRequest>,
2055) -> TransferCandidate<CpOperation> {
2056 TransferCandidate {
2057 payload: CpOperation::LocalToRemote {
2058 source: source.clone(),
2059 target: target.clone(),
2060 encryption,
2061 },
2062 source: source.display().to_string(),
2063 target: target.to_string(),
2064 relative_path,
2065 modified: metadata
2066 .modified()
2067 .ok()
2068 .and_then(|time| time.try_into().ok()),
2069 size_bytes: Some(metadata.len()),
2070 }
2071}
2072
2073fn collect_local_files(
2074 directory: &Path,
2075 base: &Path,
2076 files: &mut Vec<(PathBuf, String, std::fs::Metadata)>,
2077) -> rc_core::Result<()> {
2078 let mut entries = std::fs::read_dir(directory)?.collect::<std::io::Result<Vec<_>>>()?;
2079 entries.sort_by_key(std::fs::DirEntry::file_name);
2080 for entry in entries {
2081 let path = entry.path();
2082 let file_type = entry.file_type()?;
2083 if file_type.is_symlink() {
2084 return Err(Error::InvalidPath(format!(
2085 "Symbolic links are not supported in recursive copy: {}",
2086 path.display()
2087 )));
2088 }
2089 let metadata = entry.metadata()?;
2090 if file_type.is_file() {
2091 let relative = path
2092 .strip_prefix(base)
2093 .map_err(|error| Error::InvalidPath(error.to_string()))?
2094 .to_string_lossy()
2095 .to_string();
2096 files.push((path, relative, metadata));
2097 } else if file_type.is_dir() {
2098 collect_local_files(&path, base, files)?;
2099 }
2100 }
2101 Ok(())
2102}
2103
2104#[allow(clippy::too_many_arguments)]
2105async fn build_remote_candidates(
2106 source: &RemotePath,
2107 target: &ParsedPath,
2108 target_is_container: bool,
2109 recursive: bool,
2110 multiple_sources: bool,
2111 encryption: Option<ObjectEncryptionRequest>,
2112 source_customer_key: Option<&SseCustomerKey>,
2113 alias_manager: &AliasManager,
2114 planning_clients: &mut HashMap<String, Arc<S3Client>>,
2115 candidates: &mut Vec<TransferCandidate<CpOperation>>,
2116 key_policy: ObjectKeyPolicy,
2117) -> rc_core::Result<()> {
2118 let is_prefix = source.key.is_empty() || source.key.ends_with('/') || recursive;
2119 let client = planning_client(planning_clients, alias_manager, &source.alias).await?;
2120
2121 if is_prefix {
2122 if !recursive {
2123 return Err(Error::InvalidPath(
2124 "Remote prefix copy requires -r/--recursive".to_string(),
2125 ));
2126 }
2127 if !target_is_container {
2128 return Err(Error::InvalidPath(
2129 "Recursive copy requires a directory or remote prefix destination ending in '/'"
2130 .to_string(),
2131 ));
2132 }
2133
2134 let listing_source = recursive_listing_source(source);
2135 if let ParsedPath::Remote(target) = target
2136 && remote_copy_scopes_overlap(&listing_source, target)
2137 {
2138 return Err(Error::Conflict(format!(
2139 "Recursive source '{}' overlaps destination '{}'",
2140 listing_source, target
2141 )));
2142 }
2143
2144 let source_root = recursive_source_root(&listing_source, multiple_sources);
2145 let mut continuation_token = None;
2146 let mut seen_tokens = HashSet::new();
2147 loop {
2148 let result = client
2149 .list_objects(
2150 &listing_source,
2151 rc_core::ListOptions {
2152 recursive: true,
2153 max_keys: Some(1_000),
2154 continuation_token: continuation_token.clone(),
2155 ..Default::default()
2156 },
2157 )
2158 .await?;
2159 for object in result.items.into_iter().filter(|object| !object.is_dir) {
2160 let object_source =
2161 RemotePath::new(&listing_source.alias, &listing_source.bucket, &object.key);
2162 match target {
2163 ParsedPath::Local(target_root) => {
2164 let relative = safe_download_relative_path(
2165 &object.key,
2166 &listing_source.key,
2167 key_policy,
2168 )
2169 .map_err(Error::InvalidPath)?;
2170 let relative_string = relative.to_string_lossy().replace('\\', "/");
2171 let target_relative = if source_root.is_empty() {
2172 relative.clone()
2173 } else {
2174 PathBuf::from(&source_root).join(&relative)
2175 };
2176 let destination = safe_download_destination(target_root, &target_relative)
2177 .await
2178 .map_err(Error::InvalidPath)?;
2179 candidates.push(TransferCandidate {
2180 payload: CpOperation::RemoteToLocal {
2181 source: object_source.clone(),
2182 target: destination.clone(),
2183 },
2184 source: object_source.to_string(),
2185 target: destination.display().to_string(),
2186 relative_path: relative_string,
2187 modified: object.last_modified,
2188 size_bytes: object.size_bytes.and_then(|size| u64::try_from(size).ok()),
2189 });
2190 }
2191 ParsedPath::Remote(target) => {
2192 let (destination, relative) = recursive_remote_target(
2193 &listing_source,
2194 target,
2195 &object.key,
2196 multiple_sources,
2197 ObjectKeyPolicy::for_remote_destination(),
2198 )?;
2199 let size_bytes =
2200 object.size_bytes.and_then(|size| u64::try_from(size).ok());
2201 candidates.push(TransferCandidate {
2202 payload: CpOperation::RemoteToRemote {
2203 source: object_source.clone(),
2204 target: destination.clone(),
2205 source_info: Box::new(object.clone()),
2206 encryption: encryption.clone(),
2207 },
2208 source: object_source.to_string(),
2209 target: destination.to_string(),
2210 relative_path: relative,
2211 modified: object.last_modified,
2212 size_bytes,
2213 });
2214 }
2215 }
2216 }
2217 if !result.truncated {
2218 break;
2219 }
2220 let next_token = result.continuation_token.ok_or_else(|| {
2221 Error::InvalidPath(
2222 "Truncated object listing did not include a continuation token".to_string(),
2223 )
2224 })?;
2225 if !seen_tokens.insert(next_token.clone()) {
2226 return Err(Error::Conflict(
2227 "Object listing repeated a continuation token".to_string(),
2228 ));
2229 }
2230 continuation_token = Some(next_token);
2231 }
2232 return Ok(());
2233 }
2234
2235 let object = client
2236 .head_object_with_transfer_options(
2237 source,
2238 &TransferReadOptions {
2239 customer_key: source_customer_key.cloned(),
2240 ..TransferReadOptions::default()
2241 },
2242 )
2243 .await?;
2244 let name = source
2245 .key
2246 .rsplit('/')
2247 .next()
2248 .filter(|name| !name.is_empty())
2249 .ok_or_else(|| Error::InvalidPath(format!("Object key is empty: {source}")))?;
2250 let size_bytes = object.size_bytes.and_then(|size| u64::try_from(size).ok());
2251 let modified = object.last_modified;
2252
2253 match target {
2254 ParsedPath::Local(target) => {
2255 let destination = if target_is_container {
2256 let relative = safe_download_relative_path(name, "", key_policy)
2257 .map_err(Error::InvalidPath)?;
2258 safe_download_destination(target, &relative)
2259 .await
2260 .map_err(Error::InvalidPath)?
2261 } else {
2262 target.clone()
2263 };
2264 candidates.push(TransferCandidate {
2265 payload: CpOperation::RemoteToLocal {
2266 source: source.clone(),
2267 target: destination.clone(),
2268 },
2269 source: source.to_string(),
2270 target: destination.display().to_string(),
2271 relative_path: name.to_string(),
2272 modified,
2273 size_bytes,
2274 });
2275 }
2276 ParsedPath::Remote(target) => {
2277 let destination = if target_is_container {
2278 remote_child(target, name, ObjectKeyPolicy::for_remote_destination())?
2279 } else {
2280 normalize_remote_target(target, ObjectKeyPolicy::for_remote_destination())?
2281 };
2282 candidates.push(TransferCandidate {
2283 payload: CpOperation::RemoteToRemote {
2284 source: source.clone(),
2285 target: destination.clone(),
2286 source_info: Box::new(object),
2287 encryption,
2288 },
2289 source: source.to_string(),
2290 target: destination.to_string(),
2291 relative_path: name.to_string(),
2292 modified,
2293 size_bytes,
2294 });
2295 }
2296 }
2297 Ok(())
2298}
2299
2300async fn planning_client(
2301 clients: &mut HashMap<String, Arc<S3Client>>,
2302 alias_manager: &AliasManager,
2303 alias_name: &str,
2304) -> rc_core::Result<Arc<S3Client>> {
2305 if let Some(client) = clients.get(alias_name) {
2306 return Ok(Arc::clone(client));
2307 }
2308 let alias = alias_manager
2309 .get(alias_name)
2310 .map_err(|_| Error::AliasNotFound(alias_name.to_string()))?;
2311 let client = Arc::new(S3Client::new(alias).await?);
2312 clients.insert(alias_name.to_string(), Arc::clone(&client));
2313 Ok(client)
2314}
2315
2316fn remote_child(
2317 parent: &RemotePath,
2318 relative: &str,
2319 policy: ObjectKeyPolicy,
2320) -> rc_core::Result<RemotePath> {
2321 let parent_key = normalize_remote_prefix(&parent.key, policy)?;
2322 let relative = normalize_relative_key(relative, policy)?;
2323 let key = if parent_key.is_empty() {
2324 relative
2325 } else {
2326 format!("{parent_key}/{relative}")
2327 };
2328 Ok(RemotePath::new(&parent.alias, &parent.bucket, key))
2329}
2330
2331fn normalize_remote_target(
2332 target: &RemotePath,
2333 policy: ObjectKeyPolicy,
2334) -> rc_core::Result<RemotePath> {
2335 let key = normalize_remote_prefix(&target.key, policy)?;
2336 if target.key.ends_with('/') && !key.is_empty() {
2337 Ok(RemotePath::new(
2338 &target.alias,
2339 &target.bucket,
2340 format!("{key}/"),
2341 ))
2342 } else {
2343 Ok(RemotePath::new(&target.alias, &target.bucket, key))
2344 }
2345}
2346
2347fn normalize_remote_prefix(key: &str, policy: ObjectKeyPolicy) -> rc_core::Result<String> {
2348 if key.is_empty() {
2349 return Ok(String::new());
2350 }
2351 normalize_relative_key(key, policy)
2352}
2353
2354fn recursive_listing_source(source: &RemotePath) -> RemotePath {
2355 let key = if source.key.is_empty() || source.key.ends_with('/') {
2356 source.key.clone()
2357 } else {
2358 format!("{}/", source.key)
2359 };
2360 RemotePath::new(&source.alias, &source.bucket, key)
2361}
2362
2363fn recursive_source_root(source: &RemotePath, multiple_sources: bool) -> String {
2364 if !multiple_sources {
2365 return String::new();
2366 }
2367 source
2368 .key
2369 .trim_end_matches('/')
2370 .rsplit('/')
2371 .next()
2372 .filter(|value| !value.is_empty())
2373 .unwrap_or(&source.bucket)
2374 .to_string()
2375}
2376
2377fn recursive_remote_target(
2378 source: &RemotePath,
2379 target: &RemotePath,
2380 object_key: &str,
2381 multiple_sources: bool,
2382 policy: ObjectKeyPolicy,
2383) -> rc_core::Result<(RemotePath, String)> {
2384 let relative = object_key.strip_prefix(&source.key).ok_or_else(|| {
2385 Error::InvalidPath(format!(
2386 "Listed object key '{object_key}' is outside source prefix '{}'",
2387 source.key
2388 ))
2389 })?;
2390 if relative.is_empty() {
2391 return Err(Error::InvalidPath(format!(
2392 "Listed object key '{object_key}' does not identify a child of source prefix '{}'",
2393 source.key
2394 )));
2395 }
2396 let destination_relative = match recursive_source_root(source, multiple_sources) {
2397 root if root.is_empty() => relative.to_string(),
2398 root => format!("{root}/{relative}"),
2399 };
2400 let normalized_relative = normalize_relative_key(relative, policy)?;
2401 let normalized_destination_relative = normalize_relative_key(&destination_relative, policy)?;
2402 Ok((
2403 remote_child(target, &normalized_destination_relative, policy)?,
2404 normalized_relative,
2405 ))
2406}
2407
2408fn remote_copy_scopes_overlap(source: &RemotePath, target: &RemotePath) -> bool {
2409 if source.alias != target.alias || source.bucket != target.bucket {
2410 return false;
2411 }
2412 let source_prefix = recursive_listing_source(source).key;
2413 let target_prefix = recursive_listing_source(target).key;
2414 source_prefix.is_empty()
2415 || target_prefix.is_empty()
2416 || source_prefix.starts_with(&target_prefix)
2417 || target_prefix.starts_with(&source_prefix)
2418}
2419
2420fn parse_cp_path(path: &str, alias_manager: Option<&AliasManager>) -> rc_core::Result<ParsedPath> {
2421 let parsed = parse_path(path)?;
2422
2423 let ParsedPath::Remote(remote) = &parsed else {
2424 return Ok(parsed);
2425 };
2426
2427 if let Some(manager) = alias_manager
2428 && matches!(manager.exists(&remote.alias), Ok(true))
2429 {
2430 return Ok(parsed);
2431 }
2432
2433 if Path::new(path).exists() {
2434 return Ok(ParsedPath::Local(PathBuf::from(path)));
2435 }
2436
2437 Ok(parsed)
2438}
2439
2440async fn copy_local_to_s3_prepared(
2441 src: &Path,
2442 dst: &RemotePath,
2443 args: &CpArgs,
2444 formatter: &Formatter,
2445 encryption: Option<&ObjectEncryptionRequest>,
2446) -> ExitCode {
2447 if !src.exists() {
2449 return formatter.fail_with_suggestion(
2450 ExitCode::NotFound,
2451 &format!("Source not found: {}", src.display()),
2452 "Check the local source path and retry the copy command.",
2453 );
2454 }
2455
2456 if src.is_dir() && !args.recursive {
2458 return formatter.fail_with_suggestion(
2459 ExitCode::UsageError,
2460 "Source is a directory. Use -r/--recursive to copy directories.",
2461 "Retry with -r or --recursive to copy a directory tree.",
2462 );
2463 }
2464
2465 let alias_manager = match AliasManager::new() {
2467 Ok(am) => am,
2468 Err(e) => {
2469 formatter.error(&format!("Failed to load aliases: {e}"));
2470 return ExitCode::GeneralError;
2471 }
2472 };
2473
2474 let alias = match alias_manager.get(&dst.alias) {
2475 Ok(a) => a,
2476 Err(_) => {
2477 return formatter.fail_with_suggestion(
2478 ExitCode::NotFound,
2479 &format!("Alias '{}' not found", dst.alias),
2480 "Run `rc alias list` to inspect configured aliases or add one with `rc alias set ...`.",
2481 );
2482 }
2483 };
2484 let client = match S3Client::new(alias).await {
2485 Ok(c) => c,
2486 Err(e) => {
2487 return formatter.fail(
2488 ExitCode::NetworkError,
2489 &format!("Failed to create S3 client: {e}"),
2490 );
2491 }
2492 };
2493
2494 if src.is_file() {
2495 upload_file(&client, src, dst, args, formatter, encryption).await
2497 } else {
2498 upload_directory(&client, src, dst, args, formatter, encryption).await
2500 }
2501}
2502
2503const MULTIPART_THRESHOLD: u64 = rc_s3::multipart::DEFAULT_PART_SIZE;
2505const DOWNLOAD_PROGRESS_THRESHOLD: u64 = 4 * 1024 * 1024;
2507
2508fn update_download_progress(
2509 progress: &mut Option<ProgressBar>,
2510 output_config: &OutputConfig,
2511 bytes_downloaded: u64,
2512 total_size: Option<u64>,
2513) {
2514 let Some(total_size) = total_size else {
2515 return;
2516 };
2517
2518 if total_size < DOWNLOAD_PROGRESS_THRESHOLD {
2519 return;
2520 }
2521
2522 let progress_bar =
2523 progress.get_or_insert_with(|| ProgressBar::new(output_config.clone(), total_size));
2524 progress_bar.set_position(bytes_downloaded);
2525}
2526
2527fn print_upload_success(
2528 formatter: &Formatter,
2529 info: &rc_core::ObjectInfo,
2530 src_display: &str,
2531 dst_display: &str,
2532) {
2533 if formatter.is_json() {
2534 print_copy_json(formatter, info, src_display, dst_display);
2535 } else {
2536 let styled_src = formatter.style_file(src_display);
2537 let styled_dst = formatter.style_file(dst_display);
2538 let styled_size = formatter.style_size(&info.size_human.clone().unwrap_or_default());
2539 formatter.println(&format!("{styled_src} -> {styled_dst} ({styled_size})"));
2540 }
2541}
2542
2543async fn upload_file(
2544 client: &S3Client,
2545 src: &Path,
2546 dst: &RemotePath,
2547 args: &CpArgs,
2548 formatter: &Formatter,
2549 encryption: Option<&ObjectEncryptionRequest>,
2550) -> ExitCode {
2551 let dst_key = if dst.key.is_empty() || dst.key.ends_with('/') {
2553 let filename = src.file_name().unwrap_or_default().to_string_lossy();
2555 format!("{}{}", dst.key, filename)
2556 } else {
2557 dst.key.clone()
2558 };
2559
2560 let target = RemotePath::new(&dst.alias, &dst.bucket, &dst_key);
2561 let src_display = src.display().to_string();
2562 let dst_display = format!("{}/{}/{}", dst.alias, dst.bucket, dst_key);
2563
2564 let file_size = match std::fs::metadata(src) {
2566 Ok(m) => m.len(),
2567 Err(e) => {
2568 return formatter.fail(
2569 ExitCode::GeneralError,
2570 &format!("Failed to read {src_display}: {e}"),
2571 );
2572 }
2573 };
2574 if args.storage_class.is_some() && file_size > MULTIPART_THRESHOLD {
2575 return formatter.fail(
2576 ExitCode::UnsupportedFeature,
2577 "RustFS beta.10 does not persist storage class for multipart uploads",
2578 );
2579 }
2580 if args.dry_run {
2581 let styled_src = formatter.style_file(&src_display);
2582 let styled_dst = formatter.style_file(&dst_display);
2583 formatter.println(&format!(
2584 "Would copy: {styled_src} -> {styled_dst}{}",
2585 transfer_policy_suffix(args)
2586 ));
2587 return ExitCode::Success;
2588 }
2589
2590 let guessed_type: Option<String> = mime_guess::from_path(src)
2592 .first()
2593 .map(|m| m.essence_str().to_string());
2594 let content_type = select_upload_content_type(
2595 args.content_type.as_deref(),
2596 guessed_type.as_deref(),
2597 file_size,
2598 );
2599
2600 let progress = if file_size > MULTIPART_THRESHOLD {
2602 tracing::debug!(
2603 file_size,
2604 threshold = MULTIPART_THRESHOLD,
2605 "Using multipart upload for large file"
2606 );
2607 Some(ProgressBar::new(formatter.output_config(), file_size))
2608 } else {
2609 tracing::debug!(file_size, "Using single put_object for small file");
2610 None
2611 };
2612
2613 let upload_result = match object_write_options(
2615 &args.fidelity,
2616 content_type,
2617 encryption,
2618 args.destination_customer_key.as_ref(),
2619 args.storage_class.clone(),
2620 ) {
2621 Ok(options) => {
2622 client
2623 .put_object_from_path_with_options(&target, src, &options, |bytes_sent| {
2624 if let Some(ref pb) = progress {
2625 pb.set_position(bytes_sent);
2626 }
2627 })
2628 .await
2629 }
2630 Err(error) => Err(error),
2631 };
2632 match upload_result {
2633 Ok(info) => {
2634 if let Some(ref pb) = progress {
2635 pb.finish_and_clear();
2636 }
2637 print_upload_success(formatter, &info, &src_display, &dst_display);
2638 ExitCode::Success
2639 }
2640 Err(e) => {
2641 if let Some(ref pb) = progress {
2642 pb.finish_and_clear();
2643 }
2644 formatter.fail(
2645 exit_code_for_core_error(&e),
2646 &format!("Failed to upload {src_display}: {e}"),
2647 )
2648 }
2649 }
2650}
2651
2652fn select_upload_content_type<'a>(
2653 explicit_type: Option<&'a str>,
2654 guessed_type: Option<&'a str>,
2655 file_size: u64,
2656) -> Option<&'a str> {
2657 if file_size > MULTIPART_THRESHOLD {
2658 explicit_type
2659 } else {
2660 explicit_type.or(guessed_type)
2661 }
2662}
2663
2664async fn upload_directory(
2665 client: &S3Client,
2666 src: &Path,
2667 dst: &RemotePath,
2668 args: &CpArgs,
2669 formatter: &Formatter,
2670 encryption: Option<&ObjectEncryptionRequest>,
2671) -> ExitCode {
2672 use std::fs;
2673
2674 let mut success_count = 0;
2675 let mut error_count = 0;
2676
2677 fn walk_dir(dir: &Path, base: &Path) -> std::io::Result<Vec<(std::path::PathBuf, String)>> {
2679 let mut files = Vec::new();
2680 for entry in fs::read_dir(dir)? {
2681 let entry = entry?;
2682 let path = entry.path();
2683 if path.is_file() {
2684 let relative = path.strip_prefix(base).unwrap_or(&path);
2685 let relative_str = relative.to_string_lossy().to_string();
2686 files.push((path, relative_str));
2687 } else if path.is_dir() {
2688 files.extend(walk_dir(&path, base)?);
2689 }
2690 }
2691 Ok(files)
2692 }
2693
2694 let files = match walk_dir(src, src) {
2695 Ok(f) => f,
2696 Err(e) => {
2697 return formatter.fail(
2698 ExitCode::GeneralError,
2699 &format!("Failed to read directory: {e}"),
2700 );
2701 }
2702 };
2703
2704 for (file_path, relative_path) in files {
2705 let dst_key = if dst.key.is_empty() {
2707 relative_path.replace('\\', "/")
2708 } else if dst.key.ends_with('/') {
2709 format!("{}{}", dst.key, relative_path.replace('\\', "/"))
2710 } else {
2711 format!("{}/{}", dst.key, relative_path.replace('\\', "/"))
2712 };
2713
2714 let target = RemotePath::new(&dst.alias, &dst.bucket, &dst_key);
2715
2716 let result = upload_file(client, &file_path, &target, args, formatter, encryption).await;
2717
2718 if result == ExitCode::Success {
2719 success_count += 1;
2720 } else {
2721 error_count += 1;
2722 if !args.continue_on_error {
2723 return result;
2724 }
2725 }
2726 }
2727
2728 if error_count > 0 {
2729 formatter.warning(&format!(
2730 "Completed with errors: {success_count} succeeded, {error_count} failed"
2731 ));
2732 ExitCode::GeneralError
2733 } else {
2734 if !formatter.is_json() {
2735 formatter.success(&format!("Uploaded {success_count} file(s)."));
2736 }
2737 ExitCode::Success
2738 }
2739}
2740
2741async fn copy_s3_to_local(
2742 src: &RemotePath,
2743 dst: &Path,
2744 args: &CpArgs,
2745 formatter: &Formatter,
2746) -> ExitCode {
2747 let alias_manager = match AliasManager::new() {
2749 Ok(am) => am,
2750 Err(e) => {
2751 formatter.error(&format!("Failed to load aliases: {e}"));
2752 return ExitCode::GeneralError;
2753 }
2754 };
2755
2756 let alias = match alias_manager.get(&src.alias) {
2757 Ok(a) => a,
2758 Err(_) => {
2759 return formatter.fail_with_suggestion(
2760 ExitCode::NotFound,
2761 &format!("Alias '{}' not found", src.alias),
2762 "Run `rc alias list` to inspect configured aliases or add one with `rc alias set ...`.",
2763 );
2764 }
2765 };
2766 let client = match S3Client::new(alias).await {
2767 Ok(c) => c,
2768 Err(e) => {
2769 return formatter.fail(
2770 ExitCode::NetworkError,
2771 &format!("Failed to create S3 client: {e}"),
2772 );
2773 }
2774 };
2775
2776 let is_prefix = src.key.is_empty() || src.key.ends_with('/');
2778
2779 if is_prefix || args.recursive {
2780 download_prefix(&client, src, dst, args, formatter).await
2782 } else {
2783 download_file(&client, src, dst, args, formatter).await
2785 }
2786}
2787
2788pub(super) async fn download_file(
2789 client: &S3Client,
2790 src: &RemotePath,
2791 dst: &Path,
2792 args: &CpArgs,
2793 formatter: &Formatter,
2794) -> ExitCode {
2795 let src_display = format!("{}/{}/{}", src.alias, src.bucket, src.key);
2796
2797 let dst_path = if dst.is_dir() || dst.to_string_lossy().ends_with('/') {
2799 let filename = src.key.rsplit('/').next().unwrap_or(&src.key);
2800 let filename = match safe_download_relative_path(filename, "", args.local_key_policy()) {
2801 Ok(filename) => filename,
2802 Err(error) => {
2803 return formatter.fail(
2804 ExitCode::UsageError,
2805 &format!("Unsafe object key '{}': {error}", src.key),
2806 );
2807 }
2808 };
2809 dst.join(filename)
2810 } else {
2811 dst.to_path_buf()
2812 };
2813
2814 let dst_display = dst_path.display().to_string();
2815
2816 if args.dry_run {
2817 let styled_src = formatter.style_file(&src_display);
2818 let styled_dst = formatter.style_file(&dst_display);
2819 formatter.println(&format!("Would copy: {styled_src} -> {styled_dst}"));
2820 return ExitCode::Success;
2821 }
2822
2823 if dst_path.exists() && !args.overwrite {
2825 return formatter.fail_with_suggestion(
2826 ExitCode::Conflict,
2827 &format!("Destination exists: {dst_display}. Use --overwrite to replace."),
2828 "Retry with --overwrite if replacing the destination file is intended.",
2829 );
2830 }
2831
2832 if let Some(parent) = dst_path.parent()
2834 && !parent.exists()
2835 && let Err(e) = std::fs::create_dir_all(parent)
2836 {
2837 return formatter.fail(
2838 ExitCode::GeneralError,
2839 &format!("Failed to create directory: {e}"),
2840 );
2841 }
2842
2843 let output_config = formatter.output_config();
2844 let mut progress = None;
2845
2846 let result = client
2848 .download_object_to_path_with_transfer_options(
2849 src,
2850 &dst_path,
2851 &TransferReadOptions {
2852 customer_key: args.source_customer_key.clone(),
2853 ..TransferReadOptions::default()
2854 },
2855 |bytes_downloaded, total_size| {
2856 update_download_progress(
2857 &mut progress,
2858 &output_config,
2859 bytes_downloaded,
2860 total_size,
2861 );
2862 },
2863 )
2864 .await;
2865
2866 if let Some(ref pb) = progress {
2867 pb.finish_and_clear();
2868 }
2869
2870 match result {
2871 Ok(size) => {
2872 let size = size as i64;
2873
2874 if formatter.is_json() {
2875 let output = CpOutput {
2876 status: "success",
2877 source: src_display,
2878 target: dst_display,
2879 size_bytes: Some(size),
2880 size_human: Some(humansize::format_size(size as u64, humansize::BINARY)),
2881 version_id: None,
2882 source_version_id: None,
2883 };
2884 formatter.json(&output);
2885 } else {
2886 let styled_src = formatter.style_file(&src_display);
2887 let styled_dst = formatter.style_file(&dst_display);
2888 let styled_size =
2889 formatter.style_size(&humansize::format_size(size as u64, humansize::BINARY));
2890 formatter.println(&format!("{styled_src} -> {styled_dst} ({styled_size})"));
2891 }
2892 ExitCode::Success
2893 }
2894 Err(e) => {
2895 let err_str = e.to_string();
2896 if err_str.contains("NotFound") || err_str.contains("NoSuchKey") {
2897 formatter.fail_with_suggestion(
2898 ExitCode::NotFound,
2899 &format!("Object not found: {src_display}"),
2900 "Check the object key and bucket path, then retry the copy command.",
2901 )
2902 } else {
2903 formatter.fail(
2904 ExitCode::NetworkError,
2905 &format!("Failed to download {src_display}: {e}"),
2906 )
2907 }
2908 }
2909 }
2910}
2911
2912async fn download_prefix(
2913 client: &S3Client,
2914 src: &RemotePath,
2915 dst: &Path,
2916 args: &CpArgs,
2917 formatter: &Formatter,
2918) -> ExitCode {
2919 use rc_core::ListOptions;
2920
2921 let mut success_count = 0;
2922 let mut error_count = 0;
2923 let mut continuation_token: Option<String> = None;
2924
2925 loop {
2926 let options = ListOptions {
2927 recursive: true,
2928 max_keys: Some(1000),
2929 continuation_token: continuation_token.clone(),
2930 ..Default::default()
2931 };
2932
2933 match client.list_objects(src, options).await {
2934 Ok(result) => {
2935 for item in result.items {
2936 if item.is_dir {
2937 continue;
2938 }
2939
2940 let relative_path = match safe_download_relative_path(
2942 &item.key,
2943 &src.key,
2944 args.local_key_policy(),
2945 ) {
2946 Ok(path) => path,
2947 Err(error) => {
2948 error_count += 1;
2949 formatter.error(&format!(
2950 "Refusing unsafe object key '{}': {error}",
2951 item.key
2952 ));
2953 if !args.continue_on_error {
2954 return ExitCode::UsageError;
2955 }
2956 continue;
2957 }
2958 };
2959 let dst_path = match safe_download_destination(dst, &relative_path).await {
2960 Ok(path) => path,
2961 Err(error) => {
2962 error_count += 1;
2963 formatter.error(&format!(
2964 "Refusing unsafe destination for '{}': {error}",
2965 item.key
2966 ));
2967 if !args.continue_on_error {
2968 return ExitCode::UsageError;
2969 }
2970 continue;
2971 }
2972 };
2973
2974 let obj_src = RemotePath::new(&src.alias, &src.bucket, &item.key);
2975 let result = download_file(client, &obj_src, &dst_path, args, formatter).await;
2976
2977 if result == ExitCode::Success {
2978 success_count += 1;
2979 } else {
2980 error_count += 1;
2981 if !args.continue_on_error {
2982 return result;
2983 }
2984 }
2985 }
2986
2987 if result.truncated {
2988 continuation_token = result.continuation_token;
2989 } else {
2990 break;
2991 }
2992 }
2993 Err(e) => {
2994 return formatter.fail(
2995 ExitCode::NetworkError,
2996 &format!("Failed to list objects: {e}"),
2997 );
2998 }
2999 }
3000 }
3001
3002 if error_count > 0 {
3003 formatter.warning(&format!(
3004 "Completed with errors: {success_count} succeeded, {error_count} failed"
3005 ));
3006 ExitCode::GeneralError
3007 } else if success_count == 0 {
3008 formatter.warning("No objects found to download.");
3009 ExitCode::Success
3010 } else {
3011 if !formatter.is_json() {
3012 formatter.success(&format!("Downloaded {success_count} file(s)."));
3013 }
3014 ExitCode::Success
3015 }
3016}
3017
3018pub(super) fn safe_download_relative_path(
3019 key: &str,
3020 prefix: &str,
3021 policy: ObjectKeyPolicy,
3022) -> Result<PathBuf, String> {
3023 relative_local_path_from_key(key, prefix, policy).map_err(|error| error.to_string())
3024}
3025
3026pub(super) async fn safe_download_destination(
3027 root: &Path,
3028 relative: &Path,
3029) -> Result<PathBuf, String> {
3030 let mut destination = root.to_path_buf();
3031 for component in relative.components() {
3032 let std::path::Component::Normal(component) = component else {
3033 return Err("destination path contains a non-normal component".to_string());
3034 };
3035 destination.push(component);
3036 match tokio::fs::symlink_metadata(&destination).await {
3037 Ok(metadata) if metadata.file_type().is_symlink() => {
3038 return Err(format!(
3039 "destination component '{}' is a symbolic link",
3040 destination.display()
3041 ));
3042 }
3043 Ok(_) => {}
3044 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
3045 Err(error) => {
3046 return Err(format!(
3047 "failed to inspect destination '{}': {error}",
3048 destination.display()
3049 ));
3050 }
3051 }
3052 }
3053
3054 Ok(destination)
3055}
3056
3057async fn copy_s3_to_s3_prepared(
3058 src: &RemotePath,
3059 dst: &RemotePath,
3060 args: &CpArgs,
3061 formatter: &Formatter,
3062 encryption: Option<&ObjectEncryptionRequest>,
3063) -> ExitCode {
3064 if args.source_customer_key.is_some() || args.destination_customer_key.is_some() {
3065 return formatter.fail(
3066 ExitCode::UnsupportedFeature,
3067 "RustFS beta.10 server-side SSE-C copy is not compatibility-proven; tracked by rustfs/backlog#1467",
3068 );
3069 }
3070
3071 let alias_manager = match AliasManager::new() {
3072 Ok(am) => am,
3073 Err(e) => {
3074 formatter.error(&format!("Failed to load aliases: {e}"));
3075 return ExitCode::GeneralError;
3076 }
3077 };
3078
3079 let source_alias = match alias_manager.get(&src.alias) {
3082 Ok(a) => a,
3083 Err(_) => {
3084 return formatter.fail_with_suggestion(
3085 ExitCode::NotFound,
3086 &format!("Alias '{}' not found", src.alias),
3087 "Run `rc alias list` to inspect configured aliases or add one with `rc alias set ...`.",
3088 );
3089 }
3090 };
3091 let source_client = match S3Client::new(source_alias).await {
3092 Ok(c) => c,
3093 Err(e) => {
3094 return formatter.fail(
3095 ExitCode::NetworkError,
3096 &format!("Failed to create S3 client: {e}"),
3097 );
3098 }
3099 };
3100 let target_client;
3101 let target_client_ref = if src.alias == dst.alias {
3102 &source_client
3103 } else {
3104 let destination_alias = match alias_manager.get(&dst.alias) {
3105 Ok(a) => a,
3106 Err(_) => {
3107 return formatter.fail_with_suggestion(
3108 ExitCode::NotFound,
3109 &format!("Alias '{}' not found", dst.alias),
3110 "Run `rc alias list` to inspect configured aliases or add one with `rc alias set ...`.",
3111 );
3112 }
3113 };
3114 target_client = match S3Client::new(destination_alias).await {
3115 Ok(c) => c,
3116 Err(e) => {
3117 return formatter.fail(
3118 ExitCode::NetworkError,
3119 &format!("Failed to create destination S3 client: {e}"),
3120 );
3121 }
3122 };
3123 &target_client
3124 };
3125
3126 let src_display = format!("{}/{}/{}", src.alias, src.bucket, src.key);
3127 let dst_display = format!("{}/{}/{}", dst.alias, dst.bucket, dst.key);
3128
3129 if args.dry_run && args.storage_class.is_none() {
3130 let styled_src = formatter.style_file(&src_display);
3131 let styled_dst = formatter.style_file(&dst_display);
3132 formatter.println(&format!(
3133 "Would copy: {styled_src} -> {styled_dst}{}",
3134 transfer_policy_suffix(args)
3135 ));
3136 return ExitCode::Success;
3137 }
3138
3139 let source_info = match source_client.head_object(src).await {
3140 Ok(info) => info,
3141 Err(Error::NotFound(_)) => {
3142 return formatter.fail_with_suggestion(
3143 ExitCode::NotFound,
3144 &format!("Source not found: {src_display}"),
3145 "Check the source bucket and object key, then retry the copy command.",
3146 );
3147 }
3148 Err(error) => {
3149 return formatter.fail(
3150 exit_code_for_core_error(&error),
3151 &format!("Failed to inspect source object: {error}"),
3152 );
3153 }
3154 };
3155 let source_size = source_info
3156 .size_bytes
3157 .and_then(|size| u64::try_from(size).ok());
3158 if args.storage_class.is_some()
3159 && source_size.is_none_or(|size| {
3160 if src.alias == dst.alias {
3161 rc_core::requires_multipart_copy(size)
3162 } else {
3163 size > MULTIPART_THRESHOLD
3164 }
3165 })
3166 {
3167 return formatter.fail(
3168 ExitCode::UnsupportedFeature,
3169 "RustFS beta.10 does not persist storage class for multipart or unknown-size copies",
3170 );
3171 }
3172 if args.dry_run {
3173 let styled_src = formatter.style_file(&src_display);
3174 let styled_dst = formatter.style_file(&dst_display);
3175 formatter.println(&format!(
3176 "Would copy: {styled_src} -> {styled_dst}{}",
3177 transfer_policy_suffix(args)
3178 ));
3179 return ExitCode::Success;
3180 }
3181
3182 let cancellation = MultipartCopyCancellation::new();
3183 let transfer_cancellation = TransferCancellation::new();
3184 let signal_task = tokio::spawn({
3185 let cancellation = cancellation.clone();
3186 let transfer_cancellation = transfer_cancellation.clone();
3187 async move {
3188 if tokio::signal::ctrl_c().await.is_ok() {
3189 cancellation.cancel();
3190 transfer_cancellation.cancel();
3191 }
3192 }
3193 });
3194 let ignore_progress = |_: u64| {};
3195 let copy = perform_planned_remote_copy(
3196 &source_client,
3197 target_client_ref,
3198 src,
3199 dst,
3200 &source_info,
3201 encryption,
3202 &cancellation,
3203 &ignore_progress,
3204 args,
3205 );
3206 tokio::pin!(copy);
3207 let result = tokio::select! {
3208 biased;
3209 _ = transfer_cancellation.cancelled() => {
3210 let _ = copy.await;
3213 Err(Error::Interrupted("Copy interrupted".to_string()))
3214 }
3215 result = &mut copy => {
3216 if transfer_cancellation.is_cancelled() {
3217 Err(Error::Interrupted("Copy interrupted".to_string()))
3218 } else {
3219 result
3220 }
3221 }
3222 };
3223 signal_task.abort();
3224 let _ = signal_task.await;
3225
3226 match result {
3227 Ok(result) => {
3228 if formatter.is_json() {
3229 print_copy_json(formatter, &result.object, &src_display, &dst_display);
3230 } else {
3231 let styled_src = formatter.style_file(&src_display);
3232 let styled_dst = formatter.style_file(&dst_display);
3233 let styled_size =
3234 formatter.style_size(&result.object.size_human.unwrap_or_else(|| {
3235 humansize::format_size(result.bytes_copied, humansize::BINARY)
3236 }));
3237 let upload = result
3238 .upload_id
3239 .map(|upload_id| format!(" upload-id={upload_id}"))
3240 .unwrap_or_default();
3241 formatter.println(&format!(
3242 "{styled_src} -> {styled_dst} ({styled_size}){upload}"
3243 ));
3244 }
3245 ExitCode::Success
3246 }
3247 Err(error) => {
3248 if matches!(error, Error::NotFound(_) | Error::VersionNotFound { .. }) {
3249 formatter.fail_with_suggestion(
3250 ExitCode::NotFound,
3251 &format!("Source not found: {src_display}"),
3252 "Check the source bucket and object key, then retry the copy command.",
3253 )
3254 } else {
3255 formatter.fail(
3256 exit_code_for_core_error(&error),
3257 &format!("Failed to copy: {error}"),
3258 )
3259 }
3260 }
3261 }
3262}
3263
3264fn print_copy_json(formatter: &Formatter, info: &rc_core::ObjectInfo, source: &str, target: &str) {
3265 if info.version_id.is_some() || info.source_version_id.is_some() {
3266 formatter.json(&V3SuccessEnvelope::versioned_objects(VersionCopyData {
3267 operation: "copy",
3268 source: source.to_string(),
3269 target: target.to_string(),
3270 source_version_id: info.source_version_id.clone(),
3271 version_id: info.version_id.clone(),
3272 size_bytes: info.size_bytes,
3273 size_human: info.size_human.clone(),
3274 }));
3275 } else {
3276 formatter.json(&CpOutput {
3277 status: "success",
3278 source: source.to_string(),
3279 target: target.to_string(),
3280 size_bytes: info.size_bytes,
3281 size_human: info.size_human.clone(),
3282 version_id: None,
3283 source_version_id: None,
3284 });
3285 }
3286}
3287
3288fn parse_kms_target(value: &str) -> Result<(String, String), String> {
3289 let (target, key_id) = value
3290 .split_once('=')
3291 .ok_or_else(|| "Expected TARGET=KMS_KEY_ID for --enc-kms".to_string())?;
3292
3293 if target.is_empty() || key_id.is_empty() {
3294 return Err("Expected TARGET=KMS_KEY_ID for --enc-kms".to_string());
3295 }
3296
3297 Ok((target.to_string(), key_id.to_string()))
3298}
3299
3300pub(crate) fn parse_destination_encryption(
3301 enc_s3: &[String],
3302 enc_kms: &[String],
3303 target: &ParsedPath,
3304) -> Result<Option<ObjectEncryptionRequest>, String> {
3305 if enc_s3.is_empty() && enc_kms.is_empty() {
3306 return Ok(None);
3307 }
3308
3309 let remote = match target {
3310 ParsedPath::Remote(remote) => remote,
3311 ParsedPath::Local(_) => {
3312 return Err("Destination encryption flags must reference a remote destination".into());
3313 }
3314 };
3315
3316 let target_display = remote.to_string();
3317 let s3_matches = enc_s3.iter().any(|value| value == &target_display);
3318 let kms_targets = enc_kms
3319 .iter()
3320 .map(|value| parse_kms_target(value))
3321 .collect::<Result<Vec<_>, _>>()?;
3322 let kms_match = kms_targets
3323 .iter()
3324 .find(|(candidate, _)| candidate == &target_display);
3325
3326 if !enc_s3.is_empty() && !s3_matches {
3327 return Err(format!(
3328 "--enc-s3 target must exactly match the remote destination: {target_display}"
3329 ));
3330 }
3331
3332 if !enc_kms.is_empty() && kms_match.is_none() {
3333 return Err(format!(
3334 "--enc-kms target must exactly match the remote destination: {target_display}"
3335 ));
3336 }
3337
3338 match (s3_matches, kms_match) {
3339 (true, Some(_)) => Err(format!(
3340 "--enc-s3 and --enc-kms cannot target the same destination: {target_display}"
3341 )),
3342 (true, None) => Ok(Some(ObjectEncryptionRequest::SseS3)),
3343 (false, Some((_, key_id))) => Ok(Some(ObjectEncryptionRequest::SseKms {
3344 key_id: key_id.clone(),
3345 })),
3346 (false, None) => Ok(None),
3347 }
3348}
3349
3350#[cfg(test)]
3351mod tests {
3352 use super::*;
3353 use rc_core::{Alias, ConfigManager};
3354 use tempfile::TempDir;
3355
3356 fn temp_alias_manager() -> (AliasManager, TempDir) {
3357 let temp_dir = TempDir::new().expect("create temp dir");
3358 let config_path = temp_dir.path().join("config.toml");
3359 let config_manager = ConfigManager::with_path(config_path);
3360 let alias_manager = AliasManager::with_config_manager(config_manager);
3361 (alias_manager, temp_dir)
3362 }
3363
3364 #[test]
3365 fn test_parse_local_path() {
3366 let result = parse_path("./file.txt").unwrap();
3367 assert!(matches!(result, ParsedPath::Local(_)));
3368 }
3369
3370 #[test]
3371 fn test_parse_remote_path() {
3372 let result = parse_path("myalias/bucket/file.txt").unwrap();
3373 assert!(matches!(result, ParsedPath::Remote(_)));
3374 }
3375
3376 #[test]
3377 fn test_parse_local_absolute_path() {
3378 #[cfg(unix)]
3380 let path = "/home/user/file.txt";
3381 #[cfg(windows)]
3382 let path = "C:\\Users\\user\\file.txt";
3383
3384 let result = parse_path(path).unwrap();
3385 assert!(matches!(result, ParsedPath::Local(_)));
3386 if let ParsedPath::Local(p) = result {
3387 assert!(p.is_absolute());
3388 }
3389 }
3390
3391 #[test]
3392 fn test_parse_local_relative_path() {
3393 let result = parse_path("../file.txt").unwrap();
3394 assert!(matches!(result, ParsedPath::Local(_)));
3395 }
3396
3397 #[test]
3398 fn test_parse_remote_path_bucket_only() {
3399 let result = parse_path("myalias/bucket/").unwrap();
3400 assert!(matches!(result, ParsedPath::Remote(_)));
3401 if let ParsedPath::Remote(r) = result {
3402 assert_eq!(r.alias, "myalias");
3403 assert_eq!(r.bucket, "bucket");
3404 assert!(r.key.is_empty());
3405 }
3406 }
3407
3408 #[test]
3409 fn test_parse_remote_path_with_deep_key() {
3410 let result = parse_path("myalias/bucket/dir1/dir2/file.txt").unwrap();
3411 assert!(matches!(result, ParsedPath::Remote(_)));
3412 if let ParsedPath::Remote(r) = result {
3413 assert_eq!(r.alias, "myalias");
3414 assert_eq!(r.bucket, "bucket");
3415 assert_eq!(r.key, "dir1/dir2/file.txt");
3416 }
3417 }
3418
3419 #[test]
3420 fn test_download_progress_created_for_large_transfer() {
3421 let output_config = OutputConfig::default();
3422 let mut progress = None;
3423
3424 update_download_progress(
3425 &mut progress,
3426 &output_config,
3427 1024,
3428 Some(DOWNLOAD_PROGRESS_THRESHOLD),
3429 );
3430
3431 let progress = progress.expect("large download should create progress bar");
3432 assert!(progress.is_visible());
3433 progress.finish_and_clear();
3434 }
3435
3436 #[test]
3437 fn test_download_progress_skips_small_transfer() {
3438 let output_config = OutputConfig::default();
3439 let mut progress = None;
3440
3441 update_download_progress(
3442 &mut progress,
3443 &output_config,
3444 1024,
3445 Some(DOWNLOAD_PROGRESS_THRESHOLD - 1),
3446 );
3447
3448 assert!(progress.is_none());
3449 }
3450
3451 #[test]
3452 fn test_download_progress_skips_unknown_total_size() {
3453 let output_config = OutputConfig::default();
3454 let mut progress = None;
3455
3456 update_download_progress(&mut progress, &output_config, 1024, None);
3457
3458 assert!(progress.is_none());
3459 }
3460
3461 #[test]
3462 fn test_download_progress_respects_no_progress_config() {
3463 let output_config = OutputConfig {
3464 no_progress: true,
3465 ..Default::default()
3466 };
3467 let mut progress = None;
3468
3469 update_download_progress(
3470 &mut progress,
3471 &output_config,
3472 1024,
3473 Some(DOWNLOAD_PROGRESS_THRESHOLD),
3474 );
3475
3476 let progress = progress.expect("large download should create progress state");
3477 assert!(!progress.is_visible());
3478 }
3479
3480 #[test]
3481 fn download_relative_path_preserves_safe_nested_keys() {
3482 let relative = safe_download_relative_path(
3483 "reports/2026/july/data.csv",
3484 "reports/",
3485 ObjectKeyPolicy::Logical,
3486 )
3487 .expect("safe key should resolve");
3488
3489 assert_eq!(
3490 relative,
3491 PathBuf::from("2026").join("july").join("data.csv")
3492 );
3493 }
3494
3495 #[test]
3496 fn download_relative_path_rejects_traversal_and_absolute_keys() {
3497 for key in [
3498 "reports/../../escaped",
3499 "reports/..\\..\\escaped",
3500 "/absolute/path",
3501 ] {
3502 assert!(
3503 safe_download_relative_path(key, "reports/", ObjectKeyPolicy::Logical).is_err(),
3504 "unsafe key should be rejected: {key}"
3505 );
3506 }
3507 }
3508
3509 #[test]
3510 fn download_relative_path_accepts_colon_keys_on_logical_destinations() {
3511 let relative = safe_download_relative_path(
3512 "loki/fake/deadbeef/19f6abd9af4:19f6abe0e77:499628ff",
3513 "loki/",
3514 ObjectKeyPolicy::Logical,
3515 )
3516 .expect("colon keys are valid Unix file names");
3517
3518 assert_eq!(
3519 relative,
3520 PathBuf::from("fake")
3521 .join("deadbeef")
3522 .join("19f6abd9af4:19f6abe0e77:499628ff")
3523 );
3524 }
3525
3526 #[test]
3527 fn download_relative_path_rejects_colon_keys_when_portable_names_requested() {
3528 for key in [
3529 "reports/C:/escaped",
3530 "reports/safe:stream",
3531 "reports/CON.txt",
3532 "reports/trailing.",
3533 ] {
3534 assert!(
3535 safe_download_relative_path(key, "reports/", ObjectKeyPolicy::WindowsPortable)
3536 .is_err(),
3537 "portable-unsafe key should be rejected: {key}"
3538 );
3539 }
3540 }
3541
3542 #[cfg(unix)]
3543 #[tokio::test]
3544 async fn download_destination_rejects_existing_symlink_components() {
3545 use std::os::unix::fs::symlink;
3546
3547 let root = tempfile::tempdir().expect("create destination root");
3548 let outside = tempfile::tempdir().expect("create outside directory");
3549 symlink(outside.path(), root.path().join("linked")).expect("create test symlink");
3550
3551 let result =
3552 safe_download_destination(root.path(), &PathBuf::from("linked/file.txt")).await;
3553
3554 assert!(result.is_err());
3555 }
3556
3557 #[cfg(unix)]
3558 #[tokio::test]
3559 async fn download_destination_rejects_existing_symlink_file() {
3560 use std::os::unix::fs::symlink;
3561
3562 let root = tempfile::tempdir().expect("create destination root");
3563 let outside = tempfile::tempdir().expect("create outside directory");
3564 let outside_file = outside.path().join("file.txt");
3565 std::fs::write(&outside_file, b"outside").expect("write outside file");
3566 symlink(&outside_file, root.path().join("file.txt")).expect("create test symlink");
3567
3568 let result = safe_download_destination(root.path(), &PathBuf::from("file.txt")).await;
3569
3570 assert!(result.is_err());
3571 }
3572
3573 #[cfg(unix)]
3574 #[test]
3575 fn recursive_source_collection_rejects_symbolic_links() {
3576 use std::os::unix::fs::symlink;
3577
3578 let root = tempfile::tempdir().expect("create source root");
3579 let outside = tempfile::NamedTempFile::new().expect("create outside file");
3580 symlink(outside.path(), root.path().join("linked.txt")).expect("create source symlink");
3581 let mut files = Vec::new();
3582
3583 let result = collect_local_files(root.path(), root.path(), &mut files);
3584
3585 assert!(result.is_err());
3586 assert!(files.is_empty());
3587 }
3588
3589 #[test]
3590 fn test_select_upload_content_type_uses_guess_for_small_files() {
3591 let selected =
3592 select_upload_content_type(None, Some("text/plain"), MULTIPART_THRESHOLD - 1);
3593
3594 assert_eq!(selected, Some("text/plain"));
3595 }
3596
3597 #[test]
3598 fn test_select_upload_content_type_skips_guess_for_multipart_files() {
3599 let selected =
3600 select_upload_content_type(None, Some("text/plain"), MULTIPART_THRESHOLD + 1);
3601
3602 assert_eq!(selected, None);
3603 }
3604
3605 #[test]
3606 fn test_select_upload_content_type_uses_guess_at_multipart_boundary() {
3607 let selected = select_upload_content_type(None, Some("text/plain"), MULTIPART_THRESHOLD);
3608
3609 assert_eq!(selected, Some("text/plain"));
3610 }
3611
3612 #[test]
3613 fn test_select_upload_content_type_keeps_explicit_type_for_multipart_files() {
3614 let selected = select_upload_content_type(
3615 Some("application/octet-stream"),
3616 Some("text/plain"),
3617 MULTIPART_THRESHOLD + 1,
3618 );
3619
3620 assert_eq!(selected, Some("application/octet-stream"));
3621 }
3622
3623 #[test]
3624 fn test_parse_cp_path_prefers_existing_local_path_when_alias_missing() {
3625 let (alias_manager, temp_dir) = temp_alias_manager();
3626 let full = temp_dir.path().join("issue-2094-local").join("file.txt");
3627 let full_str = full.to_string_lossy().to_string();
3628
3629 if let Some(parent) = full.parent() {
3630 std::fs::create_dir_all(parent).expect("create parent dirs");
3631 }
3632 std::fs::write(&full, b"test").expect("write local file");
3633
3634 let parsed = parse_cp_path(&full_str, Some(&alias_manager)).expect("parse path");
3635 assert!(matches!(parsed, ParsedPath::Local(_)));
3636 }
3637
3638 #[test]
3639 fn test_parse_cp_path_keeps_remote_when_alias_exists() {
3640 let (alias_manager, _temp_dir) = temp_alias_manager();
3641 alias_manager
3642 .set(Alias::new("target", "http://localhost:9000", "a", "b"))
3643 .expect("set alias");
3644
3645 let parsed = parse_cp_path("target/bucket/file.txt", Some(&alias_manager))
3646 .expect("parse remote path");
3647 assert!(matches!(parsed, ParsedPath::Remote(_)));
3648 }
3649
3650 #[test]
3651 fn test_parse_cp_path_keeps_remote_when_local_missing() {
3652 let (alias_manager, _temp_dir) = temp_alias_manager();
3653 let parsed = parse_cp_path("missing/bucket/file.txt", Some(&alias_manager))
3654 .expect("parse remote path");
3655 assert!(matches!(parsed, ParsedPath::Remote(_)));
3656 }
3657
3658 #[test]
3659 fn test_cp_args_defaults() {
3660 let args = CpArgs::single("src", "dst");
3661 assert!(args.overwrite);
3662 assert!(!args.recursive);
3663 assert!(!args.dry_run);
3664 assert_eq!(args.concurrency, None);
3665 assert_eq!(args.retry_attempts, None);
3666 let controls = build_transfer_controls(&args).expect("default controls");
3667 assert_eq!(controls.concurrency, DEFAULT_TRANSFER_CONCURRENCY);
3668 assert_eq!(controls.retry.max_attempts, DEFAULT_RETRY_ATTEMPTS);
3669 assert_eq!(
3670 controls.retry.initial_backoff_ms,
3671 DEFAULT_RETRY_INITIAL_BACKOFF_MS
3672 );
3673 assert_eq!(controls.retry.max_backoff_ms, DEFAULT_RETRY_MAX_BACKOFF_MS);
3674 }
3675
3676 #[test]
3677 fn explicit_default_transfer_controls_enable_planner_routing() {
3678 let defaults = CpArgs::single("src", "dst");
3679 assert!(!uses_transfer_planner(&defaults));
3680
3681 let mut explicit_concurrency = defaults.clone();
3682 explicit_concurrency.concurrency = Some(DEFAULT_TRANSFER_CONCURRENCY);
3683 let mut explicit_attempts = defaults.clone();
3684 explicit_attempts.retry_attempts = Some(DEFAULT_RETRY_ATTEMPTS);
3685 let mut explicit_initial_backoff = defaults.clone();
3686 explicit_initial_backoff.retry_initial_backoff_ms = Some(DEFAULT_RETRY_INITIAL_BACKOFF_MS);
3687 let mut explicit_max_backoff = defaults;
3688 explicit_max_backoff.retry_max_backoff_ms = Some(DEFAULT_RETRY_MAX_BACKOFF_MS);
3689
3690 for args in [
3691 explicit_concurrency,
3692 explicit_attempts,
3693 explicit_initial_backoff,
3694 explicit_max_backoff,
3695 ] {
3696 assert!(uses_transfer_planner(&args));
3697 }
3698 }
3699
3700 #[test]
3701 fn planned_client_aliases_are_deduplicated() {
3702 let first = TransferCandidate {
3703 payload: CpOperation::LocalToRemote {
3704 source: PathBuf::from("first.txt"),
3705 target: RemotePath::new("shared", "bucket", "first.txt"),
3706 encryption: None,
3707 },
3708 source: "first.txt".to_string(),
3709 target: "shared/bucket/first.txt".to_string(),
3710 relative_path: "first.txt".to_string(),
3711 modified: None,
3712 size_bytes: Some(1),
3713 };
3714 let second = TransferCandidate {
3715 payload: CpOperation::RemoteToLocal {
3716 source: RemotePath::new("shared", "bucket", "second.txt"),
3717 target: PathBuf::from("second.txt"),
3718 },
3719 source: "shared/bucket/second.txt".to_string(),
3720 target: "second.txt".to_string(),
3721 relative_path: "second.txt".to_string(),
3722 modified: None,
3723 size_bytes: Some(2),
3724 };
3725
3726 assert_eq!(
3727 planned_client_aliases(&[first, second])
3728 .into_iter()
3729 .collect::<Vec<_>>(),
3730 ["shared"]
3731 );
3732 }
3733
3734 #[test]
3735 fn planned_client_aliases_include_both_sides_of_a_cross_alias_copy() {
3736 let candidate = TransferCandidate {
3737 payload: CpOperation::RemoteToRemote {
3738 source: RemotePath::new("alpha", "source", "file.txt"),
3739 target: RemotePath::new("beta", "target", "file.txt"),
3740 source_info: Box::new(ObjectInfo::file("file.txt", 4)),
3741 encryption: None,
3742 },
3743 source: "alpha/source/file.txt".to_string(),
3744 target: "beta/target/file.txt".to_string(),
3745 relative_path: "file.txt".to_string(),
3746 modified: None,
3747 size_bytes: Some(4),
3748 };
3749
3750 assert_eq!(
3751 planned_client_aliases(&[candidate])
3752 .into_iter()
3753 .collect::<Vec<_>>(),
3754 ["alpha", "beta"]
3755 );
3756 }
3757
3758 #[test]
3759 fn planned_client_aliases_keep_same_alias_remote_copy_on_one_alias() {
3760 let candidate = TransferCandidate {
3761 payload: CpOperation::RemoteToRemote {
3762 source: RemotePath::new("shared", "source", "file.txt"),
3763 target: RemotePath::new("shared", "target", "file.txt"),
3764 source_info: Box::new(ObjectInfo::file("file.txt", 4)),
3765 encryption: None,
3766 },
3767 source: "shared/source/file.txt".to_string(),
3768 target: "shared/target/file.txt".to_string(),
3769 relative_path: "file.txt".to_string(),
3770 modified: None,
3771 size_bytes: Some(4),
3772 };
3773
3774 assert_eq!(
3775 planned_client_aliases(&[candidate])
3776 .into_iter()
3777 .collect::<Vec<_>>(),
3778 ["shared"]
3779 );
3780 }
3781
3782 #[tokio::test]
3783 async fn planning_client_reuses_one_connection_pool_per_alias() {
3784 let (alias_manager, _temp_dir) = temp_alias_manager();
3785 alias_manager
3786 .set(Alias::new(
3787 "shared",
3788 "http://localhost:9000",
3789 "access",
3790 "secret",
3791 ))
3792 .expect("set alias");
3793 let mut clients = HashMap::new();
3794
3795 let first = planning_client(&mut clients, &alias_manager, "shared")
3796 .await
3797 .expect("create first client");
3798 let second = planning_client(&mut clients, &alias_manager, "shared")
3799 .await
3800 .expect("reuse first client");
3801
3802 assert_eq!(clients.len(), 1);
3803 assert!(Arc::ptr_eq(&first, &second));
3804 }
3805
3806 #[test]
3807 fn planned_remote_copy_uses_candidate_size_for_multipart_guard() {
3808 let candidate = TransferCandidate {
3809 payload: CpOperation::RemoteToRemote {
3810 source: RemotePath::new("shared", "source", "large.bin"),
3811 target: RemotePath::new("shared", "target", "large.bin"),
3812 source_info: Box::new(ObjectInfo::file(
3813 "large.bin",
3814 (MAX_SINGLE_COPY_SIZE + 1) as i64,
3815 )),
3816 encryption: None,
3817 },
3818 source: "shared/source/large.bin".to_string(),
3819 target: "shared/target/large.bin".to_string(),
3820 relative_path: "large.bin".to_string(),
3821 modified: None,
3822 size_bytes: Some(MAX_SINGLE_COPY_SIZE + 1),
3823 };
3824
3825 assert!(requires_multipart_copy(candidate.size_bytes));
3826 }
3827
3828 #[test]
3829 fn planned_remote_copy_uses_exact_five_gib_boundary() {
3830 assert!(!requires_multipart_copy(Some(MAX_SINGLE_COPY_SIZE)));
3831 assert!(requires_multipart_copy(Some(MAX_SINGLE_COPY_SIZE + 1)));
3832 }
3833
3834 #[test]
3835 fn planned_progress_replaces_attempt_bytes_instead_of_double_counting() {
3836 let progress = PlannedCopyProgress::new(
3837 OutputConfig {
3838 no_progress: true,
3839 ..OutputConfig::default()
3840 },
3841 20,
3842 );
3843 let first = ("source/a".to_string(), "target/a".to_string());
3844 let second = ("source/b".to_string(), "target/b".to_string());
3845
3846 progress.set(&first, 7);
3847 progress.set(&second, 5);
3848 progress.reset(&first);
3849 progress.set(&first, 15);
3850
3851 let positions = progress
3852 .positions
3853 .lock()
3854 .expect("planned copy progress lock");
3855 assert_eq!(positions.get(&first), Some(&15));
3856 assert_eq!(positions.get(&second), Some(&5));
3857 assert_eq!(
3858 positions.values().copied().fold(0_u64, u64::saturating_add),
3859 20
3860 );
3861 }
3862
3863 #[test]
3864 fn recursive_remote_mapping_preserves_source_relative_keys() {
3865 let source = RemotePath::new("shared", "source", "src/");
3866 let target = RemotePath::new("shared", "destination", "archive/");
3867
3868 let (destination, relative) = recursive_remote_target(
3869 &source,
3870 &target,
3871 "src/nested/report.csv",
3872 false,
3873 ObjectKeyPolicy::for_remote_destination(),
3874 )
3875 .expect("map recursive object");
3876
3877 assert_eq!(relative, "nested/report.csv");
3878 assert_eq!(destination.key, "archive/nested/report.csv");
3879 }
3880
3881 #[test]
3882 fn recursive_bucket_root_mapping_keeps_the_full_object_key() {
3883 let source = RemotePath::new("shared", "source", "");
3884 let target = RemotePath::new("shared", "destination", "archive/");
3885
3886 let (destination, relative) = recursive_remote_target(
3887 &source,
3888 &target,
3889 "nested/report.csv",
3890 false,
3891 ObjectKeyPolicy::for_remote_destination(),
3892 )
3893 .expect("map bucket object");
3894
3895 assert_eq!(relative, "nested/report.csv");
3896 assert_eq!(destination.key, "archive/nested/report.csv");
3897 }
3898
3899 #[test]
3900 fn recursive_remote_mapping_rejects_unsafe_listed_keys() {
3901 let source = RemotePath::new("shared", "source", "src/");
3902 let target = RemotePath::new("shared", "destination", "archive/");
3903
3904 for object_key in [
3905 "/absolute.txt",
3906 "src/../escape.txt",
3907 "src\\escape.txt",
3908 "src/control\u{0007}.txt",
3909 ] {
3910 assert!(
3911 recursive_remote_target(
3912 &source,
3913 &target,
3914 object_key,
3915 false,
3916 ObjectKeyPolicy::for_remote_destination(),
3917 )
3918 .is_err(),
3919 "unsafe listed key should be rejected: {object_key:?}"
3920 );
3921 }
3922 }
3923
3924 #[test]
3925 fn remote_child_rejects_unsafe_relative_keys() {
3926 let target = RemotePath::new("shared", "destination", "archive/");
3927
3928 for relative in [
3929 "../escape.txt",
3930 "nested\\escape.txt",
3931 "nested/control\u{0007}.txt",
3932 ] {
3933 assert!(
3934 remote_child(&target, relative, ObjectKeyPolicy::for_remote_destination(),)
3935 .is_err(),
3936 "unsafe relative key should be rejected: {relative:?}"
3937 );
3938 }
3939 }
3940
3941 #[test]
3942 fn remote_target_rejects_absolute_prefixes() {
3943 let target = RemotePath::new("shared", "destination", "/archive/");
3944
3945 assert!(
3946 normalize_remote_target(&target, ObjectKeyPolicy::for_remote_destination()).is_err()
3947 );
3948 }
3949
3950 #[test]
3951 fn recursive_remote_overlap_is_boundary_aware_and_symmetric() {
3952 let source = RemotePath::new("shared", "bucket", "src/");
3953 let child = RemotePath::new("shared", "bucket", "src/archive/");
3954 let sibling = RemotePath::new("shared", "bucket", "src-old/");
3955 let other_bucket = RemotePath::new("shared", "other", "src/archive/");
3956
3957 assert!(remote_copy_scopes_overlap(&source, &child));
3958 assert!(remote_copy_scopes_overlap(&child, &source));
3959 assert!(!remote_copy_scopes_overlap(&source, &sibling));
3960 assert!(!remote_copy_scopes_overlap(&source, &other_bucket));
3961 }
3962
3963 #[test]
3964 fn multipart_options_require_and_preserve_planned_source_identity() {
3965 let mut source = rc_core::ObjectInfo::file("large.bin", (MAX_SINGLE_COPY_SIZE + 1) as i64);
3966 source.etag = Some("planned-etag".to_string());
3967 source.version_id = Some("source-version".to_string());
3968 source.content_type = Some("application/octet-stream".to_string());
3969 source.metadata = Some(HashMap::from([(
3970 "project".to_string(),
3971 "archive".to_string(),
3972 )]));
3973
3974 let options = multipart_options_from_source(&source).expect("multipart options");
3975
3976 assert_eq!(options.source_size, MAX_SINGLE_COPY_SIZE + 1);
3977 assert_eq!(options.source_etag, "planned-etag");
3978 assert_eq!(options.source_version_id.as_deref(), Some("source-version"));
3979 assert_eq!(
3980 options.content_type.as_deref(),
3981 Some("application/octet-stream")
3982 );
3983 assert_eq!(
3984 options.metadata.get("project").map(String::as_str),
3985 Some("archive")
3986 );
3987 assert_eq!(
3988 options.metadata.get("rc-source-etag").map(String::as_str),
3989 Some("planned-etag")
3990 );
3991 }
3992
3993 #[test]
3994 fn transfer_age_and_rewind_use_utc_cutoffs() {
3995 let now: Timestamp = "2026-07-21T12:00:00Z".parse().expect("valid UTC timestamp");
3996
3997 assert_eq!(
3998 parse_age_cutoff("1h", now).expect("valid age").to_string(),
3999 "2026-07-21T11:00:00Z"
4000 );
4001 assert_eq!(
4002 parse_rewind_cutoff("2026-07-20T08:30:00+08:00", now)
4003 .expect("valid offset timestamp")
4004 .to_string(),
4005 "2026-07-20T00:30:00Z"
4006 );
4007 }
4008
4009 #[test]
4010 fn transfer_rate_parser_supports_decimal_and_binary_units() {
4011 assert_eq!(parse_byte_rate("10MB/s").expect("decimal rate"), 10_000_000);
4012 assert_eq!(parse_byte_rate("10MiB/s").expect("binary rate"), 10_485_760);
4013 assert!(parse_byte_rate("0").is_err());
4014 assert!(parse_byte_rate("10widgets/s").is_err());
4015 }
4016
4017 #[tokio::test]
4018 async fn planner_rejects_two_sources_that_resolve_to_one_target() {
4019 let (alias_manager, temp_dir) = temp_alias_manager();
4020 let first_dir = temp_dir.path().join("first");
4021 let second_dir = temp_dir.path().join("second");
4022 std::fs::create_dir_all(&first_dir).expect("create first source dir");
4023 std::fs::create_dir_all(&second_dir).expect("create second source dir");
4024 let first = first_dir.join("same.txt");
4025 let second = second_dir.join("same.txt");
4026 std::fs::write(&first, b"first").expect("write first source");
4027 std::fs::write(&second, b"second").expect("write second source");
4028
4029 let candidates = build_transfer_candidates(
4030 &[ParsedPath::Local(first), ParsedPath::Local(second)],
4031 &ParsedPath::Remote(RemotePath::new("target", "bucket", "prefix/")),
4032 true,
4033 false,
4034 None,
4035 None,
4036 &alias_manager,
4037 ObjectKeyPolicy::Logical,
4038 )
4039 .await
4040 .expect("sources can be expanded before selection");
4041 let plan = TransferPlan::build(candidates, &TransferSelection::default());
4042 let error = validate_plan_targets(&plan).expect_err("colliding destinations must fail");
4043
4044 assert!(error.to_string().contains("same destination"));
4045 }
4046
4047 #[test]
4048 fn parse_enc_kms_target_requires_equals_separator() {
4049 let error = parse_kms_target("local/bucket/file.txt").expect_err("missing key separator");
4050 assert!(error.contains("Expected TARGET=KMS_KEY_ID"));
4051 }
4052
4053 #[test]
4054 fn destination_encryption_rejects_local_targets() {
4055 let error = parse_destination_encryption(
4056 &[String::from("./local.txt")],
4057 &[],
4058 &ParsedPath::Local(std::path::PathBuf::from("./local.txt")),
4059 )
4060 .expect_err("local target should be rejected");
4061
4062 assert!(error.contains("must reference a remote destination"));
4063 }
4064
4065 #[test]
4066 fn destination_encryption_detects_conflicting_flags_for_same_target() {
4067 let target = ParsedPath::Remote(RemotePath::new("local", "bucket", "file.txt"));
4068 let error = parse_destination_encryption(
4069 &[String::from("local/bucket/file.txt")],
4070 &[String::from("local/bucket/file.txt=kms-key")],
4071 &target,
4072 )
4073 .expect_err("same target conflict should fail");
4074
4075 assert!(error.contains("cannot target the same destination"));
4076 }
4077
4078 #[test]
4079 fn destination_encryption_rejects_unmatched_s3_target() {
4080 let target = ParsedPath::Remote(RemotePath::new("local", "bucket", "file.txt"));
4081 let error =
4082 parse_destination_encryption(&[String::from("local/bucket/typo.txt")], &[], &target)
4083 .expect_err("unmatched s3 target should fail");
4084
4085 assert!(error.contains("must exactly match the remote destination"));
4086 }
4087
4088 #[test]
4089 fn destination_encryption_rejects_unmatched_kms_target() {
4090 let target = ParsedPath::Remote(RemotePath::new("local", "bucket", "file.txt"));
4091 let error = parse_destination_encryption(
4092 &[],
4093 &[String::from("local/bucket/typo.txt=kms-key")],
4094 &target,
4095 )
4096 .expect_err("unmatched kms target should fail");
4097
4098 assert!(error.contains("must exactly match the remote destination"));
4099 }
4100
4101 #[test]
4102 fn test_cp_output_serialization() {
4103 let output = CpOutput {
4104 status: "success",
4105 source: "src/file.txt".to_string(),
4106 target: "dst/file.txt".to_string(),
4107 size_bytes: Some(1024),
4108 size_human: Some("1 KiB".to_string()),
4109 version_id: Some("destination-v2".to_string()),
4110 source_version_id: Some("source-v1".to_string()),
4111 };
4112 let json = serde_json::to_string(&output).unwrap();
4113 assert!(json.contains("\"status\":\"success\""));
4114 assert!(json.contains("\"size_bytes\":1024"));
4115 assert!(json.contains("\"version_id\":\"destination-v2\""));
4116 assert!(json.contains("\"source_version_id\":\"source-v1\""));
4117 }
4118
4119 #[test]
4120 fn test_cp_output_skips_none_fields() {
4121 let output = CpOutput {
4122 status: "success",
4123 source: "src".to_string(),
4124 target: "dst".to_string(),
4125 size_bytes: None,
4126 size_human: None,
4127 version_id: None,
4128 source_version_id: None,
4129 };
4130 let json = serde_json::to_string(&output).unwrap();
4131 assert!(!json.contains("size_bytes"));
4132 assert!(!json.contains("size_human"));
4133 assert!(!json.contains("version_id"));
4134 }
4135
4136 #[test]
4137 fn versioned_copy_output_uses_v3_and_preserves_both_version_ids() {
4138 let envelope = V3SuccessEnvelope::versioned_objects(VersionCopyData {
4139 operation: "copy",
4140 source: "src/object.txt".to_string(),
4141 target: "dst/object.txt".to_string(),
4142 source_version_id: Some("source-v1".to_string()),
4143 version_id: Some("destination-v2".to_string()),
4144 size_bytes: Some(1024),
4145 size_human: Some("1 KiB".to_string()),
4146 });
4147
4148 let json = serde_json::to_value(envelope).expect("serialize versioned copy output");
4149 assert_eq!(json["schema_version"], 3);
4150 assert_eq!(json["type"], "versioned_objects");
4151 assert_eq!(json["data"]["operation"], "copy");
4152 assert_eq!(json["data"]["source_version_id"], "source-v1");
4153 assert_eq!(json["data"]["version_id"], "destination-v2");
4154 }
4155
4156 #[tokio::test]
4157 async fn storage_class_validation_has_usage_and_unsupported_exit_codes() {
4158 for (storage_class, expected) in [
4159 ("not-a-class", ExitCode::UsageError),
4160 ("STANDARD_IA", ExitCode::UnsupportedFeature),
4161 ] {
4162 let mut args = CpArgs::single("source.txt", "local/bucket/target.txt");
4163 args.storage_class = Some(storage_class.to_string());
4164
4165 assert_eq!(execute(args, OutputConfig::default()).await, expected);
4166 }
4167 }
4168
4169 #[test]
4170 fn fidelity_direction_preflight_rejects_silent_or_beta10_unsupported_paths() {
4171 let local = ParsedPath::Local(PathBuf::from("report.json"));
4172 let source = ParsedPath::Remote(RemotePath::new("test", "source", "report.json"));
4173 let target = ParsedPath::Remote(RemotePath::new("test", "target", "report.json"));
4174
4175 let mut preserve_upload = CpArgs::single("report.json", "test/target/report.json");
4176 preserve_upload.preserve = true;
4177 assert!(matches!(
4178 validate_fidelity_directions(&preserve_upload, std::slice::from_ref(&local), &target),
4179 Err(Error::InvalidPath(_))
4180 ));
4181
4182 let mut replace = CpArgs::single("test/source/report.json", "test/target/report.json");
4183 replace.metadata_directive = Some(MetadataDirectiveArg::Replace);
4184 assert!(matches!(
4185 validate_fidelity_directions(&replace, std::slice::from_ref(&source), &target),
4186 Err(Error::UnsupportedFeature(_))
4187 ));
4188
4189 let cross_alias_source =
4190 ParsedPath::Remote(RemotePath::new("source", "source", "report.json"));
4191 let cross_alias_target =
4192 ParsedPath::Remote(RemotePath::new("destination", "target", "report.json"));
4193 assert!(
4194 validate_fidelity_directions(
4195 &replace,
4196 std::slice::from_ref(&cross_alias_source),
4197 &cross_alias_target,
4198 )
4199 .is_ok(),
4200 "cross-alias metadata REPLACE uses the upload path"
4201 );
4202
4203 let mut tags = CpArgs::single("test/source/report.json", "test/target/report.json");
4204 tags.tagging_directive = Some(TaggingDirectiveArg::Replace);
4205 tags.fidelity.tags = vec!["env=prod".to_string()];
4206 assert!(matches!(
4207 validate_fidelity_directions(&tags, std::slice::from_ref(&source), &target),
4208 Err(Error::UnsupportedFeature(_))
4209 ));
4210
4211 let mut checksum = CpArgs::single("test/source/report.json", "test/target/report.json");
4212 checksum.fidelity.checksum = Some("sha256".to_string());
4213 assert!(matches!(
4214 validate_fidelity_directions(&checksum, std::slice::from_ref(&source), &target),
4215 Err(Error::UnsupportedFeature(_))
4216 ));
4217 }
4218
4219 #[test]
4220 fn source_identity_validation_rejects_same_size_etag_changes() {
4221 let mut planned = ObjectInfo::file("report.json", 4);
4222 planned.etag = Some("planned".to_string());
4223 let mut current = planned.clone();
4224 assert!(source_identity_matches(&planned, ¤t));
4225
4226 current.etag = Some("changed".to_string());
4227 assert!(!source_identity_matches(&planned, ¤t));
4228 current.etag = None;
4229 assert!(!source_identity_matches(&planned, ¤t));
4230
4231 planned.etag = None;
4232 assert!(!source_identity_matches(&planned, ¤t));
4233 }
4234
4235 #[test]
4236 fn preserve_builds_explicit_metadata_copy_without_replacement_payload() {
4237 let mut args = CpArgs::single("test/source/report.json", "test/target/report.json");
4238 args.preserve = true;
4239 args.fidelity.retention_mode = Some("GOVERNANCE".to_string());
4240 args.fidelity.retain_until = Some("2099-01-02T03:04:05Z".to_string());
4241 args.fidelity.legal_hold = Some("ON".to_string());
4242
4243 let options = transfer_copy_options(
4244 &args,
4245 Some("source-v1".to_string()),
4246 Some("source-etag".to_string()),
4247 None,
4248 )
4249 .expect("copy policy");
4250 assert_eq!(options.metadata_directive, Some(MetadataDirective::Copy));
4251 assert!(options.destination.attributes.is_none());
4252 assert_eq!(options.source.version_id.as_deref(), Some("source-v1"));
4253 assert_eq!(options.source_etag.as_deref(), Some("source-etag"));
4254 assert!(options.destination.retention.is_some());
4255 assert_eq!(
4256 options.destination.legal_hold,
4257 Some(rc_core::LegalHoldStatus::On)
4258 );
4259 }
4260
4261 #[test]
4262 fn storage_class_plan_rejects_multipart_without_mutation() {
4263 let plan = TransferPlan::build(
4264 vec![TransferCandidate {
4265 payload: CpOperation::RemoteToRemote {
4266 source: RemotePath::new("local", "source", "large.bin"),
4267 target: RemotePath::new("local", "target", "large.bin"),
4268 source_info: Box::new(ObjectInfo::file(
4269 "large.bin",
4270 (MAX_SINGLE_COPY_SIZE + 1) as i64,
4271 )),
4272 encryption: None,
4273 },
4274 source: "local/source/large.bin".to_string(),
4275 target: "local/target/large.bin".to_string(),
4276 relative_path: "large.bin".to_string(),
4277 modified: None,
4278 size_bytes: Some(MAX_SINGLE_COPY_SIZE + 1),
4279 }],
4280 &TransferSelection::default(),
4281 );
4282
4283 assert!(matches!(
4284 validate_storage_class_plan(&plan, Some("STANDARD")),
4285 Err(Error::UnsupportedFeature(_))
4286 ));
4287 }
4288
4289 #[test]
4290 fn storage_class_plan_rejects_cross_alias_multipart_uploads() {
4291 let plan = TransferPlan::build(
4292 vec![TransferCandidate {
4293 payload: CpOperation::RemoteToRemote {
4294 source: RemotePath::new("alpha", "source", "medium.bin"),
4295 target: RemotePath::new("beta", "target", "medium.bin"),
4296 source_info: Box::new(ObjectInfo::file(
4297 "medium.bin",
4298 (MULTIPART_THRESHOLD + 1) as i64,
4299 )),
4300 encryption: None,
4301 },
4302 source: "alpha/source/medium.bin".to_string(),
4303 target: "beta/target/medium.bin".to_string(),
4304 relative_path: "medium.bin".to_string(),
4305 modified: None,
4306 size_bytes: Some(MULTIPART_THRESHOLD + 1),
4307 }],
4308 &TransferSelection::default(),
4309 );
4310
4311 assert!(matches!(
4312 validate_storage_class_plan(&plan, Some("STANDARD")),
4313 Err(Error::UnsupportedFeature(_))
4314 ));
4315 }
4316
4317 #[test]
4318 fn piped_copy_preserves_source_content_type_unless_replaced() {
4319 let mut source = ObjectInfo::file("file.txt", 4);
4320 source.content_type = Some("text/plain".to_string());
4321 let args = CpArgs::single("alpha/source/file.txt", "beta/target/file.txt");
4322
4323 let copied = piped_copy_write_options(&args, &source, None).expect("copy metadata");
4324 assert_eq!(
4325 copied
4326 .attributes
4327 .as_ref()
4328 .and_then(|value| value.content_type.as_deref()),
4329 Some("text/plain")
4330 );
4331
4332 let mut replace = args;
4333 replace.metadata_directive = Some(MetadataDirectiveArg::Replace);
4334 let replaced = piped_copy_write_options(&replace, &source, None).expect("replace metadata");
4335 assert!(
4336 replaced
4337 .attributes
4338 .as_ref()
4339 .is_none_or(|value| value.content_type.is_none())
4340 );
4341 }
4342
4343 #[test]
4344 fn piped_copy_records_the_source_etag_as_identity_metadata() {
4345 let mut source = ObjectInfo::file("file.txt", 4);
4346 source.etag = Some("source-etag".to_string());
4347 let args = CpArgs::single("alpha/source/file.txt", "beta/target/file.txt");
4348
4349 let options = piped_copy_write_options(&args, &source, None).expect("copy metadata");
4350
4351 let attributes = options.attributes.as_ref().expect("identity attributes");
4352 assert_eq!(
4353 super::super::object_identity::identity_etag_from_metadata(Some(
4354 &attributes.user_metadata
4355 ))
4356 .as_deref(),
4357 Some("source-etag"),
4358 "a later mirror --compare auto must be able to skip this object"
4359 );
4360 }
4361
4362 #[test]
4363 fn piped_copy_records_identity_even_when_metadata_is_replaced() {
4364 let mut source = ObjectInfo::file("file.txt", 4);
4365 source.etag = Some("source-etag".to_string());
4366 source.metadata = Some(HashMap::from([(
4367 "owner".to_string(),
4368 "storage".to_string(),
4369 )]));
4370 let mut args = CpArgs::single("alpha/source/file.txt", "beta/target/file.txt");
4371 args.metadata_directive = Some(MetadataDirectiveArg::Replace);
4372
4373 let options = piped_copy_write_options(&args, &source, None).expect("replace metadata");
4374
4375 let attributes = options.attributes.as_ref().expect("identity attributes");
4376 assert_eq!(
4377 super::super::object_identity::identity_etag_from_metadata(Some(
4378 &attributes.user_metadata
4379 ))
4380 .as_deref(),
4381 Some("source-etag"),
4382 "identity is rc bookkeeping, not user metadata"
4383 );
4384 assert!(
4385 !attributes.user_metadata.contains_key("owner"),
4386 "replace must still drop source user metadata"
4387 );
4388 }
4389
4390 #[test]
4391 fn piped_copy_omits_identity_when_the_source_has_no_etag() {
4392 let source = ObjectInfo::file("file.txt", 4);
4393 let args = CpArgs::single("alpha/source/file.txt", "beta/target/file.txt");
4394
4395 let options = piped_copy_write_options(&args, &source, None).expect("copy metadata");
4396
4397 assert!(
4398 options
4399 .attributes
4400 .as_ref()
4401 .is_none_or(|attributes| attributes.user_metadata.is_empty()),
4402 "without a source ETag there is no identity to record"
4403 );
4404 }
4405
4406 #[test]
4407 fn piped_copy_preserves_source_user_metadata_unless_replaced() {
4408 let mut source = ObjectInfo::file("file.txt", 4);
4409 source.metadata = Some(HashMap::from([(
4410 "owner".to_string(),
4411 "storage".to_string(),
4412 )]));
4413 let args = CpArgs::single("alpha/source/file.txt", "beta/target/file.txt");
4414
4415 let copied = piped_copy_write_options(&args, &source, None).expect("copy metadata");
4416 assert_eq!(
4417 copied
4418 .attributes
4419 .as_ref()
4420 .and_then(|value| value.user_metadata.get("owner"))
4421 .map(String::as_str),
4422 Some("storage")
4423 );
4424
4425 let mut replace = args;
4426 replace.metadata_directive = Some(MetadataDirectiveArg::Replace);
4427 let replaced = piped_copy_write_options(&replace, &source, None).expect("replace metadata");
4428 assert!(
4429 replaced
4430 .attributes
4431 .as_ref()
4432 .is_none_or(|value| value.user_metadata.is_empty())
4433 );
4434 }
4435
4436 #[test]
4437 fn get_alias_accepts_only_one_remote_source_and_local_target() {
4438 let remote = ParsedPath::Remote(RemotePath::new("local", "reports", "report.json"));
4439 let other_remote = ParsedPath::Remote(RemotePath::new("local", "reports", "other.json"));
4440 let local = ParsedPath::Local(PathBuf::from("./report.json"));
4441
4442 assert_eq!(
4443 validate_alias_direction(TransferAlias::Get, std::slice::from_ref(&remote), &local),
4444 Ok(())
4445 );
4446 assert_eq!(
4447 validate_alias_direction(TransferAlias::Get, &[remote.clone(), other_remote], &local),
4448 Err("get requires exactly one remote source and one local target")
4449 );
4450 assert_eq!(
4451 validate_alias_direction(TransferAlias::Get, std::slice::from_ref(&local), &remote),
4452 Err("get requires exactly one remote source and one local target")
4453 );
4454 }
4455
4456 #[test]
4457 fn put_alias_accepts_only_local_sources_and_remote_target() {
4458 let first_local = ParsedPath::Local(PathBuf::from("./january.csv"));
4459 let second_local = ParsedPath::Local(PathBuf::from("./february.csv"));
4460 let remote = ParsedPath::Remote(RemotePath::new("local", "reports", ""));
4461
4462 assert_eq!(
4463 validate_alias_direction(
4464 TransferAlias::Put,
4465 &[first_local.clone(), second_local],
4466 &remote
4467 ),
4468 Ok(())
4469 );
4470 assert_eq!(
4471 validate_alias_direction(
4472 TransferAlias::Put,
4473 std::slice::from_ref(&remote),
4474 &first_local
4475 ),
4476 Err("put requires one or more local sources and one remote target")
4477 );
4478 }
4479}