diff --git a/crates/cli/src/commands/cat.rs b/crates/cli/src/commands/cat.rs index 9fa3703..5cd1c13 100644 --- a/crates/cli/src/commands/cat.rs +++ b/crates/cli/src/commands/cat.rs @@ -3,9 +3,10 @@ //! Outputs the entire content of an object to stdout. use clap::Args; -use rc_core::{AliasManager, RemotePath}; +use rc_core::{AliasManager, ObjectReadOptions, ObjectStore, RemotePath}; use rc_s3::S3Client; +use crate::commands::{exit_code_for_core_error, validate_version_selector}; use crate::exit_code::ExitCode; use crate::output::{Formatter, OutputConfig}; @@ -32,12 +33,27 @@ pub struct CatArgs { pub async fn execute(args: CatArgs, output_config: OutputConfig) -> ExitCode { let formatter = Formatter::new(output_config); - if args.enc_key.is_some() || args.rewind.is_some() || args.version_id.is_some() { + if let Err(error) = + validate_version_selector(args.version_id.as_deref(), args.rewind.as_deref()) + { + return formatter.fail(ExitCode::UsageError, &error); + } + if args.enc_key.is_some() { + return formatter.fail( + ExitCode::UnsupportedFeature, + "--enc-key is not implemented for cat", + ); + } + if args.rewind.is_some() { return formatter.fail( ExitCode::UnsupportedFeature, - "--enc-key, --rewind, and --version-id are not implemented for cat", + "--rewind is not implemented for cat", ); } + let read_options = match ObjectReadOptions::for_version(args.version_id.clone()) { + Ok(options) => options, + Err(error) => return formatter.fail(ExitCode::UsageError, &error.to_string()), + }; // Parse the path let (alias_name, bucket, key) = match parse_cat_path(&args.path) { @@ -78,20 +94,20 @@ pub async fn execute(args: CatArgs, output_config: OutputConfig) -> ExitCode { // Get object content let mut stdout = tokio::io::stdout(); - match client.write_object_to(&path, &mut stdout, None).await { + match ObjectStore::write_object_to_with_options( + &client, + &path, + &read_options, + &mut stdout, + None, + ) + .await + { Ok(_) => ExitCode::Success, Err(e) => { - let err_str = e.to_string(); - if err_str.contains("NotFound") || err_str.contains("NoSuchKey") { - formatter.error(&format!("Object not found: {}", args.path)); - ExitCode::NotFound - } else if err_str.contains("AccessDenied") { - formatter.error(&format!("Access denied: {}", args.path)); - ExitCode::AuthError - } else { - formatter.error(&format!("Failed to get object: {e}")); - ExitCode::NetworkError - } + let exit_code = exit_code_for_core_error(&e); + formatter.error_with_code(exit_code, &format!("Failed to read {}: {e}", args.path)); + exit_code } } } diff --git a/crates/cli/src/commands/cp.rs b/crates/cli/src/commands/cp.rs index 0bf8339..fb357c7 100644 --- a/crates/cli/src/commands/cp.rs +++ b/crates/cli/src/commands/cp.rs @@ -11,7 +11,7 @@ use serde::Serialize; use std::path::{Path, PathBuf}; use crate::exit_code::ExitCode; -use crate::output::{Formatter, OutputConfig, ProgressBar}; +use crate::output::{Formatter, OutputConfig, ProgressBar, V3SuccessEnvelope}; const CP_AFTER_HELP: &str = "\ Examples: @@ -78,6 +78,21 @@ struct CpOutput { size_bytes: Option, #[serde(skip_serializing_if = "Option::is_none")] size_human: Option, + #[serde(skip_serializing_if = "Option::is_none")] + version_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + source_version_id: Option, +} + +#[derive(Debug, Serialize)] +struct VersionCopyData { + operation: &'static str, + source: String, + target: String, + source_version_id: Option, + version_id: Option, + size_bytes: Option, + size_human: Option, } /// Execute the cp command @@ -272,14 +287,7 @@ fn print_upload_success( dst_display: &str, ) { if formatter.is_json() { - let output = CpOutput { - status: "success", - source: src_display.to_string(), - target: dst_display.to_string(), - size_bytes: info.size_bytes, - size_human: info.size_human.clone(), - }; - formatter.json(&output); + print_copy_json(formatter, info, src_display, dst_display); } else { let styled_src = formatter.style_file(src_display); let styled_dst = formatter.style_file(dst_display); @@ -595,6 +603,8 @@ pub(super) async fn download_file( target: dst_display, size_bytes: Some(size), size_human: Some(humansize::format_size(size as u64, humansize::BINARY)), + version_id: None, + source_version_id: None, }; formatter.json(&output); } else { @@ -889,14 +899,7 @@ async fn copy_s3_to_s3( match client.copy_object(src, dst, encryption.as_ref()).await { Ok(info) => { if formatter.is_json() { - let output = CpOutput { - status: "success", - source: src_display, - target: dst_display, - size_bytes: info.size_bytes, - size_human: info.size_human, - }; - formatter.json(&output); + print_copy_json(formatter, &info, &src_display, &dst_display); } else { let styled_src = formatter.style_file(&src_display); let styled_dst = formatter.style_file(&dst_display); @@ -920,6 +923,30 @@ async fn copy_s3_to_s3( } } +fn print_copy_json(formatter: &Formatter, info: &rc_core::ObjectInfo, source: &str, target: &str) { + if info.version_id.is_some() || info.source_version_id.is_some() { + formatter.json(&V3SuccessEnvelope::versioned_objects(VersionCopyData { + operation: "copy", + source: source.to_string(), + target: target.to_string(), + source_version_id: info.source_version_id.clone(), + version_id: info.version_id.clone(), + size_bytes: info.size_bytes, + size_human: info.size_human.clone(), + })); + } else { + formatter.json(&CpOutput { + status: "success", + source: source.to_string(), + target: target.to_string(), + size_bytes: info.size_bytes, + size_human: info.size_human.clone(), + version_id: None, + source_version_id: None, + }); + } +} + fn parse_kms_target(value: &str) -> Result<(String, String), String> { let (target, key_id) = value .split_once('=') @@ -1323,10 +1350,14 @@ mod tests { target: "dst/file.txt".to_string(), size_bytes: Some(1024), size_human: Some("1 KiB".to_string()), + version_id: Some("destination-v2".to_string()), + source_version_id: Some("source-v1".to_string()), }; let json = serde_json::to_string(&output).unwrap(); assert!(json.contains("\"status\":\"success\"")); assert!(json.contains("\"size_bytes\":1024")); + assert!(json.contains("\"version_id\":\"destination-v2\"")); + assert!(json.contains("\"source_version_id\":\"source-v1\"")); } #[test] @@ -1337,9 +1368,32 @@ mod tests { target: "dst".to_string(), size_bytes: None, size_human: None, + version_id: None, + source_version_id: None, }; let json = serde_json::to_string(&output).unwrap(); assert!(!json.contains("size_bytes")); assert!(!json.contains("size_human")); + assert!(!json.contains("version_id")); + } + + #[test] + fn versioned_copy_output_uses_v3_and_preserves_both_version_ids() { + let envelope = V3SuccessEnvelope::versioned_objects(VersionCopyData { + operation: "copy", + source: "src/object.txt".to_string(), + target: "dst/object.txt".to_string(), + source_version_id: Some("source-v1".to_string()), + version_id: Some("destination-v2".to_string()), + size_bytes: Some(1024), + size_human: Some("1 KiB".to_string()), + }); + + let json = serde_json::to_value(envelope).expect("serialize versioned copy output"); + assert_eq!(json["schema_version"], 3); + assert_eq!(json["type"], "versioned_objects"); + assert_eq!(json["data"]["operation"], "copy"); + assert_eq!(json["data"]["source_version_id"], "source-v1"); + assert_eq!(json["data"]["version_id"], "destination-v2"); } } diff --git a/crates/cli/src/commands/head.rs b/crates/cli/src/commands/head.rs index 7bb1d34..8ba8f81 100644 --- a/crates/cli/src/commands/head.rs +++ b/crates/cli/src/commands/head.rs @@ -3,10 +3,11 @@ //! Outputs the first N lines (or bytes) of an object to stdout. use clap::Args; -use rc_core::{AliasManager, ObjectStore as _, RemotePath}; +use rc_core::{AliasManager, ObjectReadOptions, ObjectStore, RemotePath}; use rc_s3::S3Client; use std::io::{self, Write}; +use crate::commands::{exit_code_for_core_error, validate_version_selector}; use crate::exit_code::ExitCode; use crate::output::{Formatter, OutputConfig}; @@ -33,12 +34,13 @@ pub struct HeadArgs { pub async fn execute(args: HeadArgs, output_config: OutputConfig) -> ExitCode { let formatter = Formatter::new(output_config); - if args.version_id.is_some() { - return formatter.fail( - ExitCode::UnsupportedFeature, - "--version-id is not implemented for head", - ); + if let Err(error) = validate_version_selector(args.version_id.as_deref(), None) { + return formatter.fail(ExitCode::UsageError, &error); } + let read_options = match ObjectReadOptions::for_version(args.version_id.clone()) { + Ok(options) => options, + Err(error) => return formatter.fail(ExitCode::UsageError, &error.to_string()), + }; // Parse the path let (alias_name, bucket, key) = match parse_head_path(&args.path) { @@ -79,9 +81,14 @@ pub async fn execute(args: HeadArgs, output_config: OutputConfig) -> ExitCode { if let Some(num_bytes) = args.bytes { let mut stdout = tokio::io::stdout(); - return match client - .write_object_to(&path, &mut stdout, Some(num_bytes as u64)) - .await + return match ObjectStore::write_object_to_with_options( + &client, + &path, + &read_options, + &mut stdout, + Some(num_bytes as u64), + ) + .await { Ok(_) => ExitCode::Success, Err(e) => output_get_error(&formatter, &args.path, e), @@ -89,7 +96,7 @@ pub async fn execute(args: HeadArgs, output_config: OutputConfig) -> ExitCode { } // Get object content - match client.get_object(&path).await { + match ObjectStore::get_object_with_options(&client, &path, &read_options).await { Ok(data) => { let content = String::from_utf8_lossy(&data); let lines: Vec<&str> = content.lines().take(args.lines).collect(); @@ -105,17 +112,9 @@ pub async fn execute(args: HeadArgs, output_config: OutputConfig) -> ExitCode { } fn output_get_error(formatter: &Formatter, path: &str, error: rc_core::Error) -> ExitCode { - let message = error.to_string(); - if message.contains("NotFound") || message.contains("NoSuchKey") { - formatter.error(&format!("Object not found: {path}")); - ExitCode::NotFound - } else if message.contains("AccessDenied") { - formatter.error(&format!("Access denied: {path}")); - ExitCode::AuthError - } else { - formatter.error(&format!("Failed to get object: {error}")); - ExitCode::NetworkError - } + let exit_code = exit_code_for_core_error(&error); + formatter.error_with_code(exit_code, &format!("Failed to read {path}: {error}")); + exit_code } /// Parse head path into (alias, bucket, key) diff --git a/crates/cli/src/commands/mirror.rs b/crates/cli/src/commands/mirror.rs index 167b0bb..d083fc3 100644 --- a/crates/cli/src/commands/mirror.rs +++ b/crates/cli/src/commands/mirror.rs @@ -753,6 +753,9 @@ mod tests { storage_class: None, content_type: self.content_type.clone(), metadata: None, + version_id: None, + source_version_id: None, + is_delete_marker: None, is_dir: false, }) } diff --git a/crates/cli/src/commands/mod.rs b/crates/cli/src/commands/mod.rs index 34606ca..d5e19e0 100644 --- a/crates/cli/src/commands/mod.rs +++ b/crates/cli/src/commands/mod.rs @@ -44,6 +44,20 @@ mod tag; mod tree; mod version; +fn exit_code_for_core_error(error: &rc_core::Error) -> ExitCode { + ExitCode::from_i32(error.exit_code()).unwrap_or(ExitCode::GeneralError) +} + +fn validate_version_selector(version_id: Option<&str>, rewind: Option<&str>) -> Result<(), String> { + if version_id.is_some() && rewind.is_some() { + return Err("--version-id cannot be combined with --rewind".to_string()); + } + if version_id.is_some_and(str::is_empty) { + return Err("--version-id cannot be empty".to_string()); + } + Ok(()) +} + /// rc - Rust S3 CLI Client /// /// A command-line interface for S3-compatible object storage services. @@ -86,7 +100,17 @@ pub struct Cli { } fn parse_request_header(value: &str) -> Result { - RequestHeader::parse(value).map_err(|error| error.to_string()) + let header = RequestHeader::parse(value).map_err(|error| error.to_string())?; + if header + .name + .eq_ignore_ascii_case("x-amz-bypass-governance-retention") + { + return Err( + "Use the remove command's explicit --bypass flag for governance retention bypass" + .to_string(), + ); + } + Ok(header) } #[derive(Copy, Clone, Debug, Eq, PartialEq, ValueEnum)] @@ -593,6 +617,20 @@ mod tests { ); } + #[test] + fn cli_requires_bypass_flag_instead_of_custom_governance_header() { + let error = Cli::try_parse_from([ + "rc", + "-H", + "x-amz-bypass-governance-retention:true", + "rm", + "local/bucket/key.txt", + ]) + .expect_err("governance bypass header should require --bypass"); + + assert!(error.to_string().contains("explicit --bypass flag")); + } + #[test] fn cli_accepts_bucket_cors_subcommand() { let cli = Cli::try_parse_from(["rc", "bucket", "cors", "list", "local/my-bucket"]) @@ -1242,4 +1280,34 @@ mod tests { other => panic!("expected object command, got {:?}", other), } } + + #[test] + fn version_selector_rejects_ambiguous_or_empty_values() { + assert!(validate_version_selector(Some("v1"), None).is_ok()); + assert!(validate_version_selector(None, Some("1h")).is_ok()); + assert!(validate_version_selector(Some("v1"), Some("1h")).is_err()); + assert!(validate_version_selector(Some(""), None).is_err()); + } + + #[test] + fn versioning_errors_map_to_stable_exit_codes() { + assert_eq!( + exit_code_for_core_error(&rc_core::Error::VersionNotFound { + path: "local/bucket/key".to_string(), + version_id: "v1".to_string(), + }), + ExitCode::NotFound + ); + assert_eq!( + exit_code_for_core_error(&rc_core::Error::Auth("denied".to_string())), + ExitCode::AuthError + ); + assert_eq!( + exit_code_for_core_error(&rc_core::Error::GovernanceDenied { + path: "local/bucket/key".to_string(), + version_id: Some("v1".to_string()), + }), + ExitCode::Conflict + ); + } } diff --git a/crates/cli/src/commands/rm.rs b/crates/cli/src/commands/rm.rs index 310b94b..deb3bb1 100644 --- a/crates/cli/src/commands/rm.rs +++ b/crates/cli/src/commands/rm.rs @@ -3,18 +3,26 @@ //! Removes one or more objects from a bucket. use clap::Args; -use rc_core::{AliasManager, ListOptions, ObjectStore as _, RemotePath}; -use rc_s3::{DeleteRequestOptions, S3Client}; +use rc_core::{ + AliasManager, DeleteRequestOptions, Error, ListObjectVersionsOptions, ListOptions, ObjectStore, + ObjectVersionIdentifier, RemotePath, +}; +use rc_s3::S3Client; use serde::Serialize; use std::collections::HashSet; +use crate::commands::exit_code_for_core_error; use crate::exit_code::ExitCode; -use crate::output::{Formatter, OutputConfig}; +use crate::output::{ + Formatter, OutputConfig, V3ErrorEnvelope, V3PartialErrorEnvelope, V3SuccessEnvelope, +}; const RM_AFTER_HELP: &str = "\ Examples: rc object remove local/my-bucket/reports/2026-04.csv + rc rm local/my-bucket/reports/2026-04.csv --version-id VERSION_ID rc rm local/my-bucket/reports/ --recursive --dry-run + rc rm local/my-bucket/reports/ --recursive --versions --bypass rc object remove local/my-bucket/archive/ --recursive --force"; /// Remove objects @@ -41,11 +49,15 @@ pub struct RmArgs { #[arg(long)] pub incomplete: bool, - /// Include versions (requires versioning support) + /// Remove all matching object versions and delete markers #[arg(long)] pub versions: bool, - /// Bypass governance retention + /// Remove one exact object version + #[arg(long, value_name = "VERSION_ID")] + pub version_id: Option, + + /// Explicitly bypass Object Lock governance retention #[arg(long)] pub bypass: bool, @@ -60,51 +72,169 @@ struct RmOutput { deleted: Vec, #[serde(skip_serializing_if = "Option::is_none")] failed: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + deleted_versions: Option>, total: usize, } +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +struct RemovalRecord { + path: String, + #[serde(skip_serializing_if = "Option::is_none")] + version_id: Option, + is_delete_marker: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +struct VersionRemovalItem { + path: String, + version_id: Option, + delete_marker: bool, +} + +impl From for VersionRemovalItem { + fn from(record: RemovalRecord) -> Self { + Self { + path: record.path, + version_id: record.version_id, + delete_marker: record.is_delete_marker, + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +struct RemovalFailureRecord { + path: String, + version_id: Option, + error_type: &'static str, + message: String, +} + +#[derive(Debug)] +struct RemovalError { + code: ExitCode, + removed: Vec, + failed: Vec, +} + +impl RemovalError { + fn failed(code: ExitCode, failure: RemovalFailureRecord) -> Self { + Self { + code, + removed: Vec::new(), + failed: vec![failure], + } + } + + fn partial( + code: ExitCode, + removed: Vec, + failed: Vec, + ) -> Self { + Self { + code, + removed, + failed, + } + } +} + +type RemovalResult = Result, RemovalError>; + +#[derive(Debug, Serialize)] +struct VersionRemoveData { + operation: &'static str, + outcome: &'static str, + dry_run: bool, + planned: Vec, + removed: Vec, + failed: Vec, + summary: VersionRemoveSummary, +} + +#[derive(Debug, Serialize)] +struct VersionRemoveSummary { + planned: usize, + removed: usize, + failed: usize, +} + /// Execute the rm command pub async fn execute(args: RmArgs, output_config: OutputConfig) -> ExitCode { let formatter = Formatter::new(output_config); + let version_output = uses_version_output(&args); - if args.incomplete || args.versions || args.bypass { - return formatter.fail( + if args.incomplete { + return fail_rm( + &formatter, + version_output, ExitCode::UnsupportedFeature, - "--incomplete, --versions, and --bypass are not implemented; refusing to continue with a silently ignored destructive option", + "--incomplete is not implemented; refusing to continue with a silently ignored destructive option", ); } + if let Err(error) = validate_rm_selectors(&args) { + return fail_rm(&formatter, version_output, ExitCode::UsageError, &error); + } // Process each path let mut all_deleted = Vec::new(); let mut all_failed = Vec::new(); - let mut has_error = false; + let mut first_error_code = None; for path_str in &args.paths { - match process_rm_path(path_str, &args, &formatter).await { - Ok(deleted) => all_deleted.extend(deleted), - Err((code, failed)) => { - has_error = true; - all_failed.extend(failed); - if code != ExitCode::Success { - // Continue processing other paths unless it's a critical error - if code == ExitCode::AuthError || code == ExitCode::UsageError { - return code; - } - } + let critical = collect_removal_result( + process_rm_path(path_str, &args, &formatter).await, + &mut all_deleted, + &mut all_failed, + &mut first_error_code, + ); + if let Some(code) = critical { + if formatter.is_json() && version_output { + break; } + return code; } } // Output summary - if formatter.is_json() { + if formatter.is_json() && version_output { + let failure_count = all_failed.len(); + let data = build_version_remove_data(args.dry_run, all_deleted, all_failed); + if let Some(code) = first_error_code { + let action = if args.dry_run { + "plan removal of" + } else { + "remove" + }; + formatter.json_error(&V3PartialErrorEnvelope::versioned_objects( + code, + format!("Failed to {action} {failure_count} object version(s)"), + Some("delete_object_versions"), + data, + )); + } else { + formatter.json(&V3SuccessEnvelope::versioned_objects(data)); + } + } else if formatter.is_json() { let output = RmOutput { - status: if has_error { "partial" } else { "success" }, - deleted: all_deleted.clone(), + status: if first_error_code.is_some() { + "partial" + } else { + "success" + }, + deleted: all_deleted + .iter() + .map(|record| record.path.clone()) + .collect(), failed: if all_failed.is_empty() { None } else { - Some(all_failed) + Some(all_failed.into_iter().map(|failure| failure.path).collect()) }, + deleted_versions: all_deleted + .iter() + .any(|record| record.version_id.is_some()) + .then(|| all_deleted.clone()), total: all_deleted.len(), }; formatter.json(&output); @@ -112,58 +242,218 @@ pub async fn execute(args: RmArgs, output_config: OutputConfig) -> ExitCode { formatter.success(&format!("Removed {} object(s).", all_deleted.len())); } - if has_error { - ExitCode::GeneralError + first_error_code.unwrap_or(ExitCode::Success) +} + +fn uses_version_output(args: &RmArgs) -> bool { + args.version_id.is_some() || args.versions +} + +fn fail_rm(formatter: &Formatter, version_output: bool, code: ExitCode, message: &str) -> ExitCode { + if formatter.is_json() && version_output { + formatter.json_error(&V3ErrorEnvelope::versioned_objects( + code, + message, + Some("versioned_objects"), + )); + code } else { - ExitCode::Success + formatter.fail(code, message) } } -async fn process_rm_path( - path_str: &str, - args: &RmArgs, +fn report_rm_error( formatter: &Formatter, -) -> Result, (ExitCode, Vec)> { + args: &RmArgs, + code: ExitCode, + message: &str, + suggestion: Option<&str>, +) -> ExitCode { + if formatter.is_json() && uses_version_output(args) { + return code; + } + if let Some(suggestion) = suggestion { + formatter.fail_with_suggestion(code, message, suggestion) + } else { + formatter.fail(code, message) + } +} + +fn removal_failure( + path: impl Into, + version_id: Option, + code: ExitCode, + message: impl Into, +) -> RemovalFailureRecord { + RemovalFailureRecord { + path: path.into(), + version_id, + error_type: error_type_for_exit_code(code), + message: message.into(), + } +} + +fn error_type_for_exit_code(code: ExitCode) -> &'static str { + match code { + ExitCode::Success | ExitCode::GeneralError => "general_error", + ExitCode::UsageError => "usage_error", + ExitCode::NetworkError => "network_error", + ExitCode::AuthError => "auth_error", + ExitCode::NotFound => "not_found", + ExitCode::Conflict => "conflict", + ExitCode::UnsupportedFeature => "unsupported_feature", + ExitCode::Interrupted => "interrupted", + } +} + +fn collect_removal_result( + result: RemovalResult, + all_removed: &mut Vec, + all_failed: &mut Vec, + first_error_code: &mut Option, +) -> Option { + match result { + Ok(removed) => { + all_removed.extend(removed); + None + } + Err(error) => { + record_removal_error(first_error_code, error.code); + all_removed.extend(error.removed); + all_failed.extend(error.failed); + is_critical_removal_error(error.code).then_some(error.code) + } + } +} + +fn record_removal_error(current: &mut Option, candidate: ExitCode) { + match current { + None => *current = Some(candidate), + Some(existing) + if is_critical_removal_error(candidate) && !is_critical_removal_error(*existing) => + { + *current = Some(candidate); + } + Some(_) => {} + } +} + +fn is_critical_removal_error(code: ExitCode) -> bool { + matches!( + code, + ExitCode::AuthError | ExitCode::UsageError | ExitCode::Interrupted + ) +} + +fn build_version_remove_data( + dry_run: bool, + records: Vec, + failed: Vec, +) -> VersionRemoveData { + let records = records + .into_iter() + .map(VersionRemovalItem::from) + .collect::>(); + let (planned, removed) = if dry_run { + (records, Vec::new()) + } else { + (Vec::new(), records) + }; + let outcome = if !failed.is_empty() { + if !dry_run && !removed.is_empty() { + "partial" + } else { + "failed" + } + } else if dry_run && !planned.is_empty() { + "planned" + } else if removed.is_empty() { + "empty" + } else { + "success" + }; + let summary = VersionRemoveSummary { + planned: planned.len(), + removed: removed.len(), + failed: failed.len(), + }; + + VersionRemoveData { + operation: "remove", + outcome, + dry_run, + planned, + removed, + failed, + summary, + } +} + +async fn process_rm_path(path_str: &str, args: &RmArgs, formatter: &Formatter) -> RemovalResult { // Parse the path let (alias_name, bucket, key) = match parse_rm_path(path_str) { Ok(parsed) => parsed, Err(e) => { - let code = formatter.fail_with_suggestion( + let code = report_rm_error( + formatter, + args, ExitCode::UsageError, &e, - "Use a remote path in the form alias/bucket[/key] before retrying the remove command.", + Some( + "Use a remote path in the form alias/bucket[/key] before retrying the remove command.", + ), ); - return Err((code, vec![path_str.to_string()])); + return Err(RemovalError::failed( + code, + removal_failure(path_str, args.version_id.clone(), code, e), + )); } }; if let Err(error) = validate_removal_scope(&key, args.recursive) { - let code = formatter.fail_with_suggestion( + let code = report_rm_error( + formatter, + args, ExitCode::UsageError, &error, - "Add --recursive only after verifying the bucket or prefix to remove.", + Some("Add --recursive only after verifying the bucket or prefix to remove."), ); - return Err((code, vec![path_str.to_string()])); + return Err(RemovalError::failed( + code, + removal_failure(path_str, args.version_id.clone(), code, error), + )); } // Load alias let alias_manager = match AliasManager::new() { Ok(am) => am, Err(e) => { - formatter.error(&format!("Failed to load aliases: {e}")); - return Err((ExitCode::GeneralError, vec![])); + let message = format!("Failed to load aliases: {e}"); + let code = report_rm_error(formatter, args, ExitCode::GeneralError, &message, None); + return Err(RemovalError::failed( + code, + removal_failure(path_str, args.version_id.clone(), code, message), + )); } }; let alias = match alias_manager.get(&alias_name) { Ok(a) => a, Err(_) => { - let code = formatter.fail_with_suggestion( + let message = format!("Alias '{alias_name}' not found"); + let code = report_rm_error( + formatter, + args, ExitCode::NotFound, - &format!("Alias '{alias_name}' not found"), - "Run `rc alias list` to inspect configured aliases or add one with `rc alias set ...`.", + &message, + Some( + "Run `rc alias list` to inspect configured aliases or add one with `rc alias set ...`.", + ), ); - return Err((code, vec![])); + return Err(RemovalError::failed( + code, + removal_failure(path_str, args.version_id.clone(), code, message), + )); } }; @@ -171,15 +461,18 @@ async fn process_rm_path( let client = match S3Client::new(alias).await { Ok(c) => c, Err(e) => { - let code = formatter.fail( - ExitCode::NetworkError, - &format!("Failed to create S3 client: {e}"), - ); - return Err((code, vec![])); + let message = format!("Failed to create S3 client: {e}"); + let code = report_rm_error(formatter, args, ExitCode::NetworkError, &message, None); + return Err(RemovalError::failed( + code, + removal_failure(path_str, args.version_id.clone(), code, message), + )); } }; - if args.recursive { + if args.versions { + delete_versions(&client, &alias_name, &bucket, &key, args, formatter).await + } else if args.recursive { delete_recursive(&client, &alias_name, &bucket, &key, args, formatter).await } else { // Delete single object @@ -194,51 +487,82 @@ async fn delete_single( key: &str, args: &RmArgs, formatter: &Formatter, -) -> Result, (ExitCode, Vec)> { +) -> RemovalResult { let path = RemotePath::new(alias_name, bucket, key); let full_path = format!("{alias_name}/{bucket}/{key}"); if args.dry_run { - let styled_path = formatter.style_file(&full_path); - formatter.println(&format!("Would remove: {styled_path}")); - return Ok(vec![full_path]); + if !formatter.is_json() { + let styled_path = formatter.style_file(&full_path); + let suffix = args + .version_id + .as_deref() + .map(|version| format!(" (version {version})")) + .unwrap_or_default(); + formatter.println(&format!("Would remove: {styled_path}{suffix}")); + } + return Ok(vec![RemovalRecord { + path: full_path, + version_id: args.version_id.clone(), + is_delete_marker: false, + }]); } - match client - .delete_object_with_options(&path, delete_request_options(args)) - .await + match ObjectStore::delete_object_with_options(client, &path, delete_request_options(args)).await { - Ok(()) => { + Ok(deleted) => { if !formatter.is_json() { let styled_path = formatter.style_file(&full_path); - formatter.println(&format!("Removed: {styled_path}")); + let version = deleted + .version_id + .as_deref() + .map(|version_id| format!(" (version {version_id})")) + .unwrap_or_default(); + let marker = if deleted.is_delete_marker { + " [delete marker]" + } else { + "" + }; + formatter.println(&format!("Removed: {styled_path}{version}{marker}")); } - Ok(vec![full_path]) + Ok(vec![RemovalRecord { + path: full_path, + version_id: deleted.version_id.or_else(|| args.version_id.clone()), + is_delete_marker: deleted.is_delete_marker, + }]) } Err(e) => { - let err_str = e.to_string(); - if err_str.contains("NotFound") || err_str.contains("NoSuchKey") { + if matches!( + &e, + Error::NotFound(_) | Error::VersionNotFound { .. } | Error::DeleteMarker { .. } + ) { if args.force { // Force mode: ignore not found errors Ok(vec![]) } else { - let code = formatter.fail_with_suggestion( + let message = e.to_string(); + let code = report_rm_error( + formatter, + args, ExitCode::NotFound, - &format!("Object not found: {full_path}"), - "Check the object key or retry with --force if missing objects are acceptable.", + &message, + Some( + "Check the object key or retry with --force if missing objects are acceptable.", + ), ); - Err((code, vec![full_path])) + Err(RemovalError::failed( + code, + removal_failure(full_path, args.version_id.clone(), code, message), + )) } - } else if err_str.contains("AccessDenied") { - let code = - formatter.fail(ExitCode::AuthError, &format!("Access denied: {full_path}")); - Err((code, vec![full_path])) } else { - let code = formatter.fail( - ExitCode::NetworkError, - &format!("Failed to remove {full_path}: {e}"), - ); - Err((code, vec![full_path])) + let exit_code = exit_code_for_core_error(&e); + let message = format!("Failed to remove {full_path}: {e}"); + let code = report_rm_error(formatter, args, exit_code, &message, None); + Err(RemovalError::failed( + code, + removal_failure(full_path, args.version_id.clone(), code, message), + )) } } } @@ -251,7 +575,7 @@ async fn delete_recursive( prefix: &str, args: &RmArgs, formatter: &Formatter, -) -> Result, (ExitCode, Vec)> { +) -> RemovalResult { let path = RemotePath::new(alias_name, bucket, prefix); // Collect all objects to delete @@ -283,18 +607,25 @@ async fn delete_recursive( Err(e) => { let err_str = e.to_string(); if err_str.contains("NotFound") || err_str.contains("NoSuchBucket") { - let code = formatter.fail_with_suggestion( + let message = format!("Bucket not found: {bucket}"); + let code = report_rm_error( + formatter, + args, ExitCode::NotFound, - &format!("Bucket not found: {bucket}"), - "Check the bucket path and retry the remove command.", + &message, + Some("Check the bucket path and retry the remove command."), ); - return Err((code, vec![])); + return Err(RemovalError::failed( + code, + removal_failure(path.to_string(), None, code, message), + )); } - let code = formatter.fail( - ExitCode::NetworkError, - &format!("Failed to list objects: {e}"), - ); - return Err((code, vec![])); + let message = format!("Failed to list objects: {e}"); + let code = report_rm_error(formatter, args, ExitCode::NetworkError, &message, None); + return Err(RemovalError::failed( + code, + removal_failure(path.to_string(), None, code, message), + )); } } } @@ -312,12 +643,18 @@ async fn delete_recursive( if args.dry_run { for key in &keys_to_delete { let full_path = format!("{alias_name}/{bucket}/{key}"); - let styled_path = formatter.style_file(&full_path); - formatter.println(&format!("Would remove: {styled_path}")); + if !formatter.is_json() { + let styled_path = formatter.style_file(&full_path); + formatter.println(&format!("Would remove: {styled_path}")); + } } return Ok(keys_to_delete .iter() - .map(|k| format!("{alias_name}/{bucket}/{k}")) + .map(|key| RemovalRecord { + path: format!("{alias_name}/{bucket}/{key}"), + version_id: None, + is_delete_marker: false, + }) .collect()); } @@ -327,19 +664,25 @@ async fn delete_recursive( let mut first_error_code = None; for key in keys_to_delete { - match delete_single(client, alias_name, bucket, &key, args, formatter).await { - Ok(paths) => deleted.extend(paths), - Err((code, paths)) => { - first_error_code.get_or_insert(code); - failed.extend(paths); - } + let critical = collect_removal_result( + delete_single(client, alias_name, bucket, &key, args, formatter).await, + &mut deleted, + &mut failed, + &mut first_error_code, + ); + if critical.is_some() { + break; } } return if failed.is_empty() { Ok(deleted) } else { - Err((first_error_code.unwrap_or(ExitCode::GeneralError), failed)) + Err(RemovalError::partial( + first_error_code.unwrap_or(ExitCode::GeneralError), + deleted, + failed, + )) }; } @@ -363,33 +706,348 @@ async fn delete_recursive( let styled_path = formatter.style_file(&full_path); formatter.println(&format!("Removed: {styled_path}")); } - deleted.push(full_path); + deleted.push(RemovalRecord { + path: full_path, + version_id: None, + is_delete_marker: false, + }); } - failed.extend( - failed_keys - .into_iter() - .map(|key| format!("{alias_name}/{bucket}/{key}")), - ); + failed.extend(failed_keys.into_iter().map(|key| { + let path = format!("{alias_name}/{bucket}/{key}"); + removal_failure( + path, + None, + ExitCode::GeneralError, + "The backend omitted the object from its delete result", + ) + })); } Err(e) => { - formatter.error_with_code( - ExitCode::GeneralError, - &format!("Failed to delete batch: {e}"), - ); + let message = format!("Failed to delete batch: {e}"); + report_rm_error(formatter, args, ExitCode::GeneralError, &message, None); for key in chunk_keys { - failed.push(format!("{alias_name}/{bucket}/{key}")); + failed.push(removal_failure( + format!("{alias_name}/{bucket}/{key}"), + None, + ExitCode::GeneralError, + message.clone(), + )); } } } } if !failed.is_empty() { - Err((ExitCode::GeneralError, failed)) + Err(RemovalError::partial( + ExitCode::GeneralError, + deleted, + failed, + )) } else { Ok(deleted) } } +async fn delete_versions( + client: &S3Client, + alias_name: &str, + bucket: &str, + key: &str, + args: &RmArgs, + formatter: &Formatter, +) -> RemovalResult { + let path = RemotePath::new(alias_name, bucket, key); + let versions = match list_versions_for_removal(client, &path, args.recursive).await { + Ok(versions) => versions, + Err(error) => { + let exit_code = exit_code_for_core_error(&error); + let message = format!("Failed to list object versions for removal: {error}"); + report_rm_error(formatter, args, exit_code, &message, None); + return Err(RemovalError::failed( + exit_code, + removal_failure(path.to_string(), None, exit_code, message), + )); + } + }; + + if versions.is_empty() { + if !args.force { + formatter.warning(&format!("No object versions found: {path}")); + } + return Ok(Vec::new()); + } + + if args.dry_run { + let records = versions + .into_iter() + .map(|version| { + let full_path = format!("{alias_name}/{bucket}/{}", version.key); + if !formatter.is_json() { + let marker = if version.is_delete_marker { + ", delete marker" + } else { + "" + }; + formatter.println(&format!( + "Would remove: {} (version {}{marker})", + formatter.style_file(&full_path), + version.version_id.as_deref().unwrap_or("null") + )); + } + RemovalRecord { + path: full_path, + version_id: version.version_id, + is_delete_marker: version.is_delete_marker, + } + }) + .collect(); + return Ok(records); + } + + let mut removed = Vec::new(); + let mut failed = Vec::new(); + let mut first_error_code = None; + + for (chunk_index, chunk) in versions.chunks(1000).enumerate() { + let requested = chunk.to_vec(); + match ObjectStore::delete_object_versions( + client, + bucket, + requested.clone(), + delete_request_options(args), + ) + .await + { + Ok(result) => { + let mut unmatched = requested.clone(); + for deleted in result.deleted { + let requested_entry = take_requested_version( + &mut unmatched, + &deleted.key, + deleted.version_id.as_deref(), + ); + let version_id = deleted.version_id.or_else(|| { + requested_entry + .as_ref() + .and_then(|entry| entry.version_id.clone()) + }); + let is_delete_marker = deleted.is_delete_marker + || requested_entry.is_some_and(|entry| entry.is_delete_marker); + let full_path = format!("{alias_name}/{bucket}/{}", deleted.key); + if !formatter.is_json() { + let marker = if is_delete_marker { + ", delete marker" + } else { + "" + }; + formatter.println(&format!( + "Removed: {} (version {}{marker})", + formatter.style_file(&full_path), + version_id.as_deref().unwrap_or("null") + )); + } + removed.push(RemovalRecord { + path: full_path, + version_id, + is_delete_marker, + }); + } + + for failure in result.failures { + let requested_entry = take_requested_version( + &mut unmatched, + &failure.key, + failure.version_id.as_deref(), + ); + let code = exit_code_for_delete_failure( + failure.code.as_deref(), + failure.message.as_deref(), + ); + record_removal_error(&mut first_error_code, code); + let full_path = format!("{alias_name}/{bucket}/{}", failure.key); + let version_id = failure + .version_id + .or_else(|| requested_entry.and_then(|entry| entry.version_id)); + let message = failure + .message + .as_deref() + .or(failure.code.as_deref()) + .unwrap_or("unknown delete error") + .to_string(); + report_rm_error( + formatter, + args, + code, + &format!( + "Failed to remove {} (version {}): {}", + full_path, + version_id.as_deref().unwrap_or("null"), + message + ), + None, + ); + failed.push(removal_failure(full_path, version_id, code, message)); + } + + for omitted in unmatched { + let full_path = format!("{alias_name}/{bucket}/{}", omitted.key); + let message = "The backend omitted the object version from its delete result"; + record_removal_error(&mut first_error_code, ExitCode::GeneralError); + report_rm_error( + formatter, + args, + ExitCode::GeneralError, + &format!( + "Failed to remove {} (version {}): {message}", + full_path, + omitted.version_id.as_deref().unwrap_or("null") + ), + None, + ); + failed.push(removal_failure( + full_path, + omitted.version_id, + ExitCode::GeneralError, + message, + )); + } + } + Err(error) => { + let code = exit_code_for_core_error(&error); + record_removal_error(&mut first_error_code, code); + let message = format!("Failed to delete versions: {error}"); + report_rm_error(formatter, args, code, &message, None); + failed.extend(requested.into_iter().map(|entry| { + removal_failure( + format!("{alias_name}/{bucket}/{}", entry.key), + entry.version_id, + code, + message.clone(), + ) + })); + if is_critical_removal_error(code) { + failed.extend(versions.iter().skip((chunk_index + 1) * 1000).map(|entry| { + removal_failure( + format!("{alias_name}/{bucket}/{}", entry.key), + entry.version_id.clone(), + code, + format!( + "Not attempted after a critical version deletion failure: {error}" + ), + ) + })); + break; + } + } + } + } + + if failed.is_empty() { + Ok(removed) + } else { + Err(RemovalError::partial( + first_error_code.unwrap_or(ExitCode::GeneralError), + removed, + failed, + )) + } +} + +fn take_requested_version( + requested: &mut Vec, + key: &str, + version_id: Option<&str>, +) -> Option { + let position = requested.iter().position(|entry| { + entry.key == key + && match version_id { + Some(version_id) => entry.version_id.as_deref() == Some(version_id), + None => true, + } + })?; + Some(requested.remove(position)) +} + +async fn list_versions_for_removal( + client: &S3Client, + path: &RemotePath, + recursive: bool, +) -> Result, Error> { + let mut versions = Vec::new(); + let mut key_marker = None; + let mut version_id_marker = None; + + loop { + let page = ObjectStore::list_object_versions_page_with_options( + client, + path, + &ListObjectVersionsOptions { + max_keys: Some(1000), + key_marker: key_marker.clone(), + version_id_marker: version_id_marker.clone(), + }, + ) + .await?; + versions.extend( + page.items + .into_iter() + .filter(|version| recursive || version.key == path.key) + .map(|version| ObjectVersionIdentifier { + key: version.key, + version_id: Some(version.version_id), + is_delete_marker: version.is_delete_marker, + }), + ); + + if !page.truncated { + break; + } + let next_key_marker = page.continuation_token.ok_or_else(|| { + Error::Network( + "S3 returned a truncated version listing without a key marker".to_string(), + ) + })?; + let next_version_id_marker = page.version_id_marker; + if key_marker.as_deref() == Some(next_key_marker.as_str()) + && version_id_marker == next_version_id_marker + { + return Err(Error::Network( + "S3 returned a truncated version listing without advancing its markers".to_string(), + )); + } + key_marker = Some(next_key_marker); + version_id_marker = next_version_id_marker; + } + + Ok(versions) +} + +fn exit_code_for_delete_failure(code: Option<&str>, message: Option<&str>) -> ExitCode { + let normalized_message = message.unwrap_or_default().to_ascii_lowercase(); + if matches!( + code, + Some("NoSuchVersion") | Some("NoSuchKey") | Some("NotFound") + ) { + ExitCode::NotFound + } else if normalized_message.contains("governance") + || normalized_message.contains("retention") + || normalized_message.contains("object lock") + || normalized_message.contains("worm") + { + ExitCode::Conflict + } else if matches!( + code, + Some("AccessDenied") | Some("Forbidden") | Some("Unauthorized") + ) || normalized_message.contains("access denied") + || normalized_message.contains("forbidden") + || normalized_message.contains("unauthorized") + { + ExitCode::AuthError + } else { + ExitCode::GeneralError + } +} + /// Parse rm path into (alias, bucket, key) fn parse_rm_path(path: &str) -> Result<(String, String, String), String> { if path.is_empty() { @@ -425,10 +1083,31 @@ fn parse_rm_path(path: &str) -> Result<(String, String, String), String> { fn delete_request_options(args: &RmArgs) -> DeleteRequestOptions { DeleteRequestOptions { + version_id: args.version_id.clone(), + bypass_governance: args.bypass, force_delete: args.purge, } } +fn validate_rm_selectors(args: &RmArgs) -> Result<(), String> { + if args.version_id.as_deref().is_some_and(str::is_empty) { + return Err("--version-id cannot be empty".to_string()); + } + if args.version_id.is_some() && args.versions { + return Err("--version-id cannot be combined with --versions".to_string()); + } + if args.version_id.is_some() && args.recursive { + return Err("--version-id cannot be combined with --recursive".to_string()); + } + if args.version_id.is_some() && args.purge { + return Err("--version-id cannot be combined with --purge".to_string()); + } + if args.version_id.is_some() && args.paths.len() != 1 { + return Err("--version-id requires exactly one object path".to_string()); + } + Ok(()) +} + fn validate_removal_scope(key: &str, recursive: bool) -> Result<(), String> { if !recursive && (key.is_empty() || key.ends_with('/')) { return Err("Bucket and prefix removal requires the explicit --recursive flag".to_string()); @@ -518,6 +1197,7 @@ mod tests { dry_run: false, incomplete: false, versions: false, + version_id: None, bypass: false, purge: true, }; @@ -535,6 +1215,7 @@ mod tests { dry_run: false, incomplete: false, versions: false, + version_id: None, bypass: false, purge: false, }; @@ -552,6 +1233,7 @@ mod tests { dry_run: false, incomplete: false, versions: false, + version_id: None, bypass: false, purge: false, }; @@ -559,4 +1241,233 @@ mod tests { let options = delete_request_options(&args); assert!(!options.force_delete); } + + #[test] + fn delete_request_options_require_explicit_bypass_and_preserve_version() { + let args = RmArgs { + paths: vec!["test/bucket/object.txt".to_string()], + recursive: false, + force: false, + dry_run: false, + incomplete: false, + versions: false, + version_id: Some("v1".to_string()), + bypass: true, + purge: false, + }; + + let options = delete_request_options(&args); + assert_eq!(options.version_id.as_deref(), Some("v1")); + assert!(options.bypass_governance); + + let default_args = RmArgs { + bypass: false, + ..args + }; + assert!(!delete_request_options(&default_args).bypass_governance); + } + + #[test] + fn rm_rejects_conflicting_version_selectors() { + let args = RmArgs { + paths: vec!["test/bucket/object.txt".to_string()], + recursive: false, + force: false, + dry_run: false, + incomplete: false, + versions: true, + version_id: Some("v1".to_string()), + bypass: false, + purge: false, + }; + + assert!(validate_rm_selectors(&args).is_err()); + } + + #[test] + fn batch_delete_failures_distinguish_access_and_governance_denials() { + assert_eq!( + exit_code_for_delete_failure(Some("NoSuchVersion"), Some("missing")), + ExitCode::NotFound + ); + assert_eq!( + exit_code_for_delete_failure(Some("AccessDenied"), Some("policy denied")), + ExitCode::AuthError + ); + assert_eq!( + exit_code_for_delete_failure( + Some("AccessDenied"), + Some("governance retention is active") + ), + ExitCode::Conflict + ); + } + + #[test] + fn rm_output_preserves_version_and_delete_marker_fields() { + let output = RmOutput { + status: "success", + deleted: vec!["test/bucket/key.txt".to_string()], + failed: None, + deleted_versions: Some(vec![RemovalRecord { + path: "test/bucket/key.txt".to_string(), + version_id: Some("marker-v1".to_string()), + is_delete_marker: true, + }]), + total: 1, + }; + + let json = serde_json::to_value(output).expect("serialize version-aware removal"); + assert_eq!(json["deleted"][0], "test/bucket/key.txt"); + assert_eq!(json["deleted_versions"][0]["version_id"], "marker-v1"); + assert_eq!(json["deleted_versions"][0]["is_delete_marker"], true); + } + + #[test] + fn partial_version_removal_keeps_successes_and_version_aware_failures() { + let removed = RemovalRecord { + path: "test/bucket/key.txt".to_string(), + version_id: Some("v1".to_string()), + is_delete_marker: false, + }; + let failure = RemovalFailureRecord { + path: "test/bucket/key.txt".to_string(), + version_id: Some("v2".to_string()), + error_type: "conflict", + message: "Governance retention denied deletion".to_string(), + }; + let mut all_removed = Vec::new(); + let mut all_failed = Vec::new(); + let mut first_error_code = None; + + let critical = collect_removal_result( + Err(RemovalError { + code: ExitCode::Conflict, + removed: vec![removed.clone()], + failed: vec![failure.clone()], + }), + &mut all_removed, + &mut all_failed, + &mut first_error_code, + ); + + assert_eq!(all_removed, vec![removed]); + assert_eq!(all_failed, vec![failure]); + assert_eq!(first_error_code, Some(ExitCode::Conflict)); + assert_eq!(critical, None); + } + + #[test] + fn critical_removal_error_overrides_an_earlier_noncritical_error() { + let mut removed = Vec::new(); + let mut failed = Vec::new(); + let mut aggregate_code = None; + + let not_found = RemovalError::failed( + ExitCode::NotFound, + removal_failure( + "test/bucket/missing.txt", + Some("missing-v1".to_string()), + ExitCode::NotFound, + "version missing", + ), + ); + assert_eq!( + collect_removal_result( + Err(not_found), + &mut removed, + &mut failed, + &mut aggregate_code, + ), + None + ); + + let access_denied = RemovalError::failed( + ExitCode::AuthError, + removal_failure( + "test/bucket/private.txt", + Some("private-v1".to_string()), + ExitCode::AuthError, + "access denied", + ), + ); + assert_eq!( + collect_removal_result( + Err(access_denied), + &mut removed, + &mut failed, + &mut aggregate_code, + ), + Some(ExitCode::AuthError) + ); + assert_eq!(aggregate_code, Some(ExitCode::AuthError)); + assert_eq!(failed.len(), 2); + } + + #[test] + fn version_delete_reconciliation_consumes_duplicate_keys_once() { + let mut requested = vec![ + ObjectVersionIdentifier { + key: "key.txt".to_string(), + version_id: Some("v1".to_string()), + is_delete_marker: false, + }, + ObjectVersionIdentifier { + key: "key.txt".to_string(), + version_id: Some("v2".to_string()), + is_delete_marker: true, + }, + ]; + + let first = take_requested_version(&mut requested, "key.txt", None) + .expect("the first unqualified delete result must match one request"); + let second = take_requested_version(&mut requested, "key.txt", None) + .expect("the second unqualified delete result must match the remaining request"); + + assert_eq!(first.version_id.as_deref(), Some("v1")); + assert_eq!(second.version_id.as_deref(), Some("v2")); + assert!(requested.is_empty()); + } + + #[test] + fn version_remove_data_marks_dry_run_items_as_planned() { + let planned = RemovalRecord { + path: "test/bucket/key.txt".to_string(), + version_id: Some("v1".to_string()), + is_delete_marker: false, + }; + + let data = build_version_remove_data(true, vec![planned.clone()], Vec::new()); + let json = serde_json::to_value(data).expect("serialize version removal dry-run"); + + assert_eq!(json["outcome"], "planned"); + assert_eq!(json["dry_run"], true); + assert_eq!(json["planned"][0]["version_id"], "v1"); + assert_eq!(json["removed"].as_array().map(Vec::len), Some(0)); + assert_eq!(json["summary"]["planned"], 1); + } + + #[test] + fn version_remove_data_marks_dry_run_discovery_errors_as_failed() { + let planned = RemovalRecord { + path: "test/bucket/key.txt".to_string(), + version_id: Some("v1".to_string()), + is_delete_marker: false, + }; + let failure = removal_failure( + "test/bucket/private.txt", + None, + ExitCode::AuthError, + "access denied", + ); + + let data = build_version_remove_data(true, vec![planned], vec![failure]); + let json = serde_json::to_value(data).expect("serialize failed version removal dry-run"); + + assert_eq!(json["outcome"], "failed"); + assert_eq!(json["dry_run"], true); + assert_eq!(json["planned"].as_array().map(Vec::len), Some(1)); + assert_eq!(json["failed"].as_array().map(Vec::len), Some(1)); + assert_eq!(json["removed"].as_array().map(Vec::len), Some(0)); + } } diff --git a/crates/cli/src/commands/stat.rs b/crates/cli/src/commands/stat.rs index ec4a007..7f43cfb 100644 --- a/crates/cli/src/commands/stat.rs +++ b/crates/cli/src/commands/stat.rs @@ -3,13 +3,14 @@ //! Displays detailed metadata information about an object. use clap::Args; -use rc_core::{AliasManager, ObjectInfo, ObjectStore as _, RemotePath}; +use rc_core::{AliasManager, ObjectInfo, ObjectReadOptions, ObjectStore, RemotePath}; use rc_s3::S3Client; use serde::Serialize; use std::collections::BTreeMap; +use crate::commands::{exit_code_for_core_error, validate_version_selector}; use crate::exit_code::ExitCode; -use crate::output::{Formatter, OutputConfig}; +use crate::output::{Formatter, OutputConfig, V3ErrorEnvelope, V3SuccessEnvelope}; /// Show object metadata #[derive(Args, Debug)] @@ -17,7 +18,7 @@ pub struct StatArgs { /// Object path (alias/bucket/key) pub path: String, - /// Show version ID information + /// Exact object version ID to inspect #[arg(long)] pub version_id: Option, @@ -47,6 +48,28 @@ struct StatOutput { metadata: Option>, } +#[derive(Debug, Serialize)] +struct VersionStatData { + operation: &'static str, + object: VersionStatObject, +} + +#[derive(Debug, Serialize)] +struct VersionStatObject { + path: String, + bucket: String, + key: String, + version_id: Option, + delete_marker: bool, + last_modified: Option, + size_bytes: Option, + size_human: Option, + etag: Option, + content_type: Option, + storage_class: Option, + metadata: Option>, +} + /// Returns true if metadata is None or contains an empty map. fn metadata_is_none_or_empty(metadata: &Option>) -> bool { match metadata { @@ -80,20 +103,38 @@ fn build_display_metadata(info: &ObjectInfo) -> BTreeMap { /// Execute the stat command pub async fn execute(args: StatArgs, output_config: OutputConfig) -> ExitCode { let formatter = Formatter::new(output_config); + let version_output = args.version_id.is_some(); - if args.version_id.is_some() || args.rewind.is_some() { - return formatter.fail( + if let Err(error) = + validate_version_selector(args.version_id.as_deref(), args.rewind.as_deref()) + { + return fail_stat(&formatter, version_output, ExitCode::UsageError, &error); + } + if args.rewind.is_some() { + return fail_stat( + &formatter, + version_output, ExitCode::UnsupportedFeature, - "--version-id and --rewind are not implemented for stat", + "--rewind is not implemented for stat", ); } + let read_options = match ObjectReadOptions::for_version(args.version_id.clone()) { + Ok(options) => options, + Err(error) => { + return fail_stat( + &formatter, + version_output, + ExitCode::UsageError, + &error.to_string(), + ); + } + }; // Parse the path let (alias_name, bucket, key) = match parse_stat_path(&args.path) { Ok(parsed) => parsed, Err(e) => { - formatter.error(&e); - return ExitCode::UsageError; + return fail_stat(&formatter, version_output, ExitCode::UsageError, &e); } }; @@ -101,16 +142,24 @@ pub async fn execute(args: StatArgs, output_config: OutputConfig) -> ExitCode { let alias_manager = match AliasManager::new() { Ok(am) => am, Err(e) => { - formatter.error(&format!("Failed to load aliases: {e}")); - return ExitCode::GeneralError; + return fail_stat( + &formatter, + version_output, + ExitCode::GeneralError, + &format!("Failed to load aliases: {e}"), + ); } }; let alias = match alias_manager.get(&alias_name) { Ok(a) => a, Err(_) => { - formatter.error(&format!("Alias '{alias_name}' not found")); - return ExitCode::NotFound; + return fail_stat( + &formatter, + version_output, + ExitCode::NotFound, + &format!("Alias '{alias_name}' not found"), + ); } }; @@ -118,37 +167,62 @@ pub async fn execute(args: StatArgs, output_config: OutputConfig) -> ExitCode { let client = match S3Client::new(alias).await { Ok(c) => c, Err(e) => { - formatter.error(&format!("Failed to create S3 client: {e}")); - return ExitCode::NetworkError; + return fail_stat( + &formatter, + version_output, + ExitCode::NetworkError, + &format!("Failed to create S3 client: {e}"), + ); } }; let path = RemotePath::new(&alias_name, &bucket, &key); // Get object metadata - match client.head_object(&path).await { + match ObjectStore::head_object_with_options(&client, &path, &read_options).await { Ok(info) => { if formatter.is_json() { - let output = StatOutput { - name: info.key.clone(), - last_modified: info.last_modified.map(|d| d.to_string()), - size_bytes: info.size_bytes, - size_human: info.size_human.clone(), - etag: info.etag.clone(), - content_type: info.content_type.clone(), - storage_class: info.storage_class.clone(), - version_id: None, - metadata: info - .metadata - .as_ref() - .map(|m| { - m.iter() - .map(|(k, v)| (k.clone(), v.clone())) - .collect::>() - }) - .filter(|m| !m.is_empty()), - }; - formatter.json(&output); + let metadata = info + .metadata + .as_ref() + .map(|m| { + m.iter() + .map(|(k, v)| (k.clone(), v.clone())) + .collect::>() + }) + .filter(|m| !m.is_empty()); + if version_output { + formatter.json(&V3SuccessEnvelope::versioned_objects(VersionStatData { + operation: "stat", + object: VersionStatObject { + path: args.path.clone(), + bucket: bucket.clone(), + key: info.key.clone(), + version_id: info.version_id.clone().or_else(|| args.version_id.clone()), + delete_marker: info.is_delete_marker.unwrap_or(false), + last_modified: info.last_modified.map(|d| d.to_string()), + size_bytes: info.size_bytes, + size_human: info.size_human.clone(), + etag: info.etag.clone(), + content_type: info.content_type.clone(), + storage_class: info.storage_class.clone(), + metadata, + }, + })); + } else { + let output = StatOutput { + name: info.key.clone(), + last_modified: info.last_modified.map(|d| d.to_string()), + size_bytes: info.size_bytes, + size_human: info.size_human.clone(), + etag: info.etag.clone(), + content_type: info.content_type.clone(), + storage_class: info.storage_class.clone(), + version_id: info.version_id.clone(), + metadata, + }; + formatter.json(&output); + } } else { // Helper to format key-value pairs with styling let format_kv = |key: &str, value: &str| { @@ -180,6 +254,9 @@ pub async fn execute(args: StatArgs, output_config: OutputConfig) -> ExitCode { if let Some(sc) = &info.storage_class { formatter.println(&format_kv("Class", sc)); } + if let Some(version_id) = &info.version_id { + formatter.println(&format_kv("Version", version_id)); + } let display_metadata = build_display_metadata(&info); if !display_metadata.is_empty() { formatter.println(&format_kv("Metadata", "")); @@ -191,21 +268,35 @@ pub async fn execute(args: StatArgs, output_config: OutputConfig) -> ExitCode { ExitCode::Success } Err(e) => { - let err_str = e.to_string(); - if err_str.contains("NotFound") || err_str.contains("NoSuchKey") { - formatter.error(&format!("Object not found: {}", args.path)); - ExitCode::NotFound - } else if err_str.contains("AccessDenied") { - formatter.error(&format!("Access denied: {}", args.path)); - ExitCode::AuthError - } else { - formatter.error(&format!("Failed to get object metadata: {e}")); - ExitCode::NetworkError - } + let exit_code = exit_code_for_core_error(&e); + fail_stat( + &formatter, + version_output, + exit_code, + &format!("Failed to inspect {}: {e}", args.path), + ) } } } +fn fail_stat( + formatter: &Formatter, + version_output: bool, + code: ExitCode, + message: &str, +) -> ExitCode { + if formatter.is_json() && version_output { + formatter.json_error(&V3ErrorEnvelope::versioned_objects( + code, + message, + Some("head_object_version"), + )); + code + } else { + formatter.fail(code, message) + } +} + /// Parse stat path into (alias, bucket, key) fn parse_stat_path(path: &str) -> Result<(String, String, String), String> { if path.is_empty() { @@ -367,6 +458,9 @@ mod tests { storage_class: None, content_type: Some("text/plain".to_string()), metadata: Some(user_meta), + version_id: None, + source_version_id: None, + is_delete_marker: None, is_dir: false, }; diff --git a/crates/cli/src/output/formatter.rs b/crates/cli/src/output/formatter.rs index a3d5820..4bfcbdb 100644 --- a/crates/cli/src/output/formatter.rs +++ b/crates/cli/src/output/formatter.rs @@ -471,6 +471,14 @@ impl Formatter { } } + /// Output a pre-built JSON error record on stderr. + pub fn json_error(&self, value: &T) { + match serde_json::to_string_pretty(value) { + Ok(json) => eprintln!("{json}"), + Err(e) => eprintln!("Error serializing output: {e}"), + } + } + /// Print a line of text (respects quiet mode) pub fn println(&self, message: &str) { if self.config.quiet { diff --git a/crates/cli/src/output/mod.rs b/crates/cli/src/output/mod.rs index d455095..2df659c 100644 --- a/crates/cli/src/output/mod.rs +++ b/crates/cli/src/output/mod.rs @@ -5,6 +5,7 @@ mod formatter; mod progress; +mod v3; // These exports will be used in Phase 2+ when commands are implemented #[allow(unused_imports)] @@ -13,6 +14,7 @@ pub use formatter::Formatter; pub use formatter::Theme; #[allow(unused_imports)] pub use progress::ProgressBar; +pub use v3::{V3ErrorEnvelope, V3PartialErrorEnvelope, V3SuccessEnvelope}; /// Output configuration derived from CLI flags #[derive(Debug, Clone, Default)] diff --git a/crates/cli/src/output/v3.rs b/crates/cli/src/output/v3.rs new file mode 100644 index 0000000..5ef3f2e --- /dev/null +++ b/crates/cli/src/output/v3.rs @@ -0,0 +1,195 @@ +//! Shared JSON output v3 envelopes for newly versioned command families. + +use serde::Serialize; + +use crate::exit_code::ExitCode; + +const VERSIONED_OBJECTS_FAMILY: &str = "versioned_objects"; + +#[derive(Debug, Serialize)] +pub struct V3SuccessEnvelope { + schema_version: u8, + #[serde(rename = "type")] + family: &'static str, + status: &'static str, + data: T, +} + +impl V3SuccessEnvelope { + pub fn versioned_objects(data: T) -> Self { + Self { + schema_version: 3, + family: VERSIONED_OBJECTS_FAMILY, + status: "success", + data, + } + } +} + +#[derive(Debug, Serialize)] +pub struct V3ErrorEnvelope { + schema_version: u8, + #[serde(rename = "type")] + family: &'static str, + status: &'static str, + error: V3ErrorDetail, +} + +impl V3ErrorEnvelope { + pub fn versioned_objects( + code: ExitCode, + message: impl Into, + capability: Option<&str>, + ) -> Self { + Self { + schema_version: 3, + family: VERSIONED_OBJECTS_FAMILY, + status: "error", + error: V3ErrorDetail::from_exit_code(code, message.into(), capability), + } + } +} + +#[derive(Debug, Serialize)] +pub struct V3PartialErrorEnvelope { + schema_version: u8, + #[serde(rename = "type")] + family: &'static str, + status: &'static str, + error: V3ErrorDetail, + data: T, +} + +impl V3PartialErrorEnvelope { + pub fn versioned_objects( + code: ExitCode, + message: impl Into, + capability: Option<&str>, + data: T, + ) -> Self { + Self { + schema_version: 3, + family: VERSIONED_OBJECTS_FAMILY, + status: "error", + error: V3ErrorDetail::from_exit_code(code, message.into(), capability), + data, + } + } +} + +#[derive(Debug, Serialize)] +#[serde(untagged)] +enum V3ErrorDetail { + Standard(V3StandardError), + Unsupported(V3UnsupportedError), +} + +impl V3ErrorDetail { + fn from_exit_code(code: ExitCode, message: String, capability: Option<&str>) -> Self { + let (error_type, retryable, suggestion) = match code { + ExitCode::UnsupportedFeature => { + return Self::Unsupported(V3UnsupportedError { + error_type: "unsupported_feature", + message, + retryable: false, + capability: capability.unwrap_or("versioned_objects").to_string(), + server: None, + suggestion: Some( + "Verify that the target RustFS version supports this operation." + .to_string(), + ), + }); + } + ExitCode::Success => ("general_error", false, None), + ExitCode::GeneralError => ("general_error", false, None), + ExitCode::UsageError => ( + "usage_error", + false, + Some("Review the command arguments and retry.".to_string()), + ), + ExitCode::NetworkError => ( + "network_error", + true, + Some("Verify the endpoint and network connectivity, then retry.".to_string()), + ), + ExitCode::AuthError => ( + "auth_error", + false, + Some("Verify the alias credentials and permissions, then retry.".to_string()), + ), + ExitCode::NotFound => ( + "not_found", + false, + Some("Check the bucket, object key, and version ID, then retry.".to_string()), + ), + ExitCode::Conflict => ( + "conflict", + false, + Some("Review the version state or retention policy, then retry.".to_string()), + ), + ExitCode::Interrupted => ( + "interrupted", + true, + Some("Retry if the operation still needs to complete.".to_string()), + ), + }; + + Self::Standard(V3StandardError { + error_type, + message, + retryable, + suggestion, + }) + } +} + +#[derive(Debug, Serialize)] +struct V3StandardError { + #[serde(rename = "type")] + error_type: &'static str, + message: String, + retryable: bool, + suggestion: Option, +} + +#[derive(Debug, Serialize)] +struct V3UnsupportedError { + #[serde(rename = "type")] + error_type: &'static str, + message: String, + retryable: bool, + capability: String, + server: Option, + suggestion: Option, +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn versioned_success_uses_the_v3_envelope() { + let value = serde_json::to_value(V3SuccessEnvelope::versioned_objects( + serde_json::json!({ "operation": "copy" }), + )) + .expect("serialize v3 success envelope"); + + assert_eq!(value["schema_version"], 3); + assert_eq!(value["type"], "versioned_objects"); + assert_eq!(value["status"], "success"); + } + + #[test] + fn versioned_errors_use_stable_error_kinds() { + let value = serde_json::to_value(V3ErrorEnvelope::versioned_objects( + ExitCode::AuthError, + "Access denied", + None, + )) + .expect("serialize v3 error envelope"); + + assert_eq!(value["status"], "error"); + assert_eq!(value["error"]["type"], "auth_error"); + assert_eq!(value["error"]["retryable"], false); + } +} diff --git a/crates/cli/tests/fixtures/output_v3/version_operations/dry_run.json b/crates/cli/tests/fixtures/output_v3/version_operations/dry_run.json new file mode 100644 index 0000000..395de99 --- /dev/null +++ b/crates/cli/tests/fixtures/output_v3/version_operations/dry_run.json @@ -0,0 +1 @@ +{"schema_version":3,"type":"versioned_objects","status":"success","data":{"operation":"remove","outcome":"planned","dry_run":true,"planned":[{"path":"local/photos/image.jpg","version_id":"v1","delete_marker":false}],"removed":[],"failed":[],"summary":{"planned":1,"removed":0,"failed":0}}} diff --git a/crates/cli/tests/fixtures/output_v3/version_operations/empty.json b/crates/cli/tests/fixtures/output_v3/version_operations/empty.json new file mode 100644 index 0000000..880628a --- /dev/null +++ b/crates/cli/tests/fixtures/output_v3/version_operations/empty.json @@ -0,0 +1 @@ +{"schema_version":3,"type":"versioned_objects","status":"success","data":{"operation":"remove","outcome":"empty","dry_run":false,"planned":[],"removed":[],"failed":[],"summary":{"planned":0,"removed":0,"failed":0}}} diff --git a/crates/cli/tests/fixtures/output_v3/version_operations/error.json b/crates/cli/tests/fixtures/output_v3/version_operations/error.json new file mode 100644 index 0000000..275c1a8 --- /dev/null +++ b/crates/cli/tests/fixtures/output_v3/version_operations/error.json @@ -0,0 +1 @@ +{"schema_version":3,"type":"versioned_objects","status":"error","error":{"type":"conflict","message":"One object version could not be removed","retryable":false,"suggestion":null},"data":{"operation":"remove","outcome":"partial","dry_run":false,"planned":[],"removed":[{"path":"local/photos/image.jpg","version_id":"v1","delete_marker":false}],"failed":[{"path":"local/photos/image.jpg","version_id":"v2","error_type":"conflict","message":"Governance retention denied deletion"}],"summary":{"planned":0,"removed":1,"failed":1}}} diff --git a/crates/cli/tests/fixtures/output_v3/version_operations/stat.json b/crates/cli/tests/fixtures/output_v3/version_operations/stat.json new file mode 100644 index 0000000..7dfd6f7 --- /dev/null +++ b/crates/cli/tests/fixtures/output_v3/version_operations/stat.json @@ -0,0 +1 @@ +{"schema_version":3,"type":"versioned_objects","status":"success","data":{"operation":"stat","object":{"path":"local/photos/image.jpg","bucket":"photos","key":"image.jpg","version_id":"v1","delete_marker":false,"last_modified":"2026-07-21T04:00:00Z","size_bytes":1024,"size_human":"1 KiB","etag":"abc123","content_type":"image/jpeg","storage_class":"STANDARD","metadata":{"owner":"media"}}}} diff --git a/crates/cli/tests/fixtures/output_v3/version_operations/success.json b/crates/cli/tests/fixtures/output_v3/version_operations/success.json new file mode 100644 index 0000000..dbf0379 --- /dev/null +++ b/crates/cli/tests/fixtures/output_v3/version_operations/success.json @@ -0,0 +1 @@ +{"schema_version":3,"type":"versioned_objects","status":"success","data":{"operation":"copy","source":"local/photos/2026/image.jpg","target":"backup/photos/2026/image.jpg","source_version_id":"source-v1","version_id":"destination-v2","size_bytes":1024,"size_human":"1 KiB"}} diff --git a/crates/cli/tests/help_contract.rs b/crates/cli/tests/help_contract.rs index 1a0db13..ddf3c72 100644 --- a/crates/cli/tests/help_contract.rs +++ b/crates/cli/tests/help_contract.rs @@ -242,6 +242,7 @@ fn top_level_command_help_contract() { "--dry-run", "--incomplete", "--versions", + "--version-id", "--bypass", "--purge", "Examples:", @@ -323,6 +324,7 @@ fn top_level_command_help_contract() { "--dry-run", "--incomplete", "--versions", + "--version-id", "--bypass", "Examples:", "rc rm local/my-bucket/reports/ --recursive --dry-run", diff --git a/crates/cli/tests/output_schema_v3.rs b/crates/cli/tests/output_schema_v3.rs index 8572775..5e172ce 100644 --- a/crates/cli/tests/output_schema_v3.rs +++ b/crates/cli/tests/output_schema_v3.rs @@ -65,6 +65,12 @@ fn fixture_path(family: &str, case: &str) -> PathBuf { .join(format!("{case}.{extension}")) } +fn version_operation_fixture_path(case: &str) -> PathBuf { + Path::new(env!("CARGO_MANIFEST_DIR")) + .join("tests/fixtures/output_v3/version_operations") + .join(format!("{case}.json")) +} + fn snapshot_payload(contents: &str) -> Option<&str> { contents .split_once("\n---\n") @@ -124,6 +130,61 @@ fn every_v3_family_has_valid_success_empty_and_error_fixtures() { } } +#[test] +fn version_operation_success_empty_error_and_dry_run_fixtures_are_valid() { + let validator = load_validator(3); + + for case in ["success", "empty", "error", "dry_run", "stat"] { + let path = version_operation_fixture_path(case); + let value = load_json(&path); + assert_valid(&validator, &value, &path.display().to_string()); + } +} + +#[test] +fn version_remove_contract_requires_version_aware_partial_results() { + let validator = load_validator(3); + let fixture = load_json(&version_operation_fixture_path("error")); + + assert_eq!(fixture["status"], "error"); + assert_eq!(fixture["data"]["outcome"], "partial"); + assert_eq!(fixture["data"]["removed"][0]["version_id"], "v1"); + assert_eq!(fixture["data"]["failed"][0]["version_id"], "v2"); + + let mut missing_version = fixture; + missing_version["data"]["failed"][0] + .as_object_mut() + .expect("failure must be an object") + .remove("version_id"); + assert!( + !validator.is_valid(&missing_version), + "version removal failures must preserve the selected version ID field" + ); +} + +#[test] +fn version_dry_run_contract_distinguishes_planned_from_removed_items() { + let validator = load_validator(3); + let fixture = load_json(&version_operation_fixture_path("dry_run")); + + assert_eq!(fixture["data"]["dry_run"], true); + assert_eq!(fixture["data"]["outcome"], "planned"); + assert_eq!(fixture["data"]["planned"].as_array().map(Vec::len), Some(1)); + assert_eq!(fixture["data"]["removed"].as_array().map(Vec::len), Some(0)); + assert_valid(&validator, &fixture, "version removal dry-run fixture"); + + let mut contradictory = fixture; + contradictory["data"]["removed"] = serde_json::json!([{ + "path": "local/photos/image.jpg", + "version_id": "v1", + "delete_marker": false + }]); + assert!( + !validator.is_valid(&contradictory), + "dry-run records must never claim that an object version was removed" + ); +} + #[test] fn legacy_schemas_compile_and_existing_v1_golden_snapshots_remain_valid() { let v1_validator = load_validator(1); diff --git a/crates/cli/tests/versioned_objects.rs b/crates/cli/tests/versioned_objects.rs new file mode 100644 index 0000000..6818ea0 --- /dev/null +++ b/crates/cli/tests/versioned_objects.rs @@ -0,0 +1,763 @@ +use std::io::{ErrorKind, Read, Write}; +use std::net::{TcpListener, TcpStream}; +use std::path::PathBuf; +use std::process::{Command, Output}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Mutex}; +use std::thread::{self, JoinHandle}; +use std::time::Duration; + +#[derive(Debug, Clone)] +struct CapturedRequest { + method: String, + path: String, + headers: Vec<(String, String)>, + body: String, +} + +impl CapturedRequest { + fn header(&self, name: &str) -> Option<&str> { + self.headers + .iter() + .find(|(key, _)| key.eq_ignore_ascii_case(name)) + .map(|(_, value)| value.as_str()) + } +} + +#[derive(Debug, Clone, Copy)] +enum ResponseMode { + Read, + Stat, + Delete, + MissingVersion, + DeleteMarker, + AccessDenied, + GovernanceDenied, + RecursiveVersions, + PartialRecursiveVersions, +} + +struct TestServer { + authority: String, + requests: Arc>>, + stop: Arc, + handle: Option>, +} + +impl TestServer { + fn start(mode: ResponseMode) -> Self { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind local HTTP server"); + listener + .set_nonblocking(true) + .expect("set listener nonblocking"); + let authority = listener + .local_addr() + .expect("local HTTP server address") + .to_string(); + let requests = Arc::new(Mutex::new(Vec::new())); + let stop = Arc::new(AtomicBool::new(false)); + let thread_requests = Arc::clone(&requests); + let thread_stop = Arc::clone(&stop); + + let handle = thread::spawn(move || { + while !thread_stop.load(Ordering::SeqCst) { + match listener.accept() { + Ok((mut stream, _)) => { + if let Some(request) = read_request(&mut stream) { + let response = response_for(mode, &request); + thread_requests + .lock() + .expect("record request") + .push(request); + let _ = stream.write_all(response.as_bytes()); + } + } + Err(error) if error.kind() == ErrorKind::WouldBlock => { + thread::sleep(Duration::from_millis(5)); + } + Err(error) if error.kind() == ErrorKind::Interrupted => {} + Err(error) => panic!("accept test request: {error}"), + } + } + }); + + Self { + authority, + requests, + stop, + handle: Some(handle), + } + } + + fn endpoint_with_credentials(&self) -> String { + format!("http://accesskey:secretkey@{}", self.authority) + } + + fn captured_requests(&self) -> Vec { + self.requests.lock().expect("captured requests").clone() + } +} + +impl Drop for TestServer { + fn drop(&mut self) { + self.stop.store(true, Ordering::SeqCst); + if let Some(handle) = self.handle.take() { + let _ = handle.join(); + } + } +} + +fn rc_binary() -> PathBuf { + if let Ok(path) = std::env::var("CARGO_BIN_EXE_rc") { + return PathBuf::from(path); + } + + let workspace_root = PathBuf::from(env!("CARGO_MANIFEST_DIR")) + .parent() + .expect("cli crate has parent directory") + .parent() + .expect("workspace root exists") + .to_path_buf(); + let binary_name = format!("rc{}", std::env::consts::EXE_SUFFIX); + let debug_binary = workspace_root.join("target/debug").join(&binary_name); + if debug_binary.exists() { + return debug_binary; + } + workspace_root.join("target/release").join(binary_name) +} + +fn run_rc(server: Option<&TestServer>, args: &[&str]) -> Output { + let config_dir = tempfile::tempdir().expect("create config dir"); + let mut command = Command::new(rc_binary()); + command + .args(args) + .env("AWS_EC2_METADATA_DISABLED", "true") + .env("RC_CONFIG_DIR", config_dir.path()); + if let Some(server) = server { + command.env("RC_HOST_test", server.endpoint_with_credentials()); + } + command.output().expect("run rc command") +} + +fn read_request(stream: &mut TcpStream) -> Option { + stream + .set_read_timeout(Some(Duration::from_secs(2))) + .expect("set request read timeout"); + let mut buffer = Vec::new(); + let mut chunk = [0_u8; 2048]; + let header_end = loop { + match stream.read(&mut chunk) { + Ok(0) => return None, + Ok(read) => { + buffer.extend_from_slice(&chunk[..read]); + if let Some(position) = buffer.windows(4).position(|window| window == b"\r\n\r\n") { + break position + 4; + } + } + Err(error) if matches!(error.kind(), ErrorKind::WouldBlock | ErrorKind::TimedOut) => { + return None; + } + Err(_) => return None, + } + }; + + let headers_text = String::from_utf8_lossy(&buffer[..header_end]); + let mut lines = headers_text.split("\r\n"); + let request_line = lines.next()?; + let mut parts = request_line.split_whitespace(); + let method = parts.next()?.to_string(); + let path = parts.next()?.to_string(); + let headers: Vec<(String, String)> = lines + .take_while(|line| !line.is_empty()) + .filter_map(|line| line.split_once(':')) + .map(|(key, value)| (key.to_ascii_lowercase(), value.trim().to_string())) + .filter(|(key, _)| key != "authorization" && key != "x-amz-security-token") + .collect(); + let content_length = headers + .iter() + .find(|(key, _)| key == "content-length") + .and_then(|(_, value)| value.parse::().ok()) + .unwrap_or_default(); + let chunked = headers + .iter() + .any(|(key, value)| key == "transfer-encoding" && value.eq_ignore_ascii_case("chunked")); + if headers + .iter() + .any(|(key, value)| key == "expect" && value.eq_ignore_ascii_case("100-continue")) + { + stream.write_all(b"HTTP/1.1 100 Continue\r\n\r\n").ok()?; + } + if chunked { + while !buffer[header_end..] + .windows(5) + .any(|window| window == b"0\r\n\r\n") + { + let read = stream.read(&mut chunk).ok()?; + if read == 0 { + break; + } + buffer.extend_from_slice(&chunk[..read]); + } + } else { + while buffer.len().saturating_sub(header_end) < content_length { + let read = stream.read(&mut chunk).ok()?; + if read == 0 { + break; + } + buffer.extend_from_slice(&chunk[..read]); + } + } + let body_end = if chunked { + buffer.len() + } else { + (header_end + content_length).min(buffer.len()) + }; + let body = String::from_utf8_lossy(&buffer[header_end..body_end]).into_owned(); + + Some(CapturedRequest { + method, + path, + headers, + body, + }) +} + +fn response_for(mode: ResponseMode, request: &CapturedRequest) -> String { + match mode { + ResponseMode::Read => http_response( + 200, + &[("Content-Type", "text/plain"), ("x-amz-version-id", "v1")], + "old", + ), + ResponseMode::Stat => http_response( + 200, + &[ + ("Content-Type", "text/plain"), + ("x-amz-version-id", "v1"), + ("ETag", "\"etag-v1\""), + ], + "", + ), + ResponseMode::Delete => http_response( + 204, + &[("x-amz-version-id", "v1"), ("x-amz-delete-marker", "false")], + "", + ), + ResponseMode::MissingVersion => s3_error(404, "NoSuchVersion", "version missing"), + ResponseMode::DeleteMarker => http_response( + 405, + &[ + ("x-amz-error-code", "MethodNotAllowed"), + ("x-amz-version-id", "marker-v1"), + ("x-amz-delete-marker", "true"), + ], + "MethodNotAlloweddelete marker", + ), + ResponseMode::AccessDenied => s3_error(403, "AccessDenied", "policy denied"), + ResponseMode::GovernanceDenied => { + s3_error(403, "AccessDenied", "governance retention is active") + } + ResponseMode::RecursiveVersions => recursive_versions_response(request), + ResponseMode::PartialRecursiveVersions => partial_recursive_versions_response(request), + } +} + +fn recursive_versions_response(request: &CapturedRequest) -> String { + if request.method == "GET" && request.path.contains("versions") { + if request.path.contains("key-marker") { + return http_response( + 200, + &[("Content-Type", "application/xml")], + second_version_page(), + ); + } + return http_response( + 200, + &[("Content-Type", "application/xml")], + first_version_page(), + ); + } + if request.method == "POST" && request.path.contains("delete") { + return http_response( + 200, + &[("Content-Type", "application/xml")], + delete_versions_result(), + ); + } + s3_error(500, "UnexpectedRequest", "unexpected version test request") +} + +fn partial_recursive_versions_response(request: &CapturedRequest) -> String { + if request.method == "GET" && request.path.contains("versions") { + if request.path.contains("key-marker") { + return http_response( + 200, + &[("Content-Type", "application/xml")], + second_version_page(), + ); + } + return http_response( + 200, + &[("Content-Type", "application/xml")], + first_version_page(), + ); + } + if request.method == "POST" && request.path.contains("delete") { + return http_response( + 200, + &[("Content-Type", "application/xml")], + partial_delete_versions_result(), + ); + } + s3_error(500, "UnexpectedRequest", "unexpected version test request") +} + +fn http_response(status: u16, headers: &[(&str, &str)], body: &str) -> String { + let reason = match status { + 200 => "OK", + 204 => "No Content", + 401 => "Unauthorized", + 403 => "Forbidden", + 404 => "Not Found", + 405 => "Method Not Allowed", + _ => "Internal Server Error", + }; + let mut response = format!("HTTP/1.1 {status} {reason}\r\n"); + for (name, value) in headers { + response.push_str(&format!("{name}: {value}\r\n")); + } + response.push_str(&format!( + "Content-Length: {}\r\nConnection: close\r\n\r\n{body}", + body.len() + )); + response +} + +fn s3_error(status: u16, code: &str, message: &str) -> String { + let body = format!("{code}{message}"); + http_response( + status, + &[ + ("Content-Type", "application/xml"), + ("x-amz-error-code", code), + ], + &body, + ) +} + +fn first_version_page() -> &'static str { + r#" + + bucketlogs/ + logs/a.txtmarker-v1 + 1000true + logs/a.txtv1false2026-07-20T00:00:00.000Z"etag-a"3STANDARD + logs/a.txtmarker-v1true2026-07-21T00:00:00.000Z +"# +} + +fn second_version_page() -> &'static str { + r#" + + bucketlogs/logs/a.txtmarker-v1 + 1000false + logs/b.txtv2true2026-07-21T00:00:00.000Z"etag-b"3STANDARD +"# +} + +fn delete_versions_result() -> &'static str { + r#" + + logs/a.txtv1 + logs/a.txtmarker-v1truemarker-v1 + logs/b.txtv2 +"# +} + +fn partial_delete_versions_result() -> &'static str { + r#" + + logs/a.txtv1 + logs/a.txtmarker-v1AccessDeniedgovernance retention is active + logs/b.txtv2 +"# +} + +#[test] +fn version_aware_commands_select_exact_versions_successfully() { + let cat_server = TestServer::start(ResponseMode::Read); + let cat = run_rc( + Some(&cat_server), + &["cat", "test/bucket/key.txt", "--version-id", "v1"], + ); + assert!( + cat.status.success(), + "stderr: {}", + String::from_utf8_lossy(&cat.stderr) + ); + assert_eq!(cat.stdout, b"old"); + assert!( + cat_server.captured_requests()[0] + .path + .contains("versionId=v1") + ); + + let head_server = TestServer::start(ResponseMode::Read); + let head = run_rc( + Some(&head_server), + &[ + "head", + "test/bucket/key.txt", + "--bytes", + "3", + "--version-id", + "v1", + ], + ); + assert!( + head.status.success(), + "stderr: {}", + String::from_utf8_lossy(&head.stderr) + ); + assert_eq!(head.stdout, b"old"); + assert!( + head_server.captured_requests()[0] + .path + .contains("versionId=v1") + ); + + let stat_server = TestServer::start(ResponseMode::Stat); + let stat = run_rc( + Some(&stat_server), + &[ + "--json", + "stat", + "test/bucket/key.txt", + "--version-id", + "v1", + ], + ); + assert!( + stat.status.success(), + "stderr: {}", + String::from_utf8_lossy(&stat.stderr) + ); + let stat_json: serde_json::Value = + serde_json::from_slice(&stat.stdout).expect("stat JSON output"); + assert_eq!(stat_json["schema_version"], 3); + assert_eq!(stat_json["type"], "versioned_objects"); + assert_eq!(stat_json["data"]["operation"], "stat"); + assert_eq!(stat_json["data"]["object"]["version_id"], "v1"); + assert!( + stat_server.captured_requests()[0] + .path + .contains("versionId=v1") + ); + + let rm_server = TestServer::start(ResponseMode::Delete); + let rm = run_rc( + Some(&rm_server), + &[ + "--json", + "rm", + "test/bucket/key.txt", + "--version-id", + "v1", + "--bypass", + ], + ); + assert!( + rm.status.success(), + "stderr: {}", + String::from_utf8_lossy(&rm.stderr) + ); + let rm_json: serde_json::Value = serde_json::from_slice(&rm.stdout).expect("rm JSON output"); + assert_eq!(rm_json["schema_version"], 3); + assert_eq!(rm_json["data"]["operation"], "remove"); + assert_eq!(rm_json["data"]["removed"][0]["version_id"], "v1"); + let request = &rm_server.captured_requests()[0]; + assert!(request.path.contains("versionId=v1")); + assert_eq!( + request.header("x-amz-bypass-governance-retention"), + Some("true") + ); +} + +#[test] +fn missing_versions_return_not_found_for_each_affected_command() { + let cases: &[&[&str]] = &[ + &["cat", "test/bucket/key.txt", "--version-id", "missing"], + &["head", "test/bucket/key.txt", "--version-id", "missing"], + &["stat", "test/bucket/key.txt", "--version-id", "missing"], + &["rm", "test/bucket/key.txt", "--version-id", "missing"], + ]; + + for args in cases { + let server = TestServer::start(ResponseMode::MissingVersion); + let output = run_rc(Some(&server), args); + assert_eq!( + output.status.code(), + Some(5), + "args={args:?}, stderr={}", + String::from_utf8_lossy(&output.stderr) + ); + } +} + +#[test] +fn read_commands_report_delete_markers_distinctly() { + let cases: &[&[&str]] = &[ + &["cat", "test/bucket/key.txt", "--version-id", "marker-v1"], + &["head", "test/bucket/key.txt", "--version-id", "marker-v1"], + &["stat", "test/bucket/key.txt", "--version-id", "marker-v1"], + ]; + + for args in cases { + let server = TestServer::start(ResponseMode::DeleteMarker); + let output = run_rc(Some(&server), args); + assert_eq!( + output.status.code(), + Some(5), + "args={args:?}, stderr={}", + String::from_utf8_lossy(&output.stderr) + ); + assert!( + String::from_utf8_lossy(&output.stderr).contains("delete marker"), + "args={args:?}, stderr={}", + String::from_utf8_lossy(&output.stderr) + ); + } +} + +#[test] +fn invalid_version_selectors_return_usage_for_each_affected_command() { + let cases: &[&[&str]] = &[ + &[ + "cat", + "test/bucket/key.txt", + "--version-id", + "v1", + "--rewind", + "1h", + ], + &["head", "test/bucket/key.txt", "--version-id", ""], + &[ + "stat", + "test/bucket/key.txt", + "--version-id", + "v1", + "--rewind", + "1h", + ], + &[ + "rm", + "test/bucket/key.txt", + "--version-id", + "v1", + "--versions", + ], + ]; + + for args in cases { + let output = run_rc(None, args); + assert_eq!( + output.status.code(), + Some(2), + "args={args:?}, stderr={}", + String::from_utf8_lossy(&output.stderr) + ); + } +} + +#[test] +fn version_selector_json_error_is_one_v3_record_on_stderr() { + let output = run_rc( + None, + &[ + "--json", + "rm", + "test/bucket/key.txt", + "--version-id", + "v1", + "--versions", + ], + ); + + assert_eq!(output.status.code(), Some(2)); + assert!(output.stdout.is_empty(), "JSON errors belong on stderr"); + let payload: serde_json::Value = + serde_json::from_slice(&output.stderr).expect("stderr must be one JSON record"); + assert_eq!(payload["schema_version"], 3); + assert_eq!(payload["type"], "versioned_objects"); + assert_eq!(payload["status"], "error"); + assert_eq!(payload["error"]["type"], "usage_error"); +} + +#[test] +fn versioned_stat_json_error_is_one_v3_record_on_stderr() { + let server = TestServer::start(ResponseMode::MissingVersion); + let output = run_rc( + Some(&server), + &[ + "--json", + "stat", + "test/bucket/key.txt", + "--version-id", + "missing", + ], + ); + + assert_eq!(output.status.code(), Some(5)); + assert!(output.stdout.is_empty(), "JSON errors belong on stderr"); + let payload: serde_json::Value = + serde_json::from_slice(&output.stderr).expect("stderr must be one JSON record"); + assert_eq!(payload["schema_version"], 3); + assert_eq!(payload["type"], "versioned_objects"); + assert_eq!(payload["status"], "error"); + assert_eq!(payload["error"]["type"], "not_found"); +} + +#[test] +fn access_and_governance_denials_have_distinct_exit_codes() { + let access_server = TestServer::start(ResponseMode::AccessDenied); + let access = run_rc( + Some(&access_server), + &["rm", "test/bucket/key.txt", "--version-id", "v1"], + ); + assert_eq!(access.status.code(), Some(4)); + + let governance_server = TestServer::start(ResponseMode::GovernanceDenied); + let governance = run_rc( + Some(&governance_server), + &["rm", "test/bucket/key.txt", "--version-id", "v1"], + ); + assert_eq!(governance.status.code(), Some(6)); +} + +#[test] +fn recursive_version_removal_paginates_and_deletes_markers() { + let server = TestServer::start(ResponseMode::RecursiveVersions); + let output = run_rc( + Some(&server), + &[ + "--json", + "rm", + "test/bucket/logs/", + "--recursive", + "--versions", + "--bypass", + ], + ); + assert!( + output.status.success(), + "stdout: {}\nstderr: {}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + let payload: serde_json::Value = + serde_json::from_slice(&output.stdout).expect("recursive rm JSON output"); + assert_eq!(payload["schema_version"], 3); + assert_eq!(payload["data"]["summary"]["removed"], 3); + assert_eq!(payload["data"]["removed"].as_array().map(Vec::len), Some(3)); + assert!( + payload["data"]["removed"] + .as_array() + .expect("removed versions") + .iter() + .any(|entry| entry["delete_marker"] == true) + ); + + let requests = server.captured_requests(); + let list_requests: Vec<_> = requests + .iter() + .filter(|request| request.method == "GET" && request.path.contains("versions")) + .collect(); + assert_eq!(list_requests.len(), 2, "requests: {requests:#?}"); + assert!(list_requests[1].path.contains("key-marker")); + assert!(list_requests[1].path.contains("version-id-marker")); + let delete_request = requests + .iter() + .find(|request| request.method == "POST" && request.path.contains("delete")) + .expect("delete versions request"); + assert!(delete_request.body.contains("v1")); + assert!( + delete_request + .body + .contains("marker-v1") + ); + assert_eq!( + delete_request.header("x-amz-bypass-governance-retention"), + Some("true") + ); +} + +#[test] +fn partial_version_removal_preserves_successes_and_version_aware_failures() { + let server = TestServer::start(ResponseMode::PartialRecursiveVersions); + let output = run_rc( + Some(&server), + &[ + "--json", + "rm", + "test/bucket/logs/", + "--recursive", + "--versions", + ], + ); + + assert_eq!( + output.status.code(), + Some(6), + "stdout: {}\nstderr: {}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + assert!(output.stdout.is_empty(), "partial errors belong on stderr"); + let payload: serde_json::Value = + serde_json::from_slice(&output.stderr).expect("partial error must be one JSON record"); + assert_eq!(payload["schema_version"], 3); + assert_eq!(payload["status"], "error"); + assert_eq!(payload["error"]["type"], "conflict"); + assert_eq!(payload["data"]["outcome"], "partial"); + assert_eq!(payload["data"]["summary"]["removed"], 2); + assert_eq!(payload["data"]["summary"]["failed"], 1); + assert_eq!(payload["data"]["removed"][0]["version_id"], "v1"); + assert_eq!(payload["data"]["failed"][0]["version_id"], "marker-v1"); +} + +#[test] +fn version_removal_json_dry_run_is_one_planned_v3_record() { + let server = TestServer::start(ResponseMode::RecursiveVersions); + let output = run_rc( + Some(&server), + &[ + "--json", + "rm", + "test/bucket/logs/", + "--recursive", + "--versions", + "--dry-run", + ], + ); + + assert!( + output.status.success(), + "stdout: {}\nstderr: {}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + let payload: serde_json::Value = + serde_json::from_slice(&output.stdout).expect("dry-run stdout must be one JSON record"); + assert_eq!(payload["schema_version"], 3); + assert_eq!(payload["data"]["outcome"], "planned"); + assert_eq!(payload["data"]["dry_run"], true); + assert_eq!(payload["data"]["summary"]["planned"], 3); + assert_eq!(payload["data"]["removed"].as_array().map(Vec::len), Some(0)); + assert!( + server + .captured_requests() + .iter() + .all(|request| request.method != "POST"), + "dry-run must not issue delete requests" + ); +} diff --git a/crates/core/src/error.rs b/crates/core/src/error.rs index 6da506d..3d7cfa9 100644 --- a/crates/core/src/error.rs +++ b/crates/core/src/error.rs @@ -54,6 +54,33 @@ pub enum Error { #[error("Not found: {0}")] NotFound(String), + /// A requested historical object version does not exist. + #[error("Object version '{version_id}' not found: {path}")] + VersionNotFound { + /// Object path. + path: String, + /// Requested version identifier. + version_id: String, + }, + + /// A read or metadata request selected a delete marker instead of object data. + #[error("Object version '{version_id}' is a delete marker: {path}")] + DeleteMarker { + /// Object path. + path: String, + /// Delete marker version identifier. + version_id: String, + }, + + /// Object Lock governance retention rejected a delete request. + #[error("Governance retention denied deletion: {path} (version: {version_id:?})")] + GovernanceDenied { + /// Object path. + path: String, + /// Requested version identifier, when one was supplied. + version_id: Option, + }, + /// Network error (retryable) #[error("Network error: {0}")] Network(String), @@ -75,14 +102,17 @@ impl Error { /// Get the appropriate exit code for this error pub const fn exit_code(&self) -> i32 { match self { - Error::InvalidPath(_) => 2, // UsageError - Error::Config(_) => 2, // UsageError - Error::Network(_) => 3, // NetworkError - Error::Auth(_) => 4, // AuthError - Error::NotFound(_) | Error::AliasNotFound(_) => 5, // NotFound - Error::Conflict(_) | Error::AliasExists(_) => 6, // Conflict - Error::UnsupportedFeature(_) => 7, // UnsupportedFeature - _ => 1, // GeneralError + Error::InvalidPath(_) => 2, // UsageError + Error::Config(_) => 2, // UsageError + Error::Network(_) => 3, // NetworkError + Error::Auth(_) => 4, // AuthError + Error::NotFound(_) + | Error::VersionNotFound { .. } + | Error::DeleteMarker { .. } + | Error::AliasNotFound(_) => 5, // NotFound + Error::Conflict(_) | Error::GovernanceDenied { .. } | Error::AliasExists(_) => 6, // Conflict + Error::UnsupportedFeature(_) => 7, // UnsupportedFeature + _ => 1, // GeneralError } } } @@ -113,4 +143,27 @@ mod tests { let err = Error::InvalidPath("/bad/path".into()); assert_eq!(err.to_string(), "Invalid path: /bad/path"); } + + #[test] + fn versioning_errors_are_distinct_and_stable() { + let missing = Error::VersionNotFound { + path: "local/bucket/object.txt".to_string(), + version_id: "v1".to_string(), + }; + let marker = Error::DeleteMarker { + path: "local/bucket/object.txt".to_string(), + version_id: "v2".to_string(), + }; + let governance = Error::GovernanceDenied { + path: "local/bucket/object.txt".to_string(), + version_id: Some("v3".to_string()), + }; + + assert_eq!(missing.exit_code(), 5); + assert!(missing.to_string().contains("version 'v1'")); + assert_eq!(marker.exit_code(), 5); + assert!(marker.to_string().contains("delete marker")); + assert_eq!(governance.exit_code(), 6); + assert!(governance.to_string().contains("Governance retention")); + } } diff --git a/crates/core/src/lib.rs b/crates/core/src/lib.rs index 906f091..502df68 100644 --- a/crates/core/src/lib.rs +++ b/crates/core/src/lib.rs @@ -47,6 +47,8 @@ pub use select::{ SelectSseCustomerOptions, }; pub use traits::{ - BucketNotification, Capabilities, ListOptions, ListResult, NotificationTarget, ObjectInfo, - ObjectStore, ObjectVersion, ObjectVersionListResult, + BucketNotification, Capabilities, DeleteObjectFailure, DeleteObjectsResult, + DeleteRequestOptions, DeletedObject, ListObjectVersionsOptions, ListOptions, ListResult, + NotificationTarget, ObjectInfo, ObjectReadOptions, ObjectStore, ObjectVersion, + ObjectVersionIdentifier, ObjectVersionListResult, }; diff --git a/crates/core/src/traits.rs b/crates/core/src/traits.rs index 34fec47..0f9d267 100644 --- a/crates/core/src/traits.rs +++ b/crates/core/src/traits.rs @@ -8,11 +8,11 @@ use std::collections::HashMap; use async_trait::async_trait; use jiff::Timestamp; use serde::{Deserialize, Serialize}; -use tokio::io::AsyncWrite; +use tokio::io::{AsyncWrite, AsyncWriteExt}; use crate::cors::CorsRule; use crate::encryption::{BucketEncryption, ObjectEncryptionRequest}; -use crate::error::Result; +use crate::error::{Error, Result}; use crate::lifecycle::LifecycleRule; use crate::path::RemotePath; use crate::replication::ReplicationConfiguration; @@ -64,6 +64,94 @@ pub struct ObjectVersionListResult { pub version_id_marker: Option, } +/// Options for selecting an object version during read and metadata operations. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct ObjectReadOptions { + /// Exact object version to select. `None` selects the current object. + pub version_id: Option, +} + +impl ObjectReadOptions { + /// Build read options while rejecting ambiguous empty version identifiers. + pub fn for_version(version_id: Option) -> Result { + if version_id.as_deref().is_some_and(str::is_empty) { + return Err(Error::InvalidPath("Version ID cannot be empty".to_string())); + } + Ok(Self { version_id }) + } +} + +/// Pagination options for listing object versions and delete markers. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct ListObjectVersionsOptions { + /// Maximum number of entries to return. + pub max_keys: Option, + /// Key marker returned by the previous page. + pub key_marker: Option, + /// Version marker returned by the previous page. + pub version_id_marker: Option, +} + +/// Request-level options for object deletion. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct DeleteRequestOptions { + /// Exact object version to delete. `None` targets the current object state. + pub version_id: Option, + /// Explicitly bypass Object Lock governance retention. + pub bypass_governance: bool, + /// Ask RustFS to permanently delete data instead of creating delete markers. + pub force_delete: bool, +} + +/// An object key and optional historical version selected for deletion. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct ObjectVersionIdentifier { + /// Object key. + pub key: String, + /// Exact version to delete, when version-aware deletion is requested. + #[serde(skip_serializing_if = "Option::is_none")] + pub version_id: Option, + /// Whether the selected version was listed as a delete marker. + pub is_delete_marker: bool, +} + +/// A version-aware delete result returned by the object store. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct DeletedObject { + /// Deleted object key. + pub key: String, + /// Deleted object version, when reported by the backend. + #[serde(skip_serializing_if = "Option::is_none")] + pub version_id: Option, + /// Whether the deleted entry is a delete marker. + pub is_delete_marker: bool, +} + +/// A per-object failure returned by a multi-object delete request. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct DeleteObjectFailure { + /// Object key that could not be deleted. + pub key: String, + /// Requested version, when present. + #[serde(skip_serializing_if = "Option::is_none")] + pub version_id: Option, + /// S3 error code, when provided. + #[serde(skip_serializing_if = "Option::is_none")] + pub code: Option, + /// Backend error message, when provided. + #[serde(skip_serializing_if = "Option::is_none")] + pub message: Option, +} + +/// Result of deleting multiple version-aware object identifiers. +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct DeleteObjectsResult { + /// Successfully deleted entries. + pub deleted: Vec, + /// Entries rejected by the backend. + pub failures: Vec, +} + /// Metadata for an object or bucket #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ObjectInfo { @@ -98,6 +186,18 @@ pub struct ObjectInfo { #[serde(skip_serializing_if = "Option::is_none")] pub metadata: Option>, + /// Object version selected or created by the operation. + #[serde(skip_serializing_if = "Option::is_none")] + pub version_id: Option, + + /// Source object version used by a copy operation. + #[serde(skip_serializing_if = "Option::is_none")] + pub source_version_id: Option, + + /// Whether the selected version is a delete marker. + #[serde(skip_serializing_if = "Option::is_none")] + pub is_delete_marker: Option, + /// Whether this is a directory/prefix pub is_dir: bool, } @@ -114,6 +214,9 @@ impl ObjectInfo { storage_class: None, content_type: None, metadata: None, + version_id: None, + source_version_id: None, + is_delete_marker: None, is_dir: false, } } @@ -129,6 +232,9 @@ impl ObjectInfo { storage_class: None, content_type: None, metadata: None, + version_id: None, + source_version_id: None, + is_delete_marker: None, is_dir: true, } } @@ -144,6 +250,9 @@ impl ObjectInfo { storage_class: None, content_type: None, metadata: None, + version_id: None, + source_version_id: None, + is_delete_marker: None, is_dir: true, } } @@ -262,6 +371,20 @@ pub trait ObjectStore: Send + Sync { /// Get object metadata async fn head_object(&self, path: &RemotePath) -> Result; + /// Get metadata for the current object or an exact historical version. + async fn head_object_with_options( + &self, + path: &RemotePath, + options: &ObjectReadOptions, + ) -> Result { + if options.version_id.is_some() { + return Err(Error::UnsupportedFeature( + "Exact-version metadata reads are not implemented by this object store".to_string(), + )); + } + self.head_object(path).await + } + /// Check if a bucket exists async fn bucket_exists(&self, bucket: &str) -> Result; @@ -277,6 +400,38 @@ pub trait ObjectStore: Send + Sync { /// Get object content as bytes async fn get_object(&self, path: &RemotePath) -> Result>; + /// Get current object content or an exact historical version as bytes. + async fn get_object_with_options( + &self, + path: &RemotePath, + options: &ObjectReadOptions, + ) -> Result> { + if options.version_id.is_some() { + return Err(Error::UnsupportedFeature( + "Exact-version object reads are not implemented by this object store".to_string(), + )); + } + self.get_object(path).await + } + + /// Stream current object content or an exact historical version to a writer. + async fn write_object_to_with_options( + &self, + path: &RemotePath, + options: &ObjectReadOptions, + writer: &mut (dyn AsyncWrite + Send + Unpin), + max_bytes: Option, + ) -> Result { + let data = self.get_object_with_options(path, options).await?; + let write_len = max_bytes + .and_then(|limit| usize::try_from(limit).ok()) + .map(|limit| limit.min(data.len())) + .unwrap_or(data.len()); + writer.write_all(&data[..write_len]).await?; + writer.flush().await?; + Ok(write_len as u64) + } + /// Upload object from bytes async fn put_object( &self, @@ -289,9 +444,41 @@ pub trait ObjectStore: Send + Sync { /// Delete an object async fn delete_object(&self, path: &RemotePath) -> Result<()>; + /// Delete the current object state or one exact version with explicit request options. + async fn delete_object_with_options( + &self, + path: &RemotePath, + options: DeleteRequestOptions, + ) -> Result { + if options.version_id.is_some() || options.bypass_governance || options.force_delete { + return Err(Error::UnsupportedFeature( + "Version-aware or policy-bypassing deletion is not implemented by this object store" + .to_string(), + )); + } + self.delete_object(path).await?; + Ok(DeletedObject { + key: path.key.clone(), + version_id: None, + is_delete_marker: false, + }) + } + /// Delete multiple objects (batch delete) async fn delete_objects(&self, bucket: &str, keys: Vec) -> Result>; + /// Delete exact object versions and delete markers in one request. + async fn delete_object_versions( + &self, + _bucket: &str, + _objects: Vec, + _options: DeleteRequestOptions, + ) -> Result { + Err(Error::UnsupportedFeature( + "Multi-object version deletion is not implemented by this object store".to_string(), + )) + } + /// Copy object within S3 (server-side copy) async fn copy_object( &self, @@ -336,6 +523,17 @@ pub trait ObjectStore: Send + Sync { max_keys: Option, ) -> Result>; + /// List one page of object versions and delete markers with both S3 pagination markers. + async fn list_object_versions_page_with_options( + &self, + _path: &RemotePath, + _options: &ListObjectVersionsOptions, + ) -> Result { + Err(Error::UnsupportedFeature( + "Paginated object version listing is not implemented by this object store".to_string(), + )) + } + /// Get object tags async fn get_object_tags( &self, @@ -451,6 +649,47 @@ mod tests { assert_eq!(info.key, "test.txt"); assert_eq!(info.size_bytes, Some(1024)); assert!(!info.is_dir); + assert_eq!(info.version_id, None); + assert_eq!(info.source_version_id, None); + assert_eq!(info.is_delete_marker, None); + } + + #[test] + fn object_read_options_reject_empty_version_ids() { + let error = ObjectReadOptions::for_version(Some(String::new())) + .expect_err("empty version IDs must be rejected"); + + assert!(matches!(error, crate::Error::InvalidPath(_))); + } + + #[test] + fn versioned_delete_targets_are_serializable_for_structured_output() { + let target = ObjectVersionIdentifier { + key: "reports/a.csv".to_string(), + version_id: Some("v1".to_string()), + is_delete_marker: true, + }; + + let json = serde_json::to_value(target).expect("serialize versioned delete target"); + assert_eq!(json["key"], "reports/a.csv"); + assert_eq!(json["version_id"], "v1"); + assert_eq!(json["is_delete_marker"], true); + } + + #[test] + fn object_info_version_fields_are_optional_and_serializable() { + let current = serde_json::to_value(ObjectInfo::file("current.txt", 1)) + .expect("serialize current object info"); + assert!(current.get("version_id").is_none()); + assert!(current.get("source_version_id").is_none()); + assert!(current.get("is_delete_marker").is_none()); + + let mut copied = ObjectInfo::file("copy.txt", 1); + copied.version_id = Some("destination-v2".to_string()); + copied.source_version_id = Some("source-v1".to_string()); + let copied = serde_json::to_value(copied).expect("serialize copy object info"); + assert_eq!(copied["version_id"], "destination-v2"); + assert_eq!(copied["source_version_id"], "source-v1"); } #[test] diff --git a/crates/s3/src/client.rs b/crates/s3/src/client.rs index 4bd14fc..9f6a386 100644 --- a/crates/s3/src/client.rs +++ b/crates/s3/src/client.rs @@ -27,11 +27,14 @@ use http_body::Frame; use http_body_util::StreamBody; use jiff::Timestamp; use quick_xml::de::from_str as from_xml_str; +pub use rc_core::DeleteRequestOptions; use rc_core::{ - Alias, BucketEncryption, BucketNotification, Capabilities, CorsRule, Error, LifecycleRule, - ListOptions, ListResult, NotificationTarget, ObjectEncryptionRequest, ObjectInfo, ObjectStore, - ObjectVersion, ObjectVersionListResult, RemotePath, ReplicationConfiguration, RequestHeader, - Result, SelectOptions, global_request_headers, + Alias, BucketEncryption, BucketNotification, Capabilities, CorsRule, DeleteObjectFailure, + DeleteObjectsResult, DeletedObject, Error, LifecycleRule, ListObjectVersionsOptions, + ListOptions, ListResult, NotificationTarget, ObjectEncryptionRequest, ObjectInfo, + ObjectReadOptions, ObjectStore, ObjectVersion, ObjectVersionIdentifier, + ObjectVersionListResult, RemotePath, ReplicationConfiguration, RequestHeader, Result, + SelectOptions, global_request_headers, }; use reqwest::Method; use reqwest::header::{CONTENT_TYPE, HeaderMap, HeaderName, HeaderValue}; @@ -856,13 +859,6 @@ pub struct S3Client { request_headers: Vec, } -/// Request-level options for delete operations. -#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] -pub struct DeleteRequestOptions { - /// Ask RustFS to permanently delete data instead of creating delete markers. - pub force_delete: bool, -} - #[derive(Debug, Clone)] struct CustomHeaderInterceptor { headers: Vec, @@ -1023,12 +1019,32 @@ impl S3Client { builder = builder.version_id_marker(version_id_marker); } - let response = builder.send().await.map_err(|e| { - let err_str = Self::format_sdk_error(&e); - if err_str.contains("NotFound") || err_str.contains("NoSuchBucket") { + let response = builder.send().await.map_err(|error| { + let formatted = Self::format_sdk_error(&error); + if let aws_sdk_s3::error::SdkError::ServiceError(service_error) = &error { + let status = service_error.raw().status().as_u16(); + let code = service_error.err().code(); + if matches!(status, 401 | 403) + || matches!( + code, + Some("AccessDenied") | Some("Forbidden") | Some("Unauthorized") + ) + { + return Error::Auth(formatted); + } + if status == 404 || matches!(code, Some("NotFound") | Some("NoSuchBucket")) { + return Error::NotFound(format!("Bucket not found: {}", path.bucket)); + } + } + if formatted.contains("AccessDenied") + || formatted.contains("Forbidden") + || formatted.contains("Unauthorized") + { + Error::Auth(formatted) + } else if formatted.contains("NotFound") || formatted.contains("NoSuchBucket") { Error::NotFound(format!("Bucket not found: {}", path.bucket)) } else { - Error::Network(err_str) + Error::Network(formatted) } })?; @@ -1080,23 +1096,37 @@ impl S3Client { pub async fn get_object_with_progress( &self, path: &RemotePath, - mut on_progress: impl FnMut(u64, Option) + Send, + on_progress: impl FnMut(u64, Option) + Send, ) -> Result> { - let response = self - .inner - .get_object() - .bucket(&path.bucket) - .key(&path.key) - .send() + self.get_object_with_progress_and_options(path, &ObjectReadOptions::default(), on_progress) .await - .map_err(|e| { - let err_str = e.to_string(); - if err_str.contains("NotFound") || err_str.contains("NoSuchKey") { - Error::NotFound(path.to_string()) - } else { - Error::Network(err_str) - } - })?; + } + + /// Download an exact object version and report downloaded bytes after each chunk. + pub async fn get_object_with_progress_and_options( + &self, + path: &RemotePath, + options: &ObjectReadOptions, + mut on_progress: impl FnMut(u64, Option) + Send, + ) -> Result> { + let mut request = self.inner.get_object().bucket(&path.bucket).key(&path.key); + if let Some(version_id) = &options.version_id { + request = request.version_id(version_id); + } + let response = request.send().await.map_err(|error| { + Self::map_object_request_error(&error, path, options.version_id.as_deref()) + })?; + + if response.delete_marker().unwrap_or(false) { + return Err(Error::DeleteMarker { + path: path.to_string(), + version_id: response + .version_id() + .or(options.version_id.as_deref()) + .unwrap_or("unknown") + .to_string(), + }); + } let content_length = response .content_length() @@ -1233,24 +1263,50 @@ impl S3Client { max_bytes: Option, ) -> Result where - W: AsyncWrite + Unpin + Send, + W: AsyncWrite + Unpin + Send + ?Sized, + { + self.write_object_to_with_options(path, &ObjectReadOptions::default(), writer, max_bytes) + .await + } + + /// Stream the current object or an exact historical version to a writer. + pub async fn write_object_to_with_options( + &self, + path: &RemotePath, + options: &ObjectReadOptions, + writer: &mut W, + max_bytes: Option, + ) -> Result + where + W: AsyncWrite + Unpin + Send + ?Sized, { if matches!(max_bytes, Some(0)) { + if options.version_id.is_some() { + self.head_object_with_options(path, options).await?; + } return Ok(0); } let mut request = self.inner.get_object().bucket(&path.bucket).key(&path.key); + if let Some(version_id) = &options.version_id { + request = request.version_id(version_id); + } if let Some(max_bytes) = max_bytes { request = request.range(format!("bytes=0-{}", max_bytes - 1)); } let response = request.send().await.map_err(|error| { - let message = error.to_string(); - if message.contains("NotFound") || message.contains("NoSuchKey") { - Error::NotFound(path.to_string()) - } else { - Error::Network(message) - } + Self::map_object_request_error(&error, path, options.version_id.as_deref()) })?; + if response.delete_marker().unwrap_or(false) { + return Err(Error::DeleteMarker { + path: path.to_string(), + version_id: response + .version_id() + .or(options.version_id.as_deref()) + .unwrap_or("unknown") + .to_string(), + }); + } let mut body = response.body; let mut bytes_written = 0u64; @@ -1305,6 +1361,7 @@ impl S3Client { info.etag = response .e_tag() .map(|etag| etag.trim_matches('"').to_string()); + info.version_id = response.version_id().map(ToString::to_string); info.last_modified = Some(jiff::Timestamp::now()); Ok(info) } @@ -1329,18 +1386,37 @@ impl S3Client { Ok(()) } - /// Delete an object with RustFS-specific request options. + /// Delete an object with version, governance, and RustFS force-delete options. + /// + /// This compatibility wrapper preserves the original unit result. Call + /// [`Self::delete_object_with_result`] when version-aware response fields are needed. pub async fn delete_object_with_options( &self, path: &RemotePath, options: DeleteRequestOptions, ) -> Result<()> { - let mut request = self + self.delete_object_with_result(path, options).await?; + Ok(()) + } + + /// Delete an object and preserve the returned version and delete-marker fields. + pub async fn delete_object_with_result( + &self, + path: &RemotePath, + options: DeleteRequestOptions, + ) -> Result { + let mut builder = self .inner .delete_object() .bucket(&path.bucket) - .key(&path.key) - .customize(); + .key(&path.key); + if let Some(version_id) = &options.version_id { + builder = builder.version_id(version_id); + } + if options.bypass_governance { + builder = builder.bypass_governance_retention(true); + } + let mut request = builder.customize(); if options.force_delete { request = request.mutate_request(|request| { @@ -1350,67 +1426,105 @@ impl S3Client { }); } - request.send().await.map_err(|e| { - let err_str = Self::format_sdk_error(&e); - let is_missing_key = if let aws_sdk_s3::error::SdkError::ServiceError(service_err) = &e - { - let code = service_err.err().code().or_else(|| { - service_err - .raw() - .headers() - .get("x-amz-error-code") - .and_then(|value| std::str::from_utf8(value.as_bytes()).ok()) - }); - matches!(code, Some("NoSuchKey") | Some("NotFound")) - || service_err.raw().status().as_u16() == 404 - } else { - err_str.contains("NotFound") || err_str.contains("NoSuchKey") - }; - - if is_missing_key { - Error::NotFound(path.to_string()) - } else { - Error::Network(err_str) - } + let response = request.send().await.map_err(|error| { + Self::map_object_request_error(&error, path, options.version_id.as_deref()) })?; - Ok(()) + Ok(DeletedObject { + key: path.key.clone(), + version_id: response + .version_id() + .or(options.version_id.as_deref()) + .map(ToString::to_string), + is_delete_marker: response.delete_marker().unwrap_or(false), + }) } - /// Delete multiple objects with RustFS-specific request options. + /// Delete multiple objects with governance and RustFS force-delete options. pub async fn delete_objects_with_options( &self, bucket: &str, keys: Vec, options: DeleteRequestOptions, ) -> Result> { - use aws_sdk_s3::types::{Delete, ObjectIdentifier}; - if keys.is_empty() { return Ok(vec![]); } + if options.version_id.is_some() { + return Err(Error::InvalidPath( + "Batch key deletion cannot apply one version ID to multiple objects".to_string(), + )); + } - let objects: Vec = - keys.iter() - .map(|key| { - ObjectIdentifier::builder().key(key).build().map_err(|e| { - Error::General(format!("invalid delete object identifier: {e}")) - }) - }) - .collect::>>()?; + let identifiers = keys + .into_iter() + .map(|key| ObjectVersionIdentifier { + key, + version_id: None, + is_delete_marker: false, + }) + .collect(); + let result = self + .delete_object_versions_with_options(bucket, identifiers, options) + .await?; + if !result.failures.is_empty() { + let error_keys: Vec<&str> = result + .failures + .iter() + .map(|failure| failure.key.as_str()) + .collect(); + tracing::warn!("Failed to delete some objects: {:?}", error_keys); + } + + Ok(result + .deleted + .into_iter() + .map(|deleted| deleted.key) + .collect()) + } + + /// Delete exact object versions and delete markers with optional governance bypass. + pub async fn delete_object_versions_with_options( + &self, + bucket: &str, + objects: Vec, + options: DeleteRequestOptions, + ) -> Result { + use aws_sdk_s3::types::{Delete, ObjectIdentifier}; + + if objects.is_empty() { + return Ok(DeleteObjectsResult::default()); + } + if options.version_id.is_some() { + return Err(Error::InvalidPath( + "Multi-object version deletion requires version IDs on each object identifier" + .to_string(), + )); + } + + let sdk_objects = objects + .iter() + .map(|object| { + let mut builder = ObjectIdentifier::builder().key(&object.key); + if let Some(version_id) = &object.version_id { + builder = builder.version_id(version_id); + } + builder.build().map_err(|error| { + Error::General(format!("invalid delete object identifier: {error}")) + }) + }) + .collect::>>()?; let delete = Delete::builder() - .set_objects(Some(objects)) + .set_objects(Some(sdk_objects)) .build() - .map_err(|e| Error::General(e.to_string()))?; - - let mut request = self - .inner - .delete_objects() - .bucket(bucket) - .delete(delete) - .customize(); + .map_err(|error| Error::General(error.to_string()))?; + let mut builder = self.inner.delete_objects().bucket(bucket).delete(delete); + if options.bypass_governance { + builder = builder.bypass_governance_retention(true); + } + let mut request = builder.customize(); if options.force_delete { request = request.mutate_request(|request| { request @@ -1419,27 +1533,43 @@ impl S3Client { }); } - let response = request - .send() - .await - .map_err(|e| Error::Network(e.to_string()))?; + let first = &objects[0]; + let error_path = RemotePath::new(&self.alias.name, bucket, &first.key); + let response = request.send().await.map_err(|error| { + Self::map_object_request_error(&error, &error_path, first.version_id.as_deref()) + })?; - let deleted: Vec = response + let deleted = response .deleted() .iter() - .filter_map(|d| d.key().map(|k| k.to_string())) + .filter_map(|entry| { + let key = entry.key()?.to_string(); + let version_id = entry + .version_id() + .or(entry.delete_marker_version_id()) + .map(ToString::to_string); + let requested_marker = objects.iter().any(|object| { + object.key == key && object.version_id == version_id && object.is_delete_marker + }); + Some(DeletedObject { + key, + version_id, + is_delete_marker: entry.delete_marker().unwrap_or(false) || requested_marker, + }) + }) + .collect(); + let failures = response + .errors() + .iter() + .map(|entry| DeleteObjectFailure { + key: entry.key().unwrap_or_default().to_string(), + version_id: entry.version_id().map(ToString::to_string), + code: entry.code().map(ToString::to_string), + message: entry.message().map(ToString::to_string), + }) .collect(); - if !response.errors().is_empty() { - let error_keys: Vec = response - .errors() - .iter() - .filter_map(|e| e.key().map(|k| k.to_string())) - .collect(); - tracing::warn!("Failed to delete some objects: {:?}", error_keys); - } - - Ok(deleted) + Ok(DeleteObjectsResult { deleted, failures }) } /// Format AWS SDK error into a detailed error message @@ -1471,6 +1601,110 @@ impl S3Client { } } + fn map_object_request_error( + error: &aws_sdk_s3::error::SdkError, + path: &RemotePath, + requested_version: Option<&str>, + ) -> Error + where + E: ProvideErrorMetadata + std::fmt::Display, + { + let formatted = Self::format_sdk_error(error); + + if let aws_sdk_s3::error::SdkError::ServiceError(service_error) = error { + let raw = service_error.raw(); + let code = service_error.err().code().or_else(|| { + raw.headers() + .get("x-amz-error-code") + .and_then(|value| std::str::from_utf8(value.as_bytes()).ok()) + }); + let status = raw.status().as_u16(); + if status == 401 || matches!(code, Some("Unauthorized")) { + return Error::Auth(formatted); + } + + let version_header = raw + .headers() + .get("x-amz-version-id") + .and_then(|value| std::str::from_utf8(value.as_bytes()).ok()); + let is_delete_marker = raw + .headers() + .get("x-amz-delete-marker") + .and_then(|value| std::str::from_utf8(value.as_bytes()).ok()) + .is_some_and(|value| value.eq_ignore_ascii_case("true")); + + let service_message = service_error.err().message().unwrap_or(formatted.as_str()); + if Self::is_retention_denial(service_message) { + return Error::GovernanceDenied { + path: path.to_string(), + version_id: requested_version.map(ToString::to_string), + }; + } + + if matches!(code, Some("AccessDenied") | Some("Forbidden")) || status == 403 { + return Error::Auth(formatted); + } + + if is_delete_marker { + return Error::DeleteMarker { + path: path.to_string(), + version_id: version_header + .or(requested_version) + .unwrap_or("unknown") + .to_string(), + }; + } + + if matches!(code, Some("NoSuchVersion")) + || (requested_version.is_some() + && (matches!(code, Some("NoSuchKey") | Some("NotFound")) || status == 404)) + { + return Error::VersionNotFound { + path: path.to_string(), + version_id: requested_version.unwrap_or("unknown").to_string(), + }; + } + + if matches!(code, Some("NoSuchKey") | Some("NotFound")) || status == 404 { + return Error::NotFound(path.to_string()); + } + } + + if Self::is_retention_denial(&formatted) { + return Error::GovernanceDenied { + path: path.to_string(), + version_id: requested_version.map(ToString::to_string), + }; + } + if formatted.contains("AccessDenied") + || formatted.contains("Forbidden") + || formatted.contains("Unauthorized") + { + return Error::Auth(formatted); + } + if formatted.contains("NoSuchVersion") + || (requested_version.is_some() + && (formatted.contains("NoSuchKey") || formatted.contains("NotFound"))) + { + return Error::VersionNotFound { + path: path.to_string(), + version_id: requested_version.unwrap_or("unknown").to_string(), + }; + } + if formatted.contains("NoSuchKey") || formatted.contains("NotFound") { + return Error::NotFound(path.to_string()); + } + Error::Network(formatted) + } + + fn is_retention_denial(message: &str) -> bool { + let normalized = message.to_ascii_lowercase(); + normalized.contains("governance") + || normalized.contains("retention") + || normalized.contains("object lock") + || normalized.contains("worm") + } + fn should_use_multipart(file_size: u64) -> bool { file_size > SINGLE_PUT_OBJECT_MAX_SIZE } @@ -1888,6 +2122,7 @@ impl S3Client { if let Some(etag) = response.e_tag() { info.etag = Some(etag.trim_matches('"').to_string()); } + info.version_id = response.version_id().map(ToString::to_string); info.last_modified = Some(jiff::Timestamp::now()); Ok(info) @@ -2058,6 +2293,7 @@ impl S3Client { if let Some(etag) = complete_response.e_tag() { info.etag = Some(etag.trim_matches('"').to_string()); } + info.version_id = complete_response.version_id().map(ToString::to_string); info.last_modified = Some(jiff::Timestamp::now()); Ok(info) @@ -2259,21 +2495,33 @@ impl ObjectStore for S3Client { } async fn head_object(&self, path: &RemotePath) -> Result { - let response = self - .inner - .head_object() - .bucket(&path.bucket) - .key(&path.key) - .send() + self.head_object_with_options(path, &ObjectReadOptions::default()) .await - .map_err(|e| { - let err_str = e.to_string(); - if err_str.contains("NotFound") || err_str.contains("NoSuchKey") { - Error::NotFound(path.to_string()) - } else { - Error::Network(err_str) - } - })?; + } + + async fn head_object_with_options( + &self, + path: &RemotePath, + options: &ObjectReadOptions, + ) -> Result { + let mut request = self.inner.head_object().bucket(&path.bucket).key(&path.key); + if let Some(version_id) = &options.version_id { + request = request.version_id(version_id); + } + let response = request.send().await.map_err(|error| { + Self::map_object_request_error(&error, path, options.version_id.as_deref()) + })?; + + if response.delete_marker().unwrap_or(false) { + return Err(Error::DeleteMarker { + path: path.to_string(), + version_id: response + .version_id() + .or(options.version_id.as_deref()) + .unwrap_or("unknown") + .to_string(), + }); + } let size = response.content_length().unwrap_or(0); let mut info = ObjectInfo::file(&path.key, size); @@ -2299,6 +2547,10 @@ impl ObjectStore for S3Client { { info.metadata = Some(meta.clone()); } + info.version_id = response + .version_id() + .or(options.version_id.as_deref()) + .map(ToString::to_string); Ok(info) } @@ -2374,6 +2626,25 @@ impl ObjectStore for S3Client { self.get_object_with_progress(path, |_, _| {}).await } + async fn get_object_with_options( + &self, + path: &RemotePath, + options: &ObjectReadOptions, + ) -> Result> { + self.get_object_with_progress_and_options(path, options, |_, _| {}) + .await + } + + async fn write_object_to_with_options( + &self, + path: &RemotePath, + options: &ObjectReadOptions, + writer: &mut (dyn AsyncWrite + Send + Unpin), + max_bytes: Option, + ) -> Result { + S3Client::write_object_to_with_options(self, path, options, writer, max_bytes).await + } + async fn put_object( &self, path: &RemotePath, @@ -2406,14 +2677,22 @@ impl ObjectStore for S3Client { if let Some(etag) = response.e_tag() { info.etag = Some(etag.trim_matches('"').to_string()); } + info.version_id = response.version_id().map(ToString::to_string); info.last_modified = Some(jiff::Timestamp::now()); Ok(info) } async fn delete_object(&self, path: &RemotePath) -> Result<()> { - self.delete_object_with_options(path, DeleteRequestOptions::default()) - .await + S3Client::delete_object_with_options(self, path, DeleteRequestOptions::default()).await + } + + async fn delete_object_with_options( + &self, + path: &RemotePath, + options: DeleteRequestOptions, + ) -> Result { + self.delete_object_with_result(path, options).await } async fn delete_objects(&self, bucket: &str, keys: Vec) -> Result> { @@ -2421,6 +2700,16 @@ impl ObjectStore for S3Client { .await } + async fn delete_object_versions( + &self, + bucket: &str, + objects: Vec, + options: DeleteRequestOptions, + ) -> Result { + self.delete_object_versions_with_options(bucket, objects, options) + .await + } + async fn copy_object( &self, src: &RemotePath, @@ -2454,6 +2743,10 @@ impl ObjectStore for S3Client { // Update etag from copy response if available let mut result = info; + if let Some(version_id) = response.version_id() { + result.version_id = Some(version_id.to_string()); + } + result.source_version_id = response.copy_source_version_id().map(ToString::to_string); if let Some(copy_result) = response.copy_object_result() && let Some(etag) = copy_result.e_tag() { @@ -2645,6 +2938,20 @@ impl ObjectStore for S3Client { Ok(self.list_object_versions_page(path, max_keys).await?.items) } + async fn list_object_versions_page_with_options( + &self, + path: &RemotePath, + options: &ListObjectVersionsOptions, + ) -> Result { + self.list_object_versions_page_with_markers( + path, + options.max_keys, + options.key_marker.as_deref(), + options.version_id_marker.as_deref(), + ) + .await + } + async fn get_object_tags( &self, path: &RemotePath, @@ -4473,13 +4780,61 @@ mod tests { let path = RemotePath::new("test", "bucket", "key.txt"); let _ = client - .delete_object_with_options(&path, DeleteRequestOptions { force_delete: true }) + .delete_object_with_options( + &path, + DeleteRequestOptions { + force_delete: true, + ..Default::default() + }, + ) .await; let request = request_receiver.expect_request(); assert_eq!(request.headers().get("x-rustfs-force-delete"), Some("true")); } + #[tokio::test] + async fn versioned_delete_sends_version_and_only_explicit_governance_bypass() { + let (client, request_receiver) = test_s3_client(None); + let path = RemotePath::new("test", "bucket", "key.txt"); + + let _ = client + .delete_object_with_options( + &path, + DeleteRequestOptions { + version_id: Some("v1".to_string()), + bypass_governance: true, + force_delete: false, + }, + ) + .await; + + let request = request_receiver.expect_request(); + assert!(request.uri().to_string().contains("versionId=v1")); + assert_eq!( + request.headers().get("x-amz-bypass-governance-retention"), + Some("true") + ); + + let (default_client, default_request_receiver) = test_s3_client(None); + let _ = default_client + .delete_object_with_options( + &path, + DeleteRequestOptions { + version_id: Some("v1".to_string()), + ..Default::default() + }, + ) + .await; + let default_request = default_request_receiver.expect_request(); + assert!( + default_request + .headers() + .get("x-amz-bypass-governance-retention") + .is_none() + ); + } + #[tokio::test] async fn conditional_mirror_writes_and_deletes_set_precondition_headers() { let put_response = http::Response::builder() @@ -4560,6 +4915,163 @@ mod tests { assert_eq!(output, b"abc"); } + #[tokio::test] + async fn get_object_with_options_selects_exact_version() { + let response = http::Response::builder() + .status(200) + .header("content-length", "3") + .header("x-amz-version-id", "v1") + .body(SdkBody::from("old")) + .expect("build versioned get response"); + let (client, request_receiver) = test_s3_client(Some(response)); + let path = RemotePath::new("test", "bucket", "key.txt"); + let options = + ObjectReadOptions::for_version(Some("v1".to_string())).expect("valid version ID"); + + let data = client + .get_object_with_options(&path, &options) + .await + .expect("read exact version"); + + let request = request_receiver.expect_request(); + assert!(request.uri().to_string().contains("versionId=v1")); + assert_eq!(data, b"old"); + } + + #[tokio::test] + async fn head_object_with_options_preserves_version_id() { + let response = http::Response::builder() + .status(200) + .header("content-length", "3") + .header("x-amz-version-id", "v1") + .body(SdkBody::from("")) + .expect("build versioned head response"); + let (client, request_receiver) = test_s3_client(Some(response)); + let path = RemotePath::new("test", "bucket", "key.txt"); + let options = + ObjectReadOptions::for_version(Some("v1".to_string())).expect("valid version ID"); + + let info = client + .head_object_with_options(&path, &options) + .await + .expect("inspect exact version"); + + let request = request_receiver.expect_request(); + assert!(request.uri().to_string().contains("versionId=v1")); + assert_eq!(info.version_id.as_deref(), Some("v1")); + } + + #[tokio::test] + async fn exact_version_errors_distinguish_missing_versions_and_delete_markers() { + let missing_response = http::Response::builder() + .status(404) + .header("x-amz-error-code", "NoSuchVersion") + .body(SdkBody::from( + "NoSuchVersionmissing", + )) + .expect("build missing version response"); + let (missing_client, _) = test_s3_client(Some(missing_response)); + let path = RemotePath::new("test", "bucket", "key.txt"); + let options = + ObjectReadOptions::for_version(Some("missing".to_string())).expect("valid version ID"); + + assert!(matches!( + missing_client + .get_object_with_options(&path, &options) + .await, + Err(Error::VersionNotFound { .. }) + )); + + let marker_response = http::Response::builder() + .status(405) + .header("x-amz-error-code", "MethodNotAllowed") + .header("x-amz-delete-marker", "true") + .header("x-amz-version-id", "marker-v1") + .body(SdkBody::from( + "MethodNotAlloweddelete marker", + )) + .expect("build delete marker response"); + let (marker_client, _) = test_s3_client(Some(marker_response)); + let marker_options = ObjectReadOptions::for_version(Some("marker-v1".to_string())) + .expect("valid version ID"); + + assert!(matches!( + marker_client + .get_object_with_options(&path, &marker_options) + .await, + Err(Error::DeleteMarker { .. }) + )); + } + + #[tokio::test] + async fn exact_version_maps_generic_missing_key_responses_to_missing_version() { + let response = http::Response::builder() + .status(404) + .header("x-amz-error-code", "NoSuchKey") + .body(SdkBody::from( + "NoSuchKeymissing", + )) + .expect("build generic missing response"); + let (client, _) = test_s3_client(Some(response)); + let path = RemotePath::new("test", "bucket", "key.txt"); + let options = + ObjectReadOptions::for_version(Some("missing".to_string())).expect("valid version ID"); + + assert!(matches!( + client.get_object_with_options(&path, &options).await, + Err(Error::VersionNotFound { .. }) + )); + } + + #[tokio::test] + async fn exact_version_maps_bare_unauthorized_status_to_auth() { + let response = http::Response::builder() + .status(401) + .body(SdkBody::from("")) + .expect("build unauthorized response"); + let (client, _) = test_s3_client(Some(response)); + let path = RemotePath::new("test", "bucket", "key.txt"); + let options = + ObjectReadOptions::for_version(Some("v1".to_string())).expect("valid version ID"); + + assert!(matches!( + client.get_object_with_options(&path, &options).await, + Err(Error::Auth(_)) + )); + } + + #[tokio::test] + async fn exact_version_maps_unauthorized_with_retention_text_to_auth() { + let response = http::Response::builder() + .status(401) + .body(SdkBody::from("governance retention is active")) + .expect("build unauthorized response"); + let (client, _) = test_s3_client(Some(response)); + let path = RemotePath::new("test", "bucket", "key.txt"); + let options = + ObjectReadOptions::for_version(Some("v1".to_string())).expect("valid version ID"); + + assert!(matches!( + client.get_object_with_options(&path, &options).await, + Err(Error::Auth(_)) + )); + } + + #[tokio::test] + async fn version_listing_maps_bare_forbidden_status_to_auth() { + let response = http::Response::builder() + .status(403) + .body(SdkBody::from("")) + .expect("build forbidden response"); + let (client, _) = test_s3_client(Some(response)); + let path = RemotePath::new("test", "bucket", "logs/"); + + assert!(matches!( + client.list_object_versions_page(&path, Some(1000)).await, + Err(Error::Auth(_)) + )); + } + #[tokio::test] async fn custom_headers_are_not_required_by_presigned_urls() { let (client, _request_receiver) = test_s3_client_with_endpoint_and_headers( @@ -4658,6 +5170,26 @@ mod tests { ); } + #[tokio::test] + async fn put_object_preserves_returned_version_id() { + let response = http::Response::builder() + .status(200) + .header("etag", "\"etag-v2\"") + .header("x-amz-version-id", "v2") + .body(SdkBody::from("")) + .expect("build versioned put response"); + let (client, _) = test_s3_client(Some(response)); + let path = RemotePath::new("test", "bucket", "file.txt"); + + let info = client + .put_object(&path, b"payload".to_vec(), Some("text/plain"), None) + .await + .expect("put versioned object"); + + assert_eq!(info.version_id.as_deref(), Some("v2")); + assert_eq!(info.etag.as_deref(), Some("etag-v2")); + } + #[tokio::test] async fn copy_object_applies_sse_kms_headers() { let response = http::Response::builder() @@ -5104,7 +5636,10 @@ mod tests { .delete_objects_with_options( "bucket", vec!["key.txt".to_string()], - DeleteRequestOptions { force_delete: true }, + DeleteRequestOptions { + force_delete: true, + ..Default::default() + }, ) .await; @@ -5135,6 +5670,57 @@ mod tests { assert!(request.headers().get("x-rustfs-force-delete").is_none()); } + #[tokio::test] + async fn delete_object_versions_preserves_versions_markers_and_bypass() { + let response = http::Response::builder() + .status(200) + .body(SdkBody::from( + r#" + + key.txtv1 + key.txtmarker-v2truemarker-v2 +"#, + )) + .expect("build version delete response"); + let (client, request_receiver) = test_s3_client(Some(response)); + + let result = client + .delete_object_versions_with_options( + "bucket", + vec![ + ObjectVersionIdentifier { + key: "key.txt".to_string(), + version_id: Some("v1".to_string()), + is_delete_marker: false, + }, + ObjectVersionIdentifier { + key: "key.txt".to_string(), + version_id: Some("marker-v2".to_string()), + is_delete_marker: true, + }, + ], + DeleteRequestOptions { + bypass_governance: true, + ..Default::default() + }, + ) + .await + .expect("delete exact versions"); + + let request = request_receiver.expect_request(); + assert_eq!( + request.headers().get("x-amz-bypass-governance-retention"), + Some("true") + ); + let body = request.body().bytes().expect("request body bytes"); + let body = std::str::from_utf8(body).expect("request body is utf8"); + assert!(body.contains("v1")); + assert!(body.contains("marker-v2")); + assert_eq!(result.deleted.len(), 2); + assert!(result.deleted[1].is_delete_marker); + assert!(result.failures.is_empty()); + } + #[tokio::test] async fn delete_objects_wrapper_uses_default_options_without_rustfs_header() { let response = http::Response::builder() diff --git a/docs/reference/rc/cp.md b/docs/reference/rc/cp.md index 375a989..99f0bb8 100644 --- a/docs/reference/rc/cp.md +++ b/docs/reference/rc/cp.md @@ -68,6 +68,8 @@ Destination encryption flags apply only to remote writes. On `rc cp`, the select The current implementation supports `SSE-S3` and `SSE-KMS`. It does not support `SSE-C`, repeated encryption selectors, or MinIO `mc`-style prefix fan-out matching beyond the exact destination argument for the current command. For shared encryption rules across commands, see [Encryption workflows](encryption.md). +When the server returns a source or destination object version ID, JSON copy output uses the output v3 `versioned_objects` envelope with `data.operation` set to `copy`. `data.source_version_id` identifies the copied source version and `data.version_id` identifies the created destination version. Copies for which the backend reports no version information retain the legacy JSON shape. + Global options shown in command syntax use the same meaning everywhere: | Option | Description | diff --git a/docs/reference/rc/output.md b/docs/reference/rc/output.md index f3ded70..7fb3151 100644 --- a/docs/reference/rc/output.md +++ b/docs/reference/rc/output.md @@ -17,7 +17,7 @@ Every v3 record contains: - `schema_version`, always the integer `3`; - `type`, identifying the command family; - `status`, either `success` or `error`; -- `data` for successful records, or `error` for failed records; +- `data` for successful records and version-removal result details, plus `error` for failed records; - optional `meta` with request and server context. Byte counts are non-negative JSON integers and timestamps use RFC 3339 date-time strings. A field is nullable only when the schema explicitly permits `null`. Server-owned fields that are unavailable on a particular RustFS version are represented as `null`, not omitted, when the field is required by the family contract. @@ -30,10 +30,22 @@ Paginated records use a `pagination` object with `truncated` and `continuation_t Streaming commands emit JSON Lines. Each non-empty line is one complete v3 record and validates independently against `output_v3.json`. A watch keepalive is a successful `watch_event` record with `data.event` set to `null` and `data.keepalive` set to `true`. Consumers must not parse an entire JSON Lines stream as one JSON document. +## Version operation records + +The `versioned_objects` family supports both paginated version listings and exact-version operation records: + +- `data.operation: stat` contains selected-version metadata under `data.object`; +- `data.operation: copy` preserves nullable source and destination version IDs; +- `data.operation: remove` separates `planned`, `removed`, and version-aware `failed` records. + +Version removal sets `data.dry_run` explicitly. Its `outcome` is `planned`, `empty`, `success`, `partial`, or `failed`. Partial and failed removals use an error envelope and retain operation data so automation can distinguish versions that were already removed from versions that require attention. + ## Errors Errors use the same family-specific `type` as successful output. An unsupported server capability has the typed error `unsupported_feature`, including the capability name and nullable server version. Other errors use the stable error kinds defined by the schema. +Non-streaming success records are written to standard output. Error records, including partial version-removal records that retain `data`, are written as one JSON document to standard error. + ## Migrating from v1 or v2 Existing v1 and v2 consumers do not need to migrate until they adopt a new command family or a command explicitly documents v3 output. @@ -47,4 +59,6 @@ When migrating: 5. For watch output, validate and process each JSON Lines record independently. 6. Ignore unknown object properties while continuing to require documented fields and types. +Commands adopting version-operation records migrate only their version-aware JSON paths. Legacy JSON remains unchanged for unversioned `stat`, `rm`, and `cp` results where the server reports no version identifiers. Scripts selecting versions must dispatch on `schema_version: 3` and read operation fields under `data`. + The golden fixtures under `crates/cli/tests/fixtures/output_v3/` provide success, empty, and error examples for every v3 family. diff --git a/docs/reference/rc/rm.md b/docs/reference/rc/rm.md index 13f0bf8..e8df852 100644 --- a/docs/reference/rc/rm.md +++ b/docs/reference/rc/rm.md @@ -19,6 +19,7 @@ rc [GLOBAL OPTIONS] rm [OPTIONS] ... | `-f, --force` | Run without interactive confirmation where required. | | `--dry-run` | Show objects that would be removed. | | `--versions` | Remove object versions where supported. | +| `--version-id ` | Remove exactly one object version without deleting sibling versions. | | `--purge` | Remove all versions and delete markers under the target where supported. | | `--bypass` | Bypass governance retention if the backend and credentials allow it. | @@ -26,7 +27,9 @@ rc [GLOBAL OPTIONS] rm [OPTIONS] ... ```bash rc rm local/reports/tmp.json +rc rm local/reports/tmp.json --version-id VERSION_ID rc object remove local/reports/tmp/ --recursive --dry-run +rc object remove local/reports/tmp/ --recursive --versions --dry-run --json rc object remove local/reports/tmp/ --recursive --force ``` @@ -34,6 +37,8 @@ rc object remove local/reports/tmp/ --recursive --force Recursive and version-aware deletions can remove many objects. Use `--dry-run` to inspect the target set before running destructive commands. +JSON output for `--version-id` and `--versions` uses the output v3 `versioned_objects` envelope. A dry run reports entries under `data.planned` with `data.dry_run` set to `true`; it never reports them as removed. A partial failure uses an error envelope whose `data.removed` and `data.failed` arrays retain the version ID for every result. Unversioned removal retains its legacy JSON shape for compatibility. + Global options shown in command syntax use the same meaning everywhere: | Option | Description | diff --git a/docs/reference/rc/stat.md b/docs/reference/rc/stat.md index 7cd6a78..a791472 100644 --- a/docs/reference/rc/stat.md +++ b/docs/reference/rc/stat.md @@ -29,6 +29,8 @@ rc object stat local/reports/summary.json --json Metadata output includes object identity and storage metadata such as size, modified time, ETag, content type, and user metadata when available. +When `--version-id` is combined with JSON output, `stat` emits the output v3 `versioned_objects` envelope with `data.operation` set to `stat`. Metadata for the selected version is nested under `data.object`. JSON output without `--version-id` retains its legacy shape. + Global options shown in command syntax use the same meaning everywhere: | Option | Description | diff --git a/schemas/output_v3.json b/schemas/output_v3.json index 55b15e0..cac3240 100644 --- a/schemas/output_v3.json +++ b/schemas/output_v3.json @@ -104,7 +104,7 @@ "storage_class": { "$ref": "#/definitions/nullableString" } } }, - "versionedObjectsData": { + "versionedObjectsListData": { "type": "object", "required": ["items", "pagination"], "properties": { @@ -115,6 +115,211 @@ "pagination": { "$ref": "#/definitions/pagination" } } }, + "versionOperationItem": { + "type": "object", + "required": ["path", "version_id", "delete_marker"], + "properties": { + "path": { "type": "string", "minLength": 1 }, + "version_id": { "$ref": "#/definitions/nullableString" }, + "delete_marker": { "type": "boolean" } + } + }, + "versionOperationFailure": { + "type": "object", + "required": ["path", "version_id", "error_type", "message"], + "properties": { + "path": { "type": "string", "minLength": 1 }, + "version_id": { "$ref": "#/definitions/nullableString" }, + "error_type": { + "type": "string", + "enum": [ + "general_error", "usage_error", "network_error", "auth_error", + "not_found", "conflict", "unsupported_feature", "interrupted" + ] + }, + "message": { "type": "string", "minLength": 1 } + } + }, + "versionRemoveSummary": { + "type": "object", + "required": ["planned", "removed", "failed"], + "properties": { + "planned": { "type": "integer", "minimum": 0 }, + "removed": { "type": "integer", "minimum": 0 }, + "failed": { "type": "integer", "minimum": 0 } + } + }, + "versionRemoveData": { + "type": "object", + "required": [ + "operation", "outcome", "dry_run", "planned", "removed", "failed", "summary" + ], + "properties": { + "operation": { "const": "remove" }, + "outcome": { + "type": "string", + "enum": ["success", "empty", "planned", "partial", "failed"] + }, + "dry_run": { "type": "boolean" }, + "planned": { + "type": "array", + "items": { "$ref": "#/definitions/versionOperationItem" } + }, + "removed": { + "type": "array", + "items": { "$ref": "#/definitions/versionOperationItem" } + }, + "failed": { + "type": "array", + "items": { "$ref": "#/definitions/versionOperationFailure" } + }, + "summary": { "$ref": "#/definitions/versionRemoveSummary" } + }, + "allOf": [ + { + "if": { + "properties": { "outcome": { "const": "planned" } }, + "required": ["outcome"] + }, + "then": { + "properties": { + "dry_run": { "const": true }, + "planned": { "minItems": 1 }, + "removed": { "maxItems": 0 }, + "failed": { "maxItems": 0 } + } + } + }, + { + "if": { + "properties": { "outcome": { "const": "partial" } }, + "required": ["outcome"] + }, + "then": { + "properties": { + "dry_run": { "const": false }, + "planned": { "maxItems": 0 }, + "removed": { "minItems": 1 }, + "failed": { "minItems": 1 } + } + } + }, + { + "if": { + "properties": { "outcome": { "const": "failed" } }, + "required": ["outcome"] + }, + "then": { + "properties": { + "removed": { "maxItems": 0 }, + "failed": { "minItems": 1 } + } + } + }, + { + "if": { + "properties": { "outcome": { "const": "success" } }, + "required": ["outcome"] + }, + "then": { + "properties": { + "dry_run": { "const": false }, + "planned": { "maxItems": 0 }, + "removed": { "minItems": 1 }, + "failed": { "maxItems": 0 } + } + } + }, + { + "if": { + "properties": { "outcome": { "const": "empty" } }, + "required": ["outcome"] + }, + "then": { + "properties": { + "planned": { "maxItems": 0 }, + "removed": { "maxItems": 0 }, + "failed": { "maxItems": 0 } + } + } + }, + { + "if": { + "properties": { "dry_run": { "const": true } }, + "required": ["dry_run"] + }, + "then": { + "properties": { "removed": { "maxItems": 0 } } + } + } + ] + }, + "versionCopyData": { + "type": "object", + "required": [ + "operation", "source", "target", "source_version_id", "version_id", + "size_bytes", "size_human" + ], + "properties": { + "operation": { "const": "copy" }, + "source": { "type": "string", "minLength": 1 }, + "target": { "type": "string", "minLength": 1 }, + "source_version_id": { "$ref": "#/definitions/nullableString" }, + "version_id": { "$ref": "#/definitions/nullableString" }, + "size_bytes": { "$ref": "#/definitions/nullableBytes" }, + "size_human": { "$ref": "#/definitions/nullableString" } + } + }, + "versionStatObject": { + "type": "object", + "required": [ + "path", "bucket", "key", "version_id", "delete_marker", "last_modified", + "size_bytes", "size_human", "etag", "content_type", "storage_class", "metadata" + ], + "properties": { + "path": { "type": "string", "minLength": 1 }, + "bucket": { "type": "string", "minLength": 1 }, + "key": { "type": "string" }, + "version_id": { "$ref": "#/definitions/nullableString" }, + "delete_marker": { "type": "boolean" }, + "last_modified": { "$ref": "#/definitions/nullableTimestamp" }, + "size_bytes": { "$ref": "#/definitions/nullableBytes" }, + "size_human": { "$ref": "#/definitions/nullableString" }, + "etag": { "$ref": "#/definitions/nullableString" }, + "content_type": { "$ref": "#/definitions/nullableString" }, + "storage_class": { "$ref": "#/definitions/nullableString" }, + "metadata": { + "oneOf": [ + { + "type": "object", + "additionalProperties": { "type": "string" } + }, + { "type": "null" } + ] + } + } + }, + "versionStatData": { + "type": "object", + "required": ["operation", "object"], + "properties": { + "operation": { "const": "stat" }, + "object": { "$ref": "#/definitions/versionStatObject" } + } + }, + "versionOperationData": { + "oneOf": [ + { "$ref": "#/definitions/versionStatData" }, + { "$ref": "#/definitions/versionCopyData" }, + { "$ref": "#/definitions/versionRemoveData" } + ] + }, + "versionedObjectsData": { + "oneOf": [ + { "$ref": "#/definitions/versionedObjectsListData" }, + { "$ref": "#/definitions/versionOperationData" } + ] + }, "retention": { "type": ["object", "null"], "required": ["mode", "retain_until"], @@ -318,8 +523,22 @@ { "$ref": "#/definitions/standardError" } ] }, + "data": { "type": "object" }, "meta": { "$ref": "#/definitions/meta" } - } + }, + "allOf": [ + { + "if": { + "properties": { "type": { "const": "versioned_objects" } }, + "required": ["type", "data"] + }, + "then": { + "properties": { + "data": { "$ref": "#/definitions/versionOperationData" } + } + } + } + ] } }, "oneOf": [