Skip to main content

rustfs_cli/commands/
cp.rs

1//! cp command - Copy objects
2//!
3//! Copies objects between local filesystem and S3, or between S3 locations.
4
5use 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/// Copy objects
57#[derive(Args, Clone)]
58#[command(after_help = CP_AFTER_HELP)]
59pub struct CpArgs {
60    /// Source paths (local paths or alias/bucket/key)
61    #[arg(required = true, num_args = 1.., value_name = "SOURCE")]
62    pub sources: Vec<String>,
63
64    /// Destination path (local path or alias/bucket/key)
65    pub target: String,
66
67    /// Copy recursively
68    #[arg(short, long)]
69    pub recursive: bool,
70
71    /// Preserve file attributes
72    #[arg(short, long, conflicts_with = "metadata_directive")]
73    pub preserve: bool,
74
75    /// Source metadata handling for remote copies
76    #[arg(long, value_enum)]
77    pub(crate) metadata_directive: Option<MetadataDirectiveArg>,
78
79    /// Source tag handling for remote copies
80    #[arg(long, value_enum)]
81    pub(crate) tagging_directive: Option<TaggingDirectiveArg>,
82
83    /// Continue on errors
84    #[arg(long)]
85    pub continue_on_error: bool,
86
87    /// Overwrite destination if it exists
88    #[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    /// Skip remote destinations that already exist
98    #[arg(long)]
99    pub skip_existing: bool,
100
101    /// Only show what would be copied (dry run)
102    #[arg(long)]
103    pub dry_run: bool,
104
105    /// Storage class for destination (S3 only)
106    #[arg(long)]
107    pub storage_class: Option<String>,
108
109    /// Content type for uploaded files
110    #[arg(long)]
111    pub content_type: Option<String>,
112
113    #[command(flatten)]
114    pub(crate) fidelity: TransferFidelityArgs,
115
116    /// Apply SSE-S3 to the remote destination path
117    #[arg(long = "enc-s3")]
118    pub enc_s3: Vec<String>,
119
120    /// Apply SSE-KMS to the remote destination path as TARGET=KMS_KEY_ID
121    #[arg(long = "enc-kms")]
122    pub enc_kms: Vec<String>,
123
124    /// Read a 32-byte SSE-C source key from a protected file
125    #[arg(long = "enc-c-source-key-file")]
126    pub enc_c_source_key_file: Option<PathBuf>,
127
128    /// Read a 32-byte SSE-C source key from the named environment variable
129    #[arg(long = "enc-c-source-key-env")]
130    pub enc_c_source_key_env: Option<String>,
131
132    /// Read a 32-byte SSE-C destination key from a protected file
133    #[arg(long = "enc-c-destination-key-file")]
134    pub enc_c_destination_key_file: Option<PathBuf>,
135
136    /// Read a 32-byte SSE-C destination key from the named environment variable
137    #[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    /// Include source-relative paths matching this glob (repeatable)
147    #[arg(long)]
148    pub include: Vec<String>,
149
150    /// Exclude source-relative paths matching this glob; exclusions always win (repeatable)
151    #[arg(long)]
152    pub exclude: Vec<String>,
153
154    /// Select objects modified more recently than this age (for example 1h or 7d)
155    #[arg(long)]
156    pub newer_than: Option<String>,
157
158    /// Select objects modified less recently than this age (for example 1h or 7d)
159    #[arg(long)]
160    pub older_than: Option<String>,
161
162    /// Select object state at or before a UTC timestamp or age
163    #[arg(long)]
164    pub rewind: Option<String>,
165
166    /// Maximum number of transfers in flight across the command
167    #[arg(long)]
168    pub concurrency: Option<usize>,
169
170    /// Aggregate transfer start rate in bytes per second (for example 10MiB/s)
171    #[arg(long)]
172    pub rate_limit: Option<String>,
173
174    /// Maximum attempts for transient transfer failures
175    #[arg(long)]
176    pub retry_attempts: Option<u32>,
177
178    /// Initial transient retry backoff in milliseconds
179    #[arg(long)]
180    pub retry_initial_backoff_ms: Option<u64>,
181
182    /// Maximum transient retry backoff in milliseconds
183    #[arg(long)]
184    pub retry_max_backoff_ms: Option<u64>,
185
186    /// Return not-found when selection produces no transfer candidates
187    #[arg(long)]
188    pub fail_empty: bool,
189
190    /// Print deterministic aggregate transfer counters (human output)
191    #[arg(long)]
192    pub summary: bool,
193
194    /// Reject object keys that cannot be created on Windows filesystems
195    #[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/// Download one remote object through the canonical copy implementation.
251#[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/// Upload one or more local paths through the canonical copy implementation.
262#[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
302/// Execute the cp command
303pub async fn execute(args: CpArgs, output_config: OutputConfig) -> ExitCode {
304    execute_with_alias(args, output_config, TransferAlias::Copy).await
305}
306
307/// Execute the mc-compatible `get` command through the canonical copy path.
308pub async fn execute_get(args: GetArgs, output_config: OutputConfig) -> ExitCode {
309    execute_with_alias(args.transfer, output_config, TransferAlias::Get).await
310}
311
312/// Execute the mc-compatible `put` command through the canonical copy path.
313pub 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    // Parse every operand before starting any transfer.
363    let mut sources = Vec::with_capacity(args.sources.len());
364    for source in &args.sources {
365        let parsed = match alias {
366            // `put` defines every source operand as a local path, including
367            // relative paths containing a slash that do not exist yet.
368            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        // `get` defines its target operand as a local path. Do not reinterpret
387        // a non-existent `directory/file` path as an alias and bucket.
388        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    // Determine copy direction.
742    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(&copy_details);
984            let copy_progress = Arc::clone(&copy_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(&copy_details);
990                let copy_progress = Arc::clone(&copy_progress);
991                async move {
992                    execute_planned_operation(
993                        item,
994                        &operation_args,
995                        &clients,
996                        &multipart_cancellation,
997                        &copy_details,
998                        &copy_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, &current) {
1303            return Err(Error::Conflict(format!(
1304                "Source changed after copy planning: {source}"
1305            )));
1306        }
1307        let options = multipart_options_from_source(&current)?;
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    // Server-side CopyObject cannot target a different alias/endpoint. Stream
1367    // through a bounded temporary file so the destination write is a normal
1368    // upload with the destination alias credentials.
1369    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, &current) {
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(&current, &after) {
1426            return Err(Error::Conflict(format!(
1427                "Source changed after copy planning: {source}"
1428            )));
1429        }
1430        let options = piped_copy_write_options(args, &current, 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, &current.etag) {
1458        (Some(planned), Some(current)) => planned == current,
1459        // Without an ETag an unversioned object has no stable read identity;
1460        // only an identical explicit version can prove that it is unchanged.
1461        (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/// Copy one object between two aliases by streaming through the client.
1470///
1471/// `rc mv` shares this path so a cross-alias move behaves exactly like a
1472/// cross-alias copy followed by a source delete, rather than reimplementing the
1473/// download/upload streaming and its source-change checks.
1474#[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    // The destination computes its own ETag, so record the source ETag the same
1537    // way `rc mirror` does. Without this a later `mirror --compare auto` cannot
1538    // tell a faithful cross-alias copy from a changed object and recopies it.
1539    // This is `rc` bookkeeping rather than user data, so it survives --metadata-directive replace.
1540    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            // Keep a missing alias as an item-level failure so continue-on-error and summaries
1623            // retain their existing behavior without reloading configuration in each worker.
1624            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    // The shared executor classifies failures and owns the exact attempt budget. Disabling adapter
1649    // retries here prevents nested retries from exceeding the command-level policy.
1650    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    // Check if source exists
2448    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 source is a directory, require recursive flag
2457    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    // Load alias and create client
2466    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        // Single file upload
2496        upload_file(&client, src, dst, args, formatter, encryption).await
2497    } else {
2498        // Directory upload
2499        upload_directory(&client, src, dst, args, formatter, encryption).await
2500    }
2501}
2502
2503/// Multipart upload threshold: files larger than this size use multipart upload.
2504const MULTIPART_THRESHOLD: u64 = rc_s3::multipart::DEFAULT_PART_SIZE;
2505/// Download progress threshold: avoid flicker for tiny downloads while surfacing meaningful waits.
2506const 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    // Determine destination key
2552    let dst_key = if dst.key.is_empty() || dst.key.ends_with('/') {
2553        // If destination is a directory, use source filename
2554        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    // Get file size for progress bar decision
2565    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    // Determine content type
2591    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    // Show progress bar for large files
2601    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    // Upload
2614    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    // Walk directory
2678    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        // Build destination key
2706        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    // Load alias and create client
2748    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    // Check if source is a prefix (directory-like)
2777    let is_prefix = src.key.is_empty() || src.key.ends_with('/');
2778
2779    if is_prefix || args.recursive {
2780        // Download multiple objects
2781        download_prefix(&client, src, dst, args, formatter).await
2782    } else {
2783        // Download single object
2784        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    // Determine destination path
2798    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    // Check if destination exists
2824    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    // Create parent directories
2833    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    // Download object
2847    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                    // Calculate relative path from prefix
2941                    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    // Same-alias copies use server-side CopyObject. Different aliases download
3080    // through a temporary file and upload with the destination credentials.
3081    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            // Do not drop an in-flight request. Multipart copy observes its own token and performs
3211            // cleanup; CopyObject is allowed to settle before the command reports interruption.
3212            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        // Use platform-appropriate absolute path
3379        #[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, &current));
4225
4226        current.etag = Some("changed".to_string());
4227        assert!(!source_identity_matches(&planned, &current));
4228        current.etag = None;
4229        assert!(!source_identity_matches(&planned, &current));
4230
4231        planned.etag = None;
4232        assert!(!source_identity_matches(&planned, &current));
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}