diff --git a/crates/cli/src/commands/ilm/rule.rs b/crates/cli/src/commands/ilm/rule.rs index 2a0e2925..9f06ef52 100644 --- a/crates/cli/src/commands/ilm/rule.rs +++ b/crates/cli/src/commands/ilm/rule.rs @@ -270,6 +270,7 @@ async fn execute_add(args: AddRuleArgs, output_config: OutputConfig) -> ExitCode Some(LifecycleExpiration { days: args.expiry_days, date: args.expiry_date, + expired_object_all_versions: None, }) } else { None @@ -311,6 +312,7 @@ async fn execute_add(args: AddRuleArgs, output_config: OutputConfig) -> ExitCode (Some(days), Some(sc)) => Some(NoncurrentVersionTransition { noncurrent_days: days, storage_class: sc.clone(), + newer_noncurrent_versions: None, }), (Some(_), None) => { formatter.error("--noncurrent-transition-storage-class is required when using --noncurrent-transition-days"); @@ -330,10 +332,15 @@ async fn execute_add(args: AddRuleArgs, output_config: OutputConfig) -> ExitCode status, prefix: args.prefix, tags: None, + object_size_greater_than: None, + object_size_less_than: None, expiration, + del_marker_expiration: None, transition, + transitions: Vec::new(), noncurrent_version_expiration, noncurrent_version_transition, + noncurrent_version_transitions: Vec::new(), abort_incomplete_multipart_upload_days: None, expired_object_delete_marker, }; @@ -415,6 +422,10 @@ async fn execute_edit(args: EditRuleArgs, output_config: OutputConfig) -> ExitCo date: args .expiry_date .or_else(|| rule.expiration.as_ref().and_then(|e| e.date.clone())), + expired_object_all_versions: rule + .expiration + .as_ref() + .and_then(|expiration| expiration.expired_object_all_versions), }); } @@ -477,6 +488,8 @@ async fn execute_edit(args: EditRuleArgs, output_config: OutputConfig) -> ExitCo .noncurrent_transition_days .unwrap_or_else(|| current.map(|c| c.noncurrent_days).unwrap_or(0)), storage_class: sc, + newer_noncurrent_versions: current + .and_then(|transition| transition.newer_noncurrent_versions), }); } @@ -526,10 +539,13 @@ async fn execute_list(args: BucketArg, output_config: OutputConfig) -> ExitCode Err(code) => return code, }; - let rules = client - .get_bucket_lifecycle(&bucket) - .await - .unwrap_or_default(); + let rules = match client.get_bucket_lifecycle(&bucket).await { + Ok(rules) => rules, + Err(error) => { + formatter.error(&format!("Failed to get lifecycle rules: {error}")); + return ExitCode::GeneralError; + } + }; if formatter.is_json() { formatter.json(&RuleListOutput { bucket, rules }); @@ -679,10 +695,13 @@ async fn execute_export(args: BucketArg, output_config: OutputConfig) -> ExitCod Err(code) => return code, }; - let rules = client - .get_bucket_lifecycle(&bucket) - .await - .unwrap_or_default(); + let rules = match client.get_bucket_lifecycle(&bucket).await { + Ok(rules) => rules, + Err(error) => { + formatter.error(&format!("Failed to get lifecycle rules: {error}")); + return ExitCode::GeneralError; + } + }; let config = LifecycleConfiguration { rules }; formatter.json(&config); @@ -784,7 +803,30 @@ fn validate_expired_delete_marker_rule(rule: &LifecycleRule) -> std::result::Res .as_ref() .is_some_and(|expiration| expiration.days.is_some() || expiration.date.is_some()), rule.tags.as_ref().is_some_and(|tags| !tags.is_empty()), - ) + )?; + + if let Some(expiration) = rule.expiration.as_ref() + && expiration.expired_object_all_versions.is_some() + && (expiration.days.is_none_or(|days| days < 1) || expiration.date.is_some()) + { + return Err( + "ExpiredObjectAllVersions requires positive expiration days and no date".to_string(), + ); + } + + if let Some(expiration) = rule.del_marker_expiration.as_ref() + && expiration.days.is_none_or(|days| days < 1) + { + return Err("delete-marker expiration days must be at least 1".to_string()); + } + + if rule.del_marker_expiration.is_some() + && rule.tags.as_ref().is_some_and(|tags| !tags.is_empty()) + { + return Err("delete-marker expiration cannot be combined with tag filters".to_string()); + } + + Ok(()) } async fn setup_client( @@ -940,13 +982,19 @@ mod tests { status: LifecycleRuleStatus::Enabled, prefix: None, tags: None, + object_size_greater_than: None, + object_size_less_than: None, expiration: Some(LifecycleExpiration { days: Some(30), date: None, + expired_object_all_versions: None, }), + del_marker_expiration: None, transition: None, + transitions: Vec::new(), noncurrent_version_expiration: None, noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), abort_incomplete_multipart_upload_days: None, expired_object_delete_marker: None, }; @@ -960,10 +1008,15 @@ mod tests { status: LifecycleRuleStatus::Enabled, prefix: None, tags: None, + object_size_greater_than: None, + object_size_less_than: None, expiration: None, + del_marker_expiration: None, transition: None, + transitions: Vec::new(), noncurrent_version_expiration: None, noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), abort_incomplete_multipart_upload_days: None, expired_object_delete_marker: None, }; @@ -1022,6 +1075,103 @@ mod tests { assert!(validate_expired_delete_marker_inputs(true, false, false).is_ok()); } + #[test] + fn test_lifecycle_extension_validation_rejects_invalid_values() { + let invalid_all_versions = LifecycleRule { + id: "invalid-all-versions".to_string(), + status: LifecycleRuleStatus::Enabled, + prefix: None, + tags: None, + object_size_greater_than: None, + object_size_less_than: None, + expiration: Some(LifecycleExpiration { + days: Some(0), + date: None, + expired_object_all_versions: Some(true), + }), + del_marker_expiration: None, + transition: None, + transitions: Vec::new(), + noncurrent_version_expiration: None, + noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), + abort_incomplete_multipart_upload_days: None, + expired_object_delete_marker: None, + }; + assert!(validate_expired_delete_marker_rule(&invalid_all_versions).is_err()); + + let all_versions_with_date = LifecycleRule { + id: "all-versions-date".to_string(), + status: LifecycleRuleStatus::Enabled, + prefix: None, + tags: None, + object_size_greater_than: None, + object_size_less_than: None, + expiration: Some(LifecycleExpiration { + days: Some(1), + date: Some("2026-01-01T00:00:00Z".to_string()), + expired_object_all_versions: Some(false), + }), + del_marker_expiration: None, + transition: None, + transitions: Vec::new(), + noncurrent_version_expiration: None, + noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), + abort_incomplete_multipart_upload_days: None, + expired_object_delete_marker: None, + }; + assert!(validate_expired_delete_marker_rule(&all_versions_with_date).is_err()); + + let invalid_delete_marker = LifecycleRule { + id: "invalid-delete-marker".to_string(), + status: LifecycleRuleStatus::Enabled, + prefix: None, + tags: None, + object_size_greater_than: None, + object_size_less_than: None, + expiration: None, + del_marker_expiration: Some(rc_core::LifecycleDelMarkerExpiration { days: None }), + transition: None, + transitions: Vec::new(), + noncurrent_version_expiration: None, + noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), + abort_incomplete_multipart_upload_days: None, + expired_object_delete_marker: None, + }; + assert!(validate_expired_delete_marker_rule(&invalid_delete_marker).is_err()); + } + + #[test] + fn test_lifecycle_import_accepts_issue_6334_shape() { + let input = r#" + { + "rules": [{ + "id": "rule-delayed-deletion", + "status": "Enabled", + "prefix": "test/", + "expiration": { + "ExpiredObjectAllVersions": true, + "DelMarkerExpiration": true, + "days": 1 + } + }] + } + "#; + + let config: LifecycleConfiguration = + serde_json::from_str(input).expect("issue lifecycle import should parse"); + assert_eq!(config.rules.len(), 1); + assert_eq!( + config.rules[0] + .del_marker_expiration + .as_ref() + .and_then(|expiration| expiration.days), + Some(1) + ); + } + #[tokio::test] async fn test_execute_import_rejects_invalid_marker_cleanup_before_alias_lookup() { let temp_dir = tempfile::tempdir().expect("create lifecycle config directory"); @@ -1032,13 +1182,19 @@ mod tests { status: LifecycleRuleStatus::Enabled, prefix: None, tags: None, + object_size_greater_than: None, + object_size_less_than: None, expiration: Some(LifecycleExpiration { days: Some(30), date: None, + expired_object_all_versions: None, }), + del_marker_expiration: None, transition: None, + transitions: Vec::new(), noncurrent_version_expiration: None, noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), abort_incomplete_multipart_upload_days: None, expired_object_delete_marker: Some(true), }], diff --git a/crates/core/src/lib.rs b/crates/core/src/lib.rs index cf69c2c0..7e1963c0 100644 --- a/crates/core/src/lib.rs +++ b/crates/core/src/lib.rs @@ -38,8 +38,9 @@ pub use cors::{CorsConfiguration, CorsRule}; pub use encryption::{BucketEncryption, ObjectEncryptionRequest}; pub use error::{Error, MultipartAbortStatus, Result}; pub use lifecycle::{ - LifecycleConfiguration, LifecycleExpiration, LifecycleRule, LifecycleRuleStatus, - LifecycleTransition, NoncurrentVersionExpiration, NoncurrentVersionTransition, + LifecycleConfiguration, LifecycleDelMarkerExpiration, LifecycleExpiration, LifecycleRule, + LifecycleRuleStatus, LifecycleTransition, NoncurrentVersionExpiration, + NoncurrentVersionTransition, }; pub use multipart_copy::{ DEFAULT_MULTIPART_COPY_PART_SIZE, MultipartCopyCancellation, MultipartCopyOptions, diff --git a/crates/core/src/lifecycle.rs b/crates/core/src/lifecycle.rs index 3ccb412c..4407805c 100644 --- a/crates/core/src/lifecycle.rs +++ b/crates/core/src/lifecycle.rs @@ -6,7 +6,8 @@ use std::collections::HashMap; use std::fmt; -use serde::{Deserialize, Serialize}; +use serde::de::Error as _; +use serde::{Deserialize, Deserializer, Serialize}; /// Full lifecycle configuration for a bucket #[derive(Debug, Clone, Serialize, Deserialize)] @@ -16,7 +17,7 @@ pub struct LifecycleConfiguration { } /// A single lifecycle rule -#[derive(Debug, Clone, Serialize, Deserialize)] +#[derive(Debug, Clone, Serialize)] #[serde(rename_all = "camelCase")] pub struct LifecycleRule { /// Rule identifier @@ -33,14 +34,41 @@ pub struct LifecycleRule { #[serde(skip_serializing_if = "Option::is_none")] pub tags: Option>, + /// Minimum object size in bytes for the rule filter. + #[serde( + skip_serializing_if = "Option::is_none", + alias = "ObjectSizeGreaterThan", + alias = "object_size_greater_than" + )] + pub object_size_greater_than: Option, + + /// Maximum object size in bytes for the rule filter. + #[serde( + skip_serializing_if = "Option::is_none", + alias = "ObjectSizeLessThan", + alias = "object_size_less_than" + )] + pub object_size_less_than: Option, + /// Expiration settings for current object versions #[serde(skip_serializing_if = "Option::is_none")] pub expiration: Option, + /// Expire delete-marker history after the configured number of days. + #[serde(skip_serializing_if = "Option::is_none")] + pub del_marker_expiration: Option, + /// Transition settings for current object versions #[serde(skip_serializing_if = "Option::is_none")] pub transition: Option, + /// Additional transition settings for current object versions. + /// + /// `transition` remains the compatibility field for the first action; + /// this collection carries any subsequent actions returned by S3. + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub transitions: Vec, + /// Expiration settings for noncurrent object versions #[serde(skip_serializing_if = "Option::is_none")] pub noncurrent_version_expiration: Option, @@ -49,6 +77,10 @@ pub struct LifecycleRule { #[serde(skip_serializing_if = "Option::is_none")] pub noncurrent_version_transition: Option, + /// Additional transition settings for noncurrent object versions. + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub noncurrent_version_transitions: Vec, + /// Days after initiation to abort incomplete multipart uploads #[serde(skip_serializing_if = "Option::is_none")] pub abort_incomplete_multipart_upload_days: Option, @@ -58,6 +90,241 @@ pub struct LifecycleRule { pub expired_object_delete_marker: Option, } +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +struct LifecycleRuleInput { + #[serde(alias = "ID")] + id: String, + #[serde(alias = "Status")] + status: LifecycleRuleStatus, + #[serde(default, alias = "Prefix")] + prefix: Option, + #[serde(default, alias = "Tags")] + tags: Option>, + #[serde( + default, + alias = "ObjectSizeGreaterThan", + alias = "object_size_greater_than" + )] + object_size_greater_than: Option, + #[serde(default, alias = "ObjectSizeLessThan", alias = "object_size_less_than")] + object_size_less_than: Option, + #[serde(default, alias = "Expiration")] + expiration: Option, + #[serde(default, alias = "Transition")] + transition: Option, + #[serde(default, alias = "Transitions")] + transitions: Vec, + #[serde(default, alias = "NoncurrentVersionExpiration")] + noncurrent_version_expiration: Option, + #[serde(default, alias = "NoncurrentVersionTransition")] + noncurrent_version_transition: Option, + #[serde(default, alias = "NoncurrentVersionTransitions")] + noncurrent_version_transitions: Vec, + #[serde(default, alias = "AbortIncompleteMultipartUploadDays")] + abort_incomplete_multipart_upload_days: Option, + #[serde( + default, + alias = "ExpiredObjectDeleteMarker", + alias = "expired_object_delete_marker" + )] + expired_object_delete_marker: Option, + #[serde( + default, + alias = "DelMarkerExpiration", + alias = "del_marker_expiration" + )] + del_marker_expiration: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +struct LifecycleExpirationInput { + #[serde(default, alias = "Days")] + days: Option, + #[serde(default, alias = "Date")] + date: Option, + #[serde( + default, + alias = "ExpiredObjectAllVersions", + alias = "expired_object_all_versions" + )] + expired_object_all_versions: Option, + #[serde( + default, + alias = "DelMarkerExpiration", + alias = "del_marker_expiration" + )] + del_marker_expiration: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(untagged)] +enum LifecycleDelMarkerInput { + Flag(bool), + Configuration(LifecycleDelMarkerExpiration), +} + +#[derive(Debug)] +struct NormalizedDelMarkerInput { + enabled: bool, + configuration: Option, +} + +impl<'de> Deserialize<'de> for LifecycleRule { + fn deserialize(deserializer: D) -> std::result::Result + where + D: Deserializer<'de>, + { + let LifecycleRuleInput { + id, + status, + prefix, + tags, + object_size_greater_than, + object_size_less_than, + expiration: expiration_input, + transition, + transitions: mut transitions_input, + noncurrent_version_expiration, + noncurrent_version_transition, + noncurrent_version_transitions: mut noncurrent_version_transitions_input, + abort_incomplete_multipart_upload_days, + expired_object_delete_marker, + del_marker_expiration: top_level_del_marker_input, + } = LifecycleRuleInput::deserialize(deserializer)?; + + let (expiration, nested_del_marker) = match expiration_input { + Some(expiration) => { + let nested_del_marker = expiration.del_marker_expiration; + ( + Some(LifecycleExpiration { + days: expiration.days, + date: expiration.date, + expired_object_all_versions: expiration.expired_object_all_versions, + }), + nested_del_marker, + ) + } + None => (None, None), + }; + + let fallback_days = expiration.as_ref().and_then(|expiration| expiration.days); + let top_level_del_marker = + normalize_del_marker_input(top_level_del_marker_input, fallback_days) + .map_err(D::Error::custom)?; + let nested_del_marker = normalize_del_marker_input(nested_del_marker, fallback_days) + .map_err(D::Error::custom)?; + let del_marker_expiration = + merge_del_marker_inputs(top_level_del_marker, nested_del_marker) + .map_err(D::Error::custom)?; + + if let Some(transition) = transition { + transitions_input.insert(0, transition); + } + let transition = transitions_input.first().cloned(); + let additional_transitions = transitions_input.into_iter().skip(1).collect(); + + if let Some(transition) = noncurrent_version_transition { + noncurrent_version_transitions_input.insert(0, transition); + } + let noncurrent_version_transition = noncurrent_version_transitions_input.first().cloned(); + let additional_noncurrent_version_transitions = noncurrent_version_transitions_input + .into_iter() + .skip(1) + .collect(); + + Ok(Self { + id, + status, + prefix, + tags, + object_size_greater_than, + object_size_less_than, + expiration, + del_marker_expiration, + transition, + transitions: additional_transitions, + noncurrent_version_expiration, + noncurrent_version_transition, + noncurrent_version_transitions: additional_noncurrent_version_transitions, + abort_incomplete_multipart_upload_days, + expired_object_delete_marker, + }) + } +} + +fn normalize_del_marker_input( + input: Option, + fallback_days: Option, +) -> std::result::Result, String> { + match input { + None => Ok(None), + Some(LifecycleDelMarkerInput::Flag(false)) => Ok(Some(NormalizedDelMarkerInput { + enabled: false, + configuration: None, + })), + Some(LifecycleDelMarkerInput::Flag(true)) => fallback_days + .map(|days| { + Some(NormalizedDelMarkerInput { + enabled: true, + configuration: Some(LifecycleDelMarkerExpiration { days: Some(days) }), + }) + }) + .ok_or_else(|| { + "DelMarkerExpiration=true requires expiration.days for compatibility input" + .to_string() + }), + Some(LifecycleDelMarkerInput::Configuration(configuration)) => { + Ok(Some(NormalizedDelMarkerInput { + enabled: true, + configuration: Some(configuration), + })) + } + } +} + +fn merge_del_marker_inputs( + top_level: Option, + nested: Option, +) -> std::result::Result, String> { + match (top_level, nested) { + (None, None) => Ok(None), + (Some(value), None) | (None, Some(value)) => Ok(value.enabled.then(|| { + value + .configuration + .unwrap_or(LifecycleDelMarkerExpiration { days: None }) + })), + (Some(primary), Some(secondary)) => { + if primary.enabled != secondary.enabled { + return Err( + "conflicting DelMarkerExpiration values in lifecycle rule: explicit false cannot be combined with an enabled value".to_string(), + ); + } + if !primary.enabled { + return Ok(None); + } + + let primary_configuration = primary + .configuration + .unwrap_or(LifecycleDelMarkerExpiration { days: None }); + let secondary_configuration = secondary + .configuration + .unwrap_or(LifecycleDelMarkerExpiration { days: None }); + if primary_configuration.days.is_some() + && secondary_configuration.days.is_some() + && primary_configuration.days != secondary_configuration.days + { + return Err("conflicting DelMarkerExpiration values in lifecycle rule".to_string()); + } + + Ok(Some(LifecycleDelMarkerExpiration { + days: primary_configuration.days.or(secondary_configuration.days), + })) + } + } +} + /// Rule status #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] pub enum LifecycleRuleStatus { @@ -88,14 +355,32 @@ impl std::str::FromStr for LifecycleRuleStatus { /// Expiration settings for current object versions #[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] pub struct LifecycleExpiration { /// Number of days after creation to expire - #[serde(skip_serializing_if = "Option::is_none")] + #[serde(skip_serializing_if = "Option::is_none", alias = "Days")] pub days: Option, /// Specific date to expire (ISO 8601 format) - #[serde(skip_serializing_if = "Option::is_none")] + #[serde(skip_serializing_if = "Option::is_none", alias = "Date")] pub date: Option, + + /// Whether all object versions should be expired together. + #[serde( + skip_serializing_if = "Option::is_none", + alias = "ExpiredObjectAllVersions", + alias = "expired_object_all_versions" + )] + pub expired_object_all_versions: Option, +} + +/// Expiration settings for delete markers and their prior object versions. +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct LifecycleDelMarkerExpiration { + /// Number of days before delete-marker history is removed. + #[serde(skip_serializing_if = "Option::is_none", alias = "Days")] + pub days: Option, } /// Transition settings for current object versions @@ -135,6 +420,10 @@ pub struct NoncurrentVersionTransition { /// Target storage class (tier name) pub storage_class: String, + + /// Maximum number of newer noncurrent versions to retain. + #[serde(skip_serializing_if = "Option::is_none")] + pub newer_noncurrent_versions: Option, } impl fmt::Display for LifecycleRule { @@ -173,13 +462,19 @@ mod tests { status: LifecycleRuleStatus::Enabled, prefix: Some("logs/".to_string()), tags: None, + object_size_greater_than: None, + object_size_less_than: None, expiration: Some(LifecycleExpiration { days: Some(30), date: None, + expired_object_all_versions: None, }), + del_marker_expiration: None, transition: None, + transitions: Vec::new(), noncurrent_version_expiration: None, noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), abort_incomplete_multipart_upload_days: Some(7), expired_object_delete_marker: None, }; @@ -217,16 +512,22 @@ mod tests { status: LifecycleRuleStatus::Enabled, prefix: None, tags: None, + object_size_greater_than: None, + object_size_less_than: None, expiration: Some(LifecycleExpiration { days: Some(365), date: None, + expired_object_all_versions: None, }), + del_marker_expiration: None, transition: None, + transitions: Vec::new(), noncurrent_version_expiration: Some(NoncurrentVersionExpiration { noncurrent_days: 30, newer_noncurrent_versions: Some(3), }), noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), abort_incomplete_multipart_upload_days: None, expired_object_delete_marker: Some(true), }], @@ -251,6 +552,7 @@ mod tests { let nvt = NoncurrentVersionTransition { noncurrent_days: 60, storage_class: "COLD_TIER".to_string(), + newer_noncurrent_versions: None, }; let json = serde_json::to_string(&nvt).unwrap(); @@ -258,4 +560,207 @@ mod tests { assert_eq!(decoded.noncurrent_days, 60); assert_eq!(decoded.storage_class, "COLD_TIER"); } + + #[test] + fn test_lifecycle_extensions_use_canonical_json_shape() { + let config = LifecycleConfiguration { + rules: vec![LifecycleRule { + id: "all-versions".to_string(), + status: LifecycleRuleStatus::Enabled, + prefix: Some("test/".to_string()), + tags: None, + object_size_greater_than: None, + object_size_less_than: None, + expiration: Some(LifecycleExpiration { + days: Some(1), + date: None, + expired_object_all_versions: Some(true), + }), + del_marker_expiration: Some(LifecycleDelMarkerExpiration { days: Some(1) }), + transition: None, + transitions: Vec::new(), + noncurrent_version_expiration: None, + noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), + abort_incomplete_multipart_upload_days: None, + expired_object_delete_marker: None, + }], + }; + + let value = serde_json::to_value(&config).expect("serialize lifecycle extensions"); + assert_eq!( + value["rules"][0]["expiration"]["expiredObjectAllVersions"], + true + ); + assert_eq!(value["rules"][0]["delMarkerExpiration"]["days"], 1); + assert!( + value["rules"][0]["expiration"] + .get("DelMarkerExpiration") + .is_none() + ); + let decoded: LifecycleConfiguration = + serde_json::from_value(value).expect("canonical lifecycle extensions should parse"); + assert_eq!( + decoded.rules[0] + .expiration + .as_ref() + .and_then(|expiration| expiration.expired_object_all_versions), + Some(true) + ); + assert_eq!( + decoded.rules[0] + .del_marker_expiration + .as_ref() + .and_then(|expiration| expiration.days), + Some(1) + ); + } + + #[test] + fn test_lifecycle_extensions_accept_issue_6334_compatibility_shape() { + let input = r#" + { + "rules": [{ + "id": "rule-delayed-deletion", + "status": "Enabled", + "prefix": "test/", + "expiration": { + "ExpiredObjectAllVersions": true, + "DelMarkerExpiration": true, + "days": 1 + } + }] + } + "#; + + let config: LifecycleConfiguration = + serde_json::from_str(input).expect("issue compatibility shape should parse"); + let rule = &config.rules[0]; + assert_eq!( + rule.expiration + .as_ref() + .and_then(|expiration| expiration.expired_object_all_versions), + Some(true) + ); + assert_eq!( + rule.del_marker_expiration + .as_ref() + .and_then(|expiration| expiration.days), + Some(1) + ); + } + + #[test] + fn test_lifecycle_extensions_reject_ambiguous_compatibility_shape() { + let input = r#" + { + "rules": [{ + "id": "conflicting", + "status": "Enabled", + "expiration": {"days": 1, "DelMarkerExpiration": true}, + "delMarkerExpiration": {"days": 2} + }] + } + "#; + + let result = serde_json::from_str::(input); + assert!(result.is_err()); + } + + #[test] + fn test_lifecycle_extensions_reject_explicit_false_conflicts() { + let inputs = [ + r#" + { + "rules": [{ + "id": "conflicting-false-top-level", + "status": "Enabled", + "expiration": {"days": 1, "DelMarkerExpiration": true}, + "delMarkerExpiration": false + }] + } + "#, + r#" + { + "rules": [{ + "id": "conflicting-false-nested", + "status": "Enabled", + "expiration": {"days": 1, "DelMarkerExpiration": false}, + "delMarkerExpiration": {"days": 1} + }] + } + "#, + ]; + + for input in inputs { + assert!( + serde_json::from_str::(input).is_err(), + "explicit false must not be silently overridden by an enabled representation" + ); + } + + let explicit_false = r#" + { + "rules": [{ + "id": "explicit-false", + "status": "Enabled", + "delMarkerExpiration": false + }] + } + "#; + let config: LifecycleConfiguration = + serde_json::from_str(explicit_false).expect("standalone false is compatible"); + assert!(config.rules[0].del_marker_expiration.is_none()); + } + + #[test] + fn test_lifecycle_rule_accepts_additional_actions_and_size_filters() { + let input = r#" + { + "rules": [{ + "id": "full-rule", + "status": "Enabled", + "objectSizeGreaterThan": 500, + "objectSizeLessThan": 64000, + "transition": {"days": 30, "storageClass": "WARM"}, + "transitions": [{"days": 60, "storageClass": "COLD"}], + "noncurrentVersionTransition": { + "noncurrentDays": 90, + "storageClass": "WARM", + "newerNoncurrentVersions": 2 + }, + "noncurrentVersionTransitions": [{ + "noncurrentDays": 180, + "storageClass": "COLD", + "newerNoncurrentVersions": 1 + }] + }] + } + "#; + + let config: LifecycleConfiguration = + serde_json::from_str(input).expect("full lifecycle rule should parse"); + let rule = &config.rules[0]; + assert_eq!(rule.object_size_greater_than, Some(500)); + assert_eq!(rule.object_size_less_than, Some(64000)); + assert_eq!(rule.transitions.len(), 1); + assert_eq!(rule.noncurrent_version_transitions.len(), 1); + assert_eq!( + rule.noncurrent_version_transition + .as_ref() + .and_then(|transition| transition.newer_noncurrent_versions), + Some(2) + ); + assert_eq!( + rule.noncurrent_version_transitions[0].newer_noncurrent_versions, + Some(1) + ); + + let serialized = + serde_json::to_value(&config).expect("full lifecycle rule should serialize"); + let decoded: LifecycleConfiguration = + serde_json::from_value(serialized).expect("full lifecycle rule should round-trip"); + assert_eq!(decoded.rules[0].transitions.len(), 1); + assert_eq!(decoded.rules[0].object_size_less_than, Some(64000)); + } } diff --git a/crates/s3/src/client.rs b/crates/s3/src/client.rs index bf7844d8..538ed290 100644 --- a/crates/s3/src/client.rs +++ b/crates/s3/src/client.rs @@ -49,10 +49,10 @@ use rc_core::{ ReplicationResyncStartResult, ReplicationResyncState, ReplicationResyncStatus, ReplicationResyncTargetStatus, RequestHeader, Result, RetentionDuration, RetentionDurationUnit, RetentionMode, SelectOptions, SseCustomerKey, TransferCopyOptions, TransferReadOptions, - global_request_headers, + global_request_headers, is_retryable_error, retry_with_backoff, }; use reqwest::Method; -use reqwest::header::{CONTENT_TYPE, HeaderMap, HeaderName, HeaderValue}; +use reqwest::header::{CONTENT_TYPE, HeaderMap, HeaderName, HeaderValue, LOCATION}; use serde::Deserialize; use sha2::{Digest, Sha256}; use std::collections::HashMap; @@ -63,6 +63,11 @@ use tokio::io::AsyncWrite; use tokio::io::AsyncWriteExt; use zeroize::Zeroizing; +use crate::lifecycle_xml::{ + build_lifecycle_configuration_xml, parse_lifecycle_configuration_xml, + validate_lifecycle_configuration_xml_response, +}; + /// Keep single-part uploads small to avoid backend incompatibilities with /// streaming aws-chunked payloads. const SINGLE_PUT_OBJECT_MAX_SIZE: u64 = crate::multipart::DEFAULT_PART_SIZE; @@ -73,6 +78,7 @@ const REPLICATION_EXTENSION_BODY_LIMIT: u64 = 1024 * 1024; const REPLICATION_CHECK_PROBE_NAMESPACE: &str = ".rustfs.sys/replication-check/"; const REPLICATION_CHECK_ERROR_LIMIT: usize = 512; const REPLICATION_CHECK_DESCRIPTION_LIMIT: usize = 1024; +const MAX_LIFECYCLE_REDIRECTS: u32 = 5; fn contains_control_characters(value: &str) -> bool { value.chars().any(char::is_control) @@ -163,6 +169,19 @@ enum BucketPolicyErrorKind { Other, } +#[derive(Debug)] +struct XmlResponse { + status: reqwest::StatusCode, + headers: HeaderMap, + body: String, +} + +#[derive(Debug, Clone, Copy)] +struct XmlRequestOptions { + include_content_md5: bool, + include_custom_headers: bool, +} + /// Custom HTTP connector using reqwest, supporting insecure TLS (skip cert verification) /// and custom CA bundles. Used when `alias.insecure = true` or `alias.ca_bundle.is_some()`. #[derive(Debug, Clone)] @@ -270,6 +289,93 @@ fn sdk_retry_config( .build()) } +fn is_lifecycle_redirect(status: reqwest::StatusCode) -> bool { + matches!( + status, + reqwest::StatusCode::MOVED_PERMANENTLY + | reqwest::StatusCode::TEMPORARY_REDIRECT + | reqwest::StatusCode::PERMANENT_REDIRECT + ) +} + +fn is_lifecycle_retryable_error(error: &Error) -> bool { + let Error::Network(message) = error else { + return is_retryable_error(error); + }; + if message.starts_with("Request failed:") || message.starts_with("Failed to read response:") { + return true; + } + if let Some(status) = message + .strip_prefix("HTTP ") + .and_then(|message| message.split(':').next()) + .and_then(|status| status.parse::().ok()) + { + return (500..600).contains(&status) + || status == reqwest::StatusCode::TOO_MANY_REQUESTS.as_u16(); + } + is_retryable_error(error) +} + +fn recognized_s3_provider(host: &str) -> Option<&'static str> { + let host = host.to_ascii_lowercase(); + [ + ("amazonaws.com.cn", "s3"), + ("amazonaws.com", "s3"), + ("aliyuncs.com", "oss"), + ] + .into_iter() + .find_map(|(suffix, service)| { + let belongs_to_provider = host == suffix || host.ends_with(&format!(".{suffix}")); + let is_service_host = host + .split('.') + .any(|label| label == service || label.starts_with(&format!("{service}-"))); + (belongs_to_provider && is_service_host).then_some(suffix) + }) +} + +fn lifecycle_redirect_host_is_trusted(current: &str, next: &str) -> bool { + current.eq_ignore_ascii_case(next) + || recognized_s3_provider(current) + .zip(recognized_s3_provider(next)) + .is_some_and(|(current, next)| current == next) +} + +fn resolve_lifecycle_redirect( + current_url: &reqwest::Url, + headers: &HeaderMap, +) -> Result { + let location = headers + .get(LOCATION) + .ok_or_else(|| Error::Network("lifecycle redirect missing Location header".to_string()))? + .to_str() + .map_err(|_| { + Error::Network("lifecycle redirect has an invalid Location header".to_string()) + })?; + let next_url = current_url + .join(location) + .map_err(|error| Error::Network(format!("invalid lifecycle redirect Location: {error}")))?; + if !matches!(next_url.scheme(), "http" | "https") + || !next_url.username().is_empty() + || next_url.password().is_some() + { + return Err(Error::Network( + "lifecycle redirect must target an HTTP(S) URL without credentials".to_string(), + )); + } + let trusted_host = current_url + .host_str() + .zip(next_url.host_str()) + .is_some_and(|(current, next)| lifecycle_redirect_host_is_trusted(current, next)); + let same_port = current_url.port_or_known_default() == next_url.port_or_known_default(); + let downgraded = current_url.scheme() == "https" && next_url.scheme() != "https"; + if !trusted_host || !same_port || downgraded { + return Err(Error::Network( + "lifecycle redirect must stay on the configured endpoint host or a recognized S3 provider, and must not downgrade HTTPS".to_string(), + )); + } + Ok(next_url) +} + fn sdk_timeout_config( config: &rc_core::alias::TimeoutConfig, ) -> Result { @@ -1387,108 +1493,66 @@ fn build_replication_configuration_xml(config: &ReplicationConfiguration) -> Str xml } -fn parse_lifecycle_filter_prefix( - filter: Option<&aws_sdk_s3::types::LifecycleRuleFilter>, -) -> Option { - filter - .and_then(|filter| filter.prefix().map(str::to_string)) - .or_else(|| filter.and_then(|filter| filter.and()?.prefix().map(str::to_string))) -} - -fn parse_lifecycle_filter_tags( - filter: Option<&aws_sdk_s3::types::LifecycleRuleFilter>, -) -> Option> { - filter - .and_then(|filter| collect_tag_map(filter.tag().map(|tag| (tag.key(), tag.value())))) - .or_else(|| { - filter.and_then(|filter| { - collect_tag_map( - filter - .and()? - .tags() - .iter() - .map(|tag| (tag.key(), tag.value())), - ) - }) - }) -} - -fn build_s3_tag(key: &str, value: &str) -> Result { - aws_sdk_s3::types::Tag::builder() - .key(key) - .value(value) - .build() - .map_err(|error| Error::General(format!("build filter tag: {error}"))) -} - -fn build_lifecycle_rule_filter( - prefix: Option<&str>, - tags: Option<&HashMap>, -) -> Result> { - let Some(tags) = tags.filter(|tags| !tags.is_empty()) else { - return Ok(prefix.map(|prefix| { - aws_sdk_s3::types::LifecycleRuleFilter::builder() - .prefix(prefix) - .build() - })); - }; - - let tag_values = sorted_tags(tags) - .into_iter() - .map(|(key, value)| build_s3_tag(key, value)) - .collect::>>()?; - - let filter = if prefix.is_some() || tag_values.len() > 1 { - let mut and_builder = aws_sdk_s3::types::LifecycleRuleAndOperator::builder(); - if let Some(prefix) = prefix { - and_builder = and_builder.prefix(prefix); - } - for tag in tag_values { - and_builder = and_builder.tags(tag); - } - aws_sdk_s3::types::LifecycleRuleFilter::builder() - .and(and_builder.build()) - .build() - } else { - aws_sdk_s3::types::LifecycleRuleFilter::builder() - .tag( - tag_values - .into_iter() - .next() - .expect("non-empty tags required to build lifecycle filter"), - ) - .build() - }; - - Ok(Some(filter)) -} - fn validate_lifecycle_rule(rule: &LifecycleRule) -> Result<()> { - if rule.expired_object_delete_marker != Some(true) { - return Ok(()); + if let Some(expiration) = rule.expiration.as_ref() + && expiration.expired_object_all_versions.is_some() + && (expiration.days.is_none_or(|days| days < 1) || expiration.date.is_some()) + { + return Err(Error::InvalidPath(format!( + "lifecycle rule '{}' requires positive expiration days without a date when ExpiredObjectAllVersions is specified", + rule.id + ))); } - if rule - .expiration - .as_ref() - .is_some_and(|expiration| expiration.days.is_some() || expiration.date.is_some()) + if let Some(expiration) = rule.del_marker_expiration.as_ref() + && expiration.days.is_none_or(|days| days < 1) { return Err(Error::InvalidPath(format!( - "lifecycle rule '{}' cannot combine current expiration days or date with expired delete-marker cleanup", + "lifecycle rule '{}' requires positive days for delete-marker expiration", rule.id ))); } - if rule.tags.as_ref().is_some_and(|tags| !tags.is_empty()) { + if rule.del_marker_expiration.is_some() + && rule.tags.as_ref().is_some_and(|tags| !tags.is_empty()) + { return Err(Error::InvalidPath(format!( - "lifecycle rule '{}' cannot combine tag filters with expired delete-marker cleanup", + "lifecycle rule '{}' cannot combine delete-marker expiration with tag filters", rule.id ))); } + if rule.expired_object_delete_marker == Some(true) { + if rule + .expiration + .as_ref() + .is_some_and(|expiration| expiration.days.is_some() || expiration.date.is_some()) + { + return Err(Error::InvalidPath(format!( + "lifecycle rule '{}' cannot combine current expiration days or date with expired delete-marker cleanup", + rule.id + ))); + } + + if rule.tags.as_ref().is_some_and(|tags| !tags.is_empty()) { + return Err(Error::InvalidPath(format!( + "lifecycle rule '{}' cannot combine tag filters with expired delete-marker cleanup", + rule.id + ))); + } + } + Ok(()) } +fn is_missing_lifecycle_configuration_error(error_text: &str) -> bool { + let normalized = error_text.to_ascii_lowercase(); + normalized.contains("nosuchlifecycleconfiguration") + || normalized.contains("lifecycle configuration is not found") + || normalized.contains("lifecycle configuration was not found") + || normalized.contains("lifecycle configuration does not exist") +} + impl HttpConnector for ReqwestConnector { fn call(&self, mut request: HttpRequest) -> HttpConnectorFuture { let client = self.client.clone(); @@ -2845,6 +2909,37 @@ impl S3Client { Ok(url) } + fn lifecycle_url(&self, bucket: &str) -> Result { + let mut url = + reqwest::Url::parse(self.alias.endpoint.trim_end_matches('/')).map_err(|e| { + Error::Network(format!("Invalid endpoint '{}': {e}", self.alias.endpoint)) + })?; + + if force_path_style_for_alias(&self.alias) { + let mut segments = url.path_segments_mut().map_err(|_| { + Error::Network(format!( + "Endpoint '{}' does not support path-style bucket operations", + self.alias.endpoint + )) + })?; + segments.pop_if_empty(); + segments.push(bucket); + } else { + let host = url + .host_str() + .ok_or_else(|| Error::Network("Missing host in lifecycle endpoint".to_string()))?; + let bucket_host = format!("{bucket}.{host}"); + url.set_host(Some(&bucket_host)).map_err(|_| { + Error::Network(format!( + "Bucket '{bucket}' cannot be used with DNS-style lifecycle requests" + )) + })?; + } + + url.set_query(Some("lifecycle=")); + Ok(url) + } + fn replication_extension_url( &self, bucket: &str, @@ -3179,6 +3274,18 @@ impl S3Client { url: &str, headers: &HeaderMap, body: &[u8], + ) -> Result { + self.sign_xml_request_for_region(method, url, headers, body, &self.alias.region) + .await + } + + async fn sign_xml_request_for_region( + &self, + method: &Method, + url: &str, + headers: &HeaderMap, + body: &[u8], + region: &str, ) -> Result { if self.alias.anonymous { return Ok(headers.clone()); @@ -3198,7 +3305,7 @@ impl S3Client { let signing_params = v4::SigningParams::builder() .identity(&identity) - .region(&self.alias.region) + .region(region) .name(S3_SERVICE_NAME) .time(std::time::SystemTime::now()) .settings(signing_settings) @@ -3240,17 +3347,139 @@ impl S3Client { url: reqwest::Url, content_type: Option<&str>, body: Option>, + ) -> Result { + self.xml_request_inner(method, url, content_type, body, false) + .await + } + + async fn xml_request_inner( + &self, + method: Method, + url: reqwest::Url, + content_type: Option<&str>, + body: Option>, + include_content_md5: bool, + ) -> Result { + let body = body.unwrap_or_default(); + let response = self + .xml_request_once( + &method, + &url, + content_type, + &body, + &self.alias.region, + XmlRequestOptions { + include_content_md5, + include_custom_headers: true, + }, + ) + .await?; + + if !response.status.is_success() { + return Err(self.xml_response_error(&response)); + } + + Ok(response.body) + } + + async fn lifecycle_xml_request( + &self, + method: Method, + url: reqwest::Url, + content_type: Option<&str>, + body: Option>, + include_content_md5: bool, ) -> Result { let body = body.unwrap_or_default(); + let retry = self.alias.retry_config(); + // Reuse the same validation as the SDK transport so malformed alias retry + // settings cannot silently disable lifecycle retries. + sdk_retry_config(&retry)?; + + let initial_region = self.alias.region.clone(); + let initial_host = url.host_str().map(str::to_owned); + retry_with_backoff( + &retry, + || { + let method = method.clone(); + let body = body.clone(); + let mut current_url = url.clone(); + let mut signing_region = initial_region.clone(); + let initial_host = initial_host.clone(); + async move { + let mut redirects = 0; + loop { + let response = self + .xml_request_once( + &method, + ¤t_url, + content_type, + &body, + &signing_region, + XmlRequestOptions { + include_content_md5, + include_custom_headers: current_url + .host_str() + .zip(initial_host.as_deref()) + .is_some_and(|(current, initial)| { + current.eq_ignore_ascii_case(initial) + }), + }, + ) + .await?; + + if is_lifecycle_redirect(response.status) { + if redirects >= MAX_LIFECYCLE_REDIRECTS { + return Err(self.xml_response_error(&response)); + } + if let Some(region) = response + .headers + .get("x-amz-bucket-region") + .and_then(|value| value.to_str().ok()) + .map(str::trim) + .filter(|value| { + !value.is_empty() && !contains_control_characters(value) + }) + { + signing_region = region.to_string(); + } + current_url = + resolve_lifecycle_redirect(¤t_url, &response.headers)?; + redirects += 1; + continue; + } + + if response.status.is_success() { + return Ok(response.body); + } + + return Err(self.xml_response_error(&response)); + } + } + }, + is_lifecycle_retryable_error, + ) + .await + } + + async fn xml_request_once( + &self, + method: &Method, + url: &reqwest::Url, + content_type: Option<&str>, + body: &[u8], + signing_region: &str, + options: XmlRequestOptions, + ) -> Result { let mut headers = HeaderMap::new(); headers.insert( "x-amz-content-sha256", - HeaderValue::from_str(&Self::sha256_hash(&body)) + HeaderValue::from_str(&Self::sha256_hash(body)) .map_err(|e| Error::Auth(format!("Invalid content hash header: {e}")))?, ); headers.insert( "host", - HeaderValue::from_str(&self.request_host(&url)?) + HeaderValue::from_str(&self.request_host(url)?) .map_err(|e| Error::Auth(format!("Invalid host header: {e}")))?, ); @@ -3262,24 +3491,36 @@ impl S3Client { ); } - for header in &self.request_headers { - let name = HeaderName::from_bytes(header.name.as_bytes()) - .map_err(|e| Error::Auth(format!("Invalid custom header name: {e}")))?; - let value = HeaderValue::from_str(&header.value) - .map_err(|e| Error::Auth(format!("Invalid custom header value: {e}")))?; - headers.insert(name, value); + // Caller-supplied x-amz headers can contain credentials; keep them on the + // configured host when a provider redirect changes the authority. + if options.include_custom_headers { + for header in &self.request_headers { + let name = HeaderName::from_bytes(header.name.as_bytes()) + .map_err(|e| Error::Auth(format!("Invalid custom header name: {e}")))?; + let value = HeaderValue::from_str(&header.value) + .map_err(|e| Error::Auth(format!("Invalid custom header value: {e}")))?; + headers.insert(name, value); + } + } + + if options.include_content_md5 { + headers.insert( + HeaderName::from_static("content-md5"), + HeaderValue::from_str(&BASE64_STANDARD.encode(Md5::digest(body))) + .map_err(|e| Error::Auth(format!("Invalid content MD5 header: {e}")))?, + ); } let signed_headers = self - .sign_xml_request(&method, url.as_str(), &headers, &body) + .sign_xml_request_for_region(method, url.as_str(), &headers, body, signing_region) .await?; - let mut request_builder = self.xml_http_client.request(method, url); + let mut request_builder = self.xml_http_client.request(method.clone(), url.clone()); for (name, value) in &signed_headers { request_builder = request_builder.header(name, value); } if !body.is_empty() { - request_builder = request_builder.body(body); + request_builder = request_builder.body(body.to_vec()); } let response = request_builder @@ -3288,20 +3529,25 @@ impl S3Client { .map_err(|e| Error::Network(format!("Request failed: {e}")))?; let status = response.status(); + let response_headers = response.headers().clone(); let text = response .text() .await .map_err(|e| Error::Network(format!("Failed to read response: {e}")))?; - if !status.is_success() { - return Err(Error::Network(format!( - "HTTP {}: {}", - status.as_u16(), - text - ))); - } + Ok(XmlResponse { + status, + headers: response_headers, + body: text, + }) + } - Ok(text) + fn xml_response_error(&self, response: &XmlResponse) -> Error { + Error::Network(self.redact_sensitive_text(format!( + "HTTP {}: {}", + response.status.as_u16(), + response.body + ))) } fn bucket_policy_error_kind( @@ -6294,232 +6540,38 @@ impl ObjectStore for S3Client { } async fn get_bucket_lifecycle(&self, bucket: &str) -> Result> { - let response = match self - .inner - .get_bucket_lifecycle_configuration() - .bucket(bucket) - .send() + let url = self.lifecycle_url(bucket)?; + let body = match self + .lifecycle_xml_request(Method::GET, url, None, None, false) .await { - Ok(resp) => resp, + Ok(body) => body, + Err(Error::Network(error_text)) + if is_missing_lifecycle_configuration_error(&error_text) => + { + return Ok(Vec::new()); + } Err(error) => { - let error_text = Self::format_sdk_error(&error); - if error_text.contains("NoSuchLifecycleConfiguration") - || error_text.contains("lifecycle configuration is not found") - { - return Ok(Vec::new()); - } - return Err(Error::General(format!( - "get_bucket_lifecycle: {error_text}" - ))); + return Err(Error::General(format!("get_bucket_lifecycle: {error}"))); } }; - let mut rules = Vec::new(); - for sdk_rule in response.rules() { - let id = sdk_rule.id().unwrap_or("").to_string(); - let status = match sdk_rule.status().as_str() { - "Enabled" => rc_core::LifecycleRuleStatus::Enabled, - _ => rc_core::LifecycleRuleStatus::Disabled, - }; - - let prefix = parse_lifecycle_filter_prefix(sdk_rule.filter()); - let tags = parse_lifecycle_filter_tags(sdk_rule.filter()); - - let expiration = sdk_rule - .expiration() - .map(|exp| rc_core::LifecycleExpiration { - days: exp.days(), - date: exp.date().map(|d| d.to_string()), - }); - - let transition = sdk_rule - .transitions() - .first() - .map(|t| rc_core::LifecycleTransition { - days: t.days(), - date: t.date().map(|d| d.to_string()), - storage_class: t - .storage_class() - .map(|sc| sc.as_str().to_string()) - .unwrap_or_default(), - }); - - let noncurrent_version_expiration = - sdk_rule.noncurrent_version_expiration().map(|nve| { - rc_core::NoncurrentVersionExpiration { - noncurrent_days: nve.noncurrent_days().unwrap_or(0), - newer_noncurrent_versions: nve.newer_noncurrent_versions(), - } - }); - - let noncurrent_version_transition = sdk_rule - .noncurrent_version_transitions() - .first() - .map(|nvt| rc_core::NoncurrentVersionTransition { - noncurrent_days: nvt.noncurrent_days().unwrap_or(0), - storage_class: nvt - .storage_class() - .map(|sc| sc.as_str().to_string()) - .unwrap_or_default(), - }); - - let abort_incomplete_multipart_upload_days = sdk_rule - .abort_incomplete_multipart_upload() - .and_then(|a| a.days_after_initiation()); - - let expired_object_delete_marker = sdk_rule - .expiration() - .and_then(|e| e.expired_object_delete_marker()) - .filter(|v| *v); - - rules.push(LifecycleRule { - id, - status, - prefix, - tags, - expiration, - transition, - noncurrent_version_expiration, - noncurrent_version_transition, - abort_incomplete_multipart_upload_days, - expired_object_delete_marker, - }); - } - - Ok(rules) + parse_lifecycle_configuration_xml(&body) } async fn set_bucket_lifecycle(&self, bucket: &str, rules: Vec) -> Result<()> { - use aws_sdk_s3::types::{ - AbortIncompleteMultipartUpload, BucketLifecycleConfiguration, ExpirationStatus, - LifecycleExpiration as SdkExpiration, LifecycleRule as SdkRule, - NoncurrentVersionExpiration as SdkNve, NoncurrentVersionTransition as SdkNvt, - Transition, TransitionStorageClass, - }; - - let mut sdk_rules = Vec::new(); - for rule in rules { - validate_lifecycle_rule(&rule)?; - let status = match rule.status { - rc_core::LifecycleRuleStatus::Enabled => ExpirationStatus::Enabled, - rc_core::LifecycleRuleStatus::Disabled => ExpirationStatus::Disabled, - }; - - let filter = build_lifecycle_rule_filter(rule.prefix.as_deref(), rule.tags.as_ref())?; - - let expiration = - if rule.expiration.is_some() || rule.expired_object_delete_marker == Some(true) { - let mut builder = SdkExpiration::builder(); - if let Some(days) = rule - .expiration - .as_ref() - .and_then(|expiration| expiration.days) - { - builder = builder.days(days); - } - if let Some(date_str) = rule - .expiration - .as_ref() - .and_then(|expiration| expiration.date.as_deref()) - && let Ok(dt) = aws_smithy_types::DateTime::from_str( - date_str, - aws_smithy_types::date_time::Format::DateTime, - ) - { - builder = builder.date(dt); - } - if let Some(true) = rule.expired_object_delete_marker { - builder = builder.expired_object_delete_marker(true); - } - Some(builder.build()) - } else { - None - }; - - let transitions = rule.transition.map(|t| { - #[allow(deprecated)] - let sc = TransitionStorageClass::from(t.storage_class.as_str()); - let mut builder = Transition::builder().storage_class(sc); - if let Some(days) = t.days { - builder = builder.days(days); - } - if let Some(ref date_str) = t.date - && let Ok(dt) = aws_smithy_types::DateTime::from_str( - date_str, - aws_smithy_types::date_time::Format::DateTime, - ) - { - builder = builder.date(dt); - } - vec![builder.build()] - }); - - let nve = rule.noncurrent_version_expiration.map(|nve| { - let mut builder = SdkNve::builder().noncurrent_days(nve.noncurrent_days); - if let Some(newer) = nve.newer_noncurrent_versions { - builder = builder.newer_noncurrent_versions(newer); - } - builder.build() - }); - - let nvt = rule.noncurrent_version_transition.map(|nvt| { - let sc = TransitionStorageClass::from(nvt.storage_class.as_str()); - let builder = SdkNvt::builder() - .noncurrent_days(nvt.noncurrent_days) - .storage_class(sc); - vec![builder.build()] - }); - - let abort = rule.abort_incomplete_multipart_upload_days.map(|days| { - AbortIncompleteMultipartUpload::builder() - .days_after_initiation(days) - .build() - }); - - let mut builder = SdkRule::builder().id(&rule.id).status(status); - if let Some(filter) = filter { - builder = builder.filter(filter); - } - if let Some(expiration) = expiration { - builder = builder.expiration(expiration); - } - if let Some(transitions) = transitions { - builder = builder.set_transitions(Some(transitions)); - } - if let Some(nve) = nve { - builder = builder.noncurrent_version_expiration(nve); - } - if let Some(nvt) = nvt { - builder = builder.set_noncurrent_version_transitions(Some(nvt)); - } - if let Some(abort) = abort { - builder = builder.abort_incomplete_multipart_upload(abort); - } - - let sdk_rule = builder - .build() - .map_err(|e| Error::General(format!("build lifecycle rule: {e}")))?; - sdk_rules.push(sdk_rule); + for rule in &rules { + validate_lifecycle_rule(rule)?; } - let config = BucketLifecycleConfiguration::builder() - .set_rules(Some(sdk_rules)) - .build() - .map_err(|e| Error::General(format!("build lifecycle config: {e}")))?; - - self.inner - .put_bucket_lifecycle_configuration() - .bucket(bucket) - .lifecycle_configuration(config) - .send() + let url = self.lifecycle_url(bucket)?; + let body = build_lifecycle_configuration_xml(&rules).into_bytes(); + let response = self + .lifecycle_xml_request(Method::PUT, url, Some("application/xml"), Some(body), true) .await - .map_err(|e| { - Error::General(format!( - "set_bucket_lifecycle: {}", - Self::format_sdk_error(&e) - )) - })?; + .map_err(|error| Error::General(format!("set_bucket_lifecycle: {error}")))?; + validate_lifecycle_configuration_xml_response(&response) + .map_err(|error| Error::General(format!("set_bucket_lifecycle: {error}")))?; Ok(()) } @@ -6737,6 +6789,7 @@ mod tests { method: String, target: String, headers: Vec<(String, String)>, + body: Vec, } fn test_s3_client( @@ -7046,6 +7099,7 @@ mod tests { method, target, headers, + body: buffer[header_end..header_end + content_length].to_vec(), } } @@ -7115,6 +7169,34 @@ mod tests { (endpoint, receiver, handle) } + fn start_lifecycle_sequence_test_server( + responses: Vec>, + ) -> ( + String, + mpsc::Receiver, + thread::JoinHandle<()>, + ) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind lifecycle test server"); + let endpoint = format!("http://{}", listener.local_addr().expect("local addr")); + let (sender, receiver) = mpsc::channel(); + + let handle = thread::spawn(move || { + for response in responses { + let (mut stream, _) = listener.accept().expect("accept lifecycle request"); + stream + .set_read_timeout(Some(Duration::from_secs(5))) + .expect("set lifecycle request timeout"); + let request = read_xml_request(&mut stream); + sender.send(request).expect("send lifecycle request"); + stream + .write_all(&response) + .expect("write lifecycle response"); + } + }); + + (endpoint, receiver, handle) + } + fn start_repeated_replication_extension_test_server( response: Vec, request_count: usize, @@ -7435,40 +7517,305 @@ mod tests { } #[test] - fn build_lifecycle_rule_filter_preserves_prefix_and_tags() { - let mut tags = HashMap::new(); - tags.insert("env".to_string(), "prod".to_string()); - tags.insert("team".to_string(), "core".to_string()); + fn missing_lifecycle_configuration_detection_does_not_hide_missing_bucket() { + assert!(is_missing_lifecycle_configuration_error( + "HTTP 404: NoSuchLifecycleConfiguration" + )); + assert!(is_missing_lifecycle_configuration_error( + "lifecycle configuration is not found" + )); + assert!(is_missing_lifecycle_configuration_error( + "The lifecycle configuration was not found" + )); + assert!(is_missing_lifecycle_configuration_error( + "HTTP 404: The lifecycle configuration does not exist" + )); + assert!(!is_missing_lifecycle_configuration_error( + "HTTP 404: NoSuchBucket" + )); + assert!(!is_missing_lifecycle_configuration_error("HTTP 404")); + } - let filter = build_lifecycle_rule_filter(Some("logs/"), Some(&tags)) - .expect("build lifecycle filter") - .expect("lifecycle filter"); + #[tokio::test] + async fn get_bucket_lifecycle_maps_only_missing_configuration_to_empty_rules() { + let missing_config_body = b"NoSuchLifecycleConfiguration"; + let mut missing_config_response = format!( + "HTTP/1.1 404 Not Found\r\ncontent-length: {}\r\nconnection: close\r\n\r\n", + missing_config_body.len() + ) + .into_bytes(); + missing_config_response.extend_from_slice(missing_config_body); + let (missing_config_endpoint, missing_config_receiver, missing_config_server) = + start_replication_extension_test_server(missing_config_response); + let (missing_config_client, _) = + test_s3_client_with_endpoint(&missing_config_endpoint, None); + let rules = ObjectStore::get_bucket_lifecycle(&missing_config_client, "bucket") + .await + .expect("missing lifecycle configuration should be empty"); + assert!(rules.is_empty()); + missing_config_receiver + .recv_timeout(Duration::from_secs(5)) + .expect("missing configuration request should be captured"); + missing_config_server + .join() + .expect("missing configuration server should finish"); + + let missing_bucket_body = b"NoSuchBucket"; + let mut missing_bucket_response = format!( + "HTTP/1.1 404 Not Found\r\ncontent-length: {}\r\nconnection: close\r\n\r\n", + missing_bucket_body.len() + ) + .into_bytes(); + missing_bucket_response.extend_from_slice(missing_bucket_body); + let (missing_bucket_endpoint, missing_bucket_receiver, missing_bucket_server) = + start_replication_extension_test_server(missing_bucket_response); + let (missing_bucket_client, _) = + test_s3_client_with_endpoint(&missing_bucket_endpoint, None); + let result = ObjectStore::get_bucket_lifecycle(&missing_bucket_client, "bucket").await; + assert!( + matches!(result, Err(Error::General(ref message)) if message.contains("NoSuchBucket")), + "missing bucket must not be reported as an empty lifecycle configuration: {result:?}" + ); + missing_bucket_receiver + .recv_timeout(Duration::from_secs(5)) + .expect("missing bucket request should be captured"); + missing_bucket_server + .join() + .expect("missing bucket server should finish"); + + let malformed_body = b""; + let mut malformed_response = format!( + "HTTP/1.1 200 OK\r\ncontent-length: {}\r\nconnection: close\r\n\r\n", + malformed_body.len() + ) + .into_bytes(); + malformed_response.extend_from_slice(malformed_body); + let (malformed_endpoint, malformed_receiver, malformed_server) = + start_replication_extension_test_server(malformed_response); + let (malformed_client, _) = test_s3_client_with_endpoint(&malformed_endpoint, None); + let result = ObjectStore::get_bucket_lifecycle(&malformed_client, "bucket").await; + assert!( + matches!(result, Err(Error::General(ref message)) if message.contains("parse bucket lifecycle xml")), + "malformed lifecycle XML must not become an empty configuration: {result:?}" + ); + malformed_receiver + .recv_timeout(Duration::from_secs(5)) + .expect("malformed lifecycle request should be captured"); + malformed_server + .join() + .expect("malformed lifecycle server should finish"); + } - assert_eq!( - parse_lifecycle_filter_prefix(Some(&filter)).as_deref(), - Some("logs/") + #[tokio::test] + async fn get_bucket_lifecycle_rejects_success_error_envelope() { + let body = b"InternalErrorunexpected"; + let response = format!( + "HTTP/1.1 200 OK\r\ncontent-length: {}\r\nconnection: close\r\n\r\n", + body.len() + ) + .into_bytes(); + let mut response = response; + response.extend_from_slice(body); + let (endpoint, receiver, server) = start_replication_extension_test_server(response); + let (client, _) = test_s3_client_with_endpoint(&endpoint, None); + + let result = ObjectStore::get_bucket_lifecycle(&client, "bucket").await; + assert!( + matches!(result, Err(Error::General(ref message)) if message.contains("unexpected lifecycle root")), + "success Error envelope must not be treated as an empty lifecycle: {result:?}" + ); + receiver + .recv_timeout(Duration::from_secs(5)) + .expect("error response request should be captured"); + server.join().expect("error response server should finish"); + } + + #[tokio::test] + async fn set_bucket_lifecycle_rejects_success_error_envelope() { + let body = b"InternalErrorunexpected"; + let response = format!( + "HTTP/1.1 200 OK\r\ncontent-length: {}\r\nconnection: close\r\n\r\n", + body.len() + ) + .into_bytes(); + let mut response = response; + response.extend_from_slice(body); + let (endpoint, receiver, server) = start_replication_extension_test_server(response); + let (client, _) = test_s3_client_with_endpoint(&endpoint, None); + let rule = LifecycleRule { + id: "error-response".to_string(), + status: rc_core::LifecycleRuleStatus::Enabled, + prefix: None, + tags: None, + object_size_greater_than: None, + object_size_less_than: None, + expiration: None, + del_marker_expiration: None, + transition: None, + transitions: Vec::new(), + noncurrent_version_expiration: None, + noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), + abort_incomplete_multipart_upload_days: None, + expired_object_delete_marker: None, + }; + + let result = ObjectStore::set_bucket_lifecycle(&client, "bucket", vec![rule]).await; + assert!( + matches!(result, Err(Error::General(ref message)) if message.contains("unexpected lifecycle root")), + "success Error envelope must not be treated as a successful PUT: {result:?}" ); - let parsed_tags = parse_lifecycle_filter_tags(Some(&filter)).expect("parsed tags"); - assert_eq!(parsed_tags.get("env").map(String::as_str), Some("prod")); - assert_eq!(parsed_tags.get("team").map(String::as_str), Some("core")); + receiver + .recv_timeout(Duration::from_secs(5)) + .expect("error response request should be captured"); + server.join().expect("error response server should finish"); + } + + #[tokio::test] + async fn lifecycle_get_retries_transient_failure() { + let success_body = br#"retryEnabled"#; + let first = b"InternalError"; + let first_response = format!( + "HTTP/1.1 500 Internal Server Error\r\ncontent-length: {}\r\nconnection: close\r\n\r\n", + first.len() + ) + .into_bytes(); + let mut first_response = first_response; + first_response.extend_from_slice(first); + let second_response = format!( + "HTTP/1.1 200 OK\r\ncontent-length: {}\r\nconnection: close\r\n\r\n", + success_body.len() + ) + .into_bytes(); + let mut second_response = second_response; + second_response.extend_from_slice(success_body); + let (endpoint, receiver, server) = + start_lifecycle_sequence_test_server(vec![first_response, second_response]); + let (mut client, _) = test_s3_client_with_endpoint(&endpoint, None); + client.alias.retry = Some(rc_core::alias::RetryConfig { + max_attempts: 2, + initial_backoff_ms: 1, + max_backoff_ms: 1, + }); + + let rules = ObjectStore::get_bucket_lifecycle(&client, "bucket") + .await + .expect("transient lifecycle failure should be retried"); + assert_eq!(rules[0].id, "retry"); + let first_request = receiver + .recv_timeout(Duration::from_secs(5)) + .expect("first lifecycle request should be captured"); + let second_request = receiver + .recv_timeout(Duration::from_secs(5)) + .expect("retry lifecycle request should be captured"); + assert_eq!(first_request.method, "GET"); + assert_eq!(second_request.method, "GET"); + server.join().expect("retry server should finish"); + } + + #[tokio::test] + async fn lifecycle_get_retries_dropped_connection() { + let success_body = br#"network-retryEnabled"#; + let success_response = format!( + "HTTP/1.1 200 OK\r\ncontent-length: {}\r\nconnection: close\r\n\r\n", + success_body.len() + ) + .into_bytes(); + let mut success_response = success_response; + success_response.extend_from_slice(success_body); + let (endpoint, receiver, server) = + start_lifecycle_sequence_test_server(vec![Vec::new(), success_response]); + let (mut client, _) = test_s3_client_with_endpoint(&endpoint, None); + client.alias.retry = Some(rc_core::alias::RetryConfig { + max_attempts: 2, + initial_backoff_ms: 1, + max_backoff_ms: 1, + }); + + let rules = ObjectStore::get_bucket_lifecycle(&client, "bucket") + .await + .expect("dropped lifecycle connection should be retried"); + assert_eq!(rules[0].id, "network-retry"); + for _ in 0..2 { + let request = receiver + .recv_timeout(Duration::from_secs(5)) + .expect("lifecycle request should be captured"); + assert_eq!(request.method, "GET"); + } + server.join().expect("network retry server should finish"); + } + + #[tokio::test] + async fn lifecycle_put_follows_redirect_and_resigns_request() { + let redirect = b"HTTP/1.1 307 Temporary Redirect\r\nlocation: /redirected?lifecycle=\r\nx-amz-bucket-region: eu-west-1\r\ncontent-length: 0\r\nconnection: close\r\n\r\n".to_vec(); + let success = b"HTTP/1.1 200 OK\r\ncontent-length: 0\r\nconnection: close\r\n\r\n".to_vec(); + let (endpoint, receiver, server) = + start_lifecycle_sequence_test_server(vec![redirect, success]); + let (mut client, _) = test_s3_client_with_endpoint(&endpoint, None); + client.alias.retry = Some(rc_core::alias::RetryConfig { + max_attempts: 1, + initial_backoff_ms: 1, + max_backoff_ms: 1, + }); + let rule = LifecycleRule { + id: "redirect".to_string(), + status: rc_core::LifecycleRuleStatus::Enabled, + prefix: None, + tags: None, + object_size_greater_than: None, + object_size_less_than: None, + expiration: None, + del_marker_expiration: None, + transition: None, + transitions: Vec::new(), + noncurrent_version_expiration: None, + noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), + abort_incomplete_multipart_upload_days: None, + expired_object_delete_marker: None, + }; + + ObjectStore::set_bucket_lifecycle(&client, "bucket", vec![rule]) + .await + .expect("lifecycle PUT should follow endpoint redirect"); + let first_request = receiver + .recv_timeout(Duration::from_secs(5)) + .expect("redirect request should be captured"); + let second_request = receiver + .recv_timeout(Duration::from_secs(5)) + .expect("redirect follow-up should be captured"); + assert_eq!(first_request.method, "PUT"); + assert_eq!(second_request.method, "PUT"); + assert_eq!(second_request.target, "/redirected?lifecycle="); + let first_signature = header_value(&first_request.headers, "authorization") + .expect("redirect request should be signed"); + let second_signature = header_value(&second_request.headers, "authorization") + .expect("redirect follow-up should be signed"); + assert_ne!(first_signature, second_signature); + assert!(second_signature.contains("/eu-west-1/s3/aws4_request")); + server.join().expect("redirect server should finish"); } #[tokio::test] async fn set_bucket_lifecycle_serializes_marker_only_expiration() { - let response = http::Response::builder() - .status(200) - .body(SdkBody::empty()) - .expect("build lifecycle response"); - let (client, request_receiver) = test_s3_client(Some(response)); + let (endpoint, request_receiver, server_handle) = start_replication_extension_test_server( + b"HTTP/1.1 200 OK\r\ncontent-length: 0\r\nconnection: close\r\n\r\n".to_vec(), + ); + let (client, _) = test_s3_client_with_endpoint(&endpoint, None); let rule = LifecycleRule { id: "marker-only".to_string(), status: rc_core::LifecycleRuleStatus::Enabled, prefix: Some(String::new()), tags: None, + object_size_greater_than: None, + object_size_less_than: None, expiration: None, + del_marker_expiration: None, transition: None, + transitions: Vec::new(), noncurrent_version_expiration: None, noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), abort_incomplete_multipart_upload_days: None, expired_object_delete_marker: Some(true), }; @@ -7477,33 +7824,41 @@ mod tests { .await .expect("set marker-only lifecycle rule"); - let request = request_receiver.expect_request(); - let body = request.body().bytes().expect("request body bytes"); - let body = std::str::from_utf8(body).expect("request body is utf8"); + let request = request_receiver + .recv_timeout(Duration::from_secs(5)) + .expect("server should capture lifecycle request"); + let body = std::str::from_utf8(&request.body).expect("request body is utf8"); + assert_eq!(request.method, "PUT"); + assert_eq!(request.target, "/bucket?lifecycle="); assert!(body.contains( "true" )); + server_handle.join().expect("server thread should finish"); } #[tokio::test] async fn set_bucket_lifecycle_serializes_noncurrent_expiration_with_marker_cleanup() { - let response = http::Response::builder() - .status(200) - .body(SdkBody::empty()) - .expect("build lifecycle response"); - let (client, request_receiver) = test_s3_client(Some(response)); + let (endpoint, request_receiver, server_handle) = start_replication_extension_test_server( + b"HTTP/1.1 200 OK\r\ncontent-length: 0\r\nconnection: close\r\n\r\n".to_vec(), + ); + let (client, _) = test_s3_client_with_endpoint(&endpoint, None); let rule = LifecycleRule { id: "noncurrent-marker".to_string(), status: rc_core::LifecycleRuleStatus::Enabled, prefix: Some(String::new()), tags: None, + object_size_greater_than: None, + object_size_less_than: None, expiration: None, + del_marker_expiration: None, transition: None, + transitions: Vec::new(), noncurrent_version_expiration: Some(rc_core::NoncurrentVersionExpiration { noncurrent_days: 1, newer_noncurrent_versions: None, }), noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), abort_incomplete_multipart_upload_days: None, expired_object_delete_marker: Some(true), }; @@ -7512,9 +7867,10 @@ mod tests { .await .expect("set noncurrent lifecycle rule with marker cleanup"); - let request = request_receiver.expect_request(); - let body = request.body().bytes().expect("request body bytes"); - let body = std::str::from_utf8(body).expect("request body is utf8"); + let request = request_receiver + .recv_timeout(Duration::from_secs(5)) + .expect("server should capture lifecycle request"); + let body = std::str::from_utf8(&request.body).expect("request body is utf8"); assert!(body.contains( "1" )); @@ -7523,14 +7879,58 @@ mod tests { )); assert!(!body.contains("")); assert!(!body.contains("")); + server_handle.join().expect("server thread should finish"); + } + + #[tokio::test] + async fn set_bucket_lifecycle_serializes_extension_expirations_and_content_md5() { + let (endpoint, request_receiver, server_handle) = start_replication_extension_test_server( + b"HTTP/1.1 200 OK\r\ncontent-length: 0\r\nconnection: close\r\n\r\n".to_vec(), + ); + let (client, _) = test_s3_client_with_endpoint(&endpoint, None); + let rule = LifecycleRule { + id: "extensions".to_string(), + status: rc_core::LifecycleRuleStatus::Enabled, + prefix: Some("test/".to_string()), + tags: None, + object_size_greater_than: None, + object_size_less_than: None, + expiration: Some(rc_core::LifecycleExpiration { + days: Some(1), + date: None, + expired_object_all_versions: Some(true), + }), + del_marker_expiration: Some(rc_core::LifecycleDelMarkerExpiration { days: Some(1) }), + transition: None, + transitions: Vec::new(), + noncurrent_version_expiration: None, + noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), + abort_incomplete_multipart_upload_days: None, + expired_object_delete_marker: None, + }; + + ObjectStore::set_bucket_lifecycle(&client, "bucket", vec![rule]) + .await + .expect("set extension lifecycle rule"); + + let request = request_receiver + .recv_timeout(Duration::from_secs(5)) + .expect("server should capture lifecycle request"); + let body = std::str::from_utf8(&request.body).expect("request body is utf8"); + assert!(body.contains("true")); + assert!(body.contains("1")); + let expected_md5 = BASE64_STANDARD.encode(Md5::digest(&request.body)); + assert_eq!( + header_value(&request.headers, "content-md5"), + Some(expected_md5.as_str()) + ); + server_handle.join().expect("server thread should finish"); } #[tokio::test] async fn get_then_set_bucket_lifecycle_preserves_marker_cleanup() { - let get_response = http::Response::builder() - .status(200) - .body(SdkBody::from( - r#" + let get_body = br#" noncurrent-marker @@ -7538,31 +7938,53 @@ mod tests { 1 true + 2 -"#, - )) - .expect("build get lifecycle response"); - let (read_client, _read_request_receiver) = test_s3_client(Some(get_response)); +"#; + let mut get_response = format!( + "HTTP/1.1 200 OK\r\ncontent-type: application/xml\r\ncontent-length: {}\r\nconnection: close\r\n\r\n", + get_body.len() + ) + .into_bytes(); + get_response.extend_from_slice(get_body); + let (get_endpoint, _get_request_receiver, get_server_handle) = + start_replication_extension_test_server(get_response); + let (read_client, _) = test_s3_client_with_endpoint(&get_endpoint, None); let rules = ObjectStore::get_bucket_lifecycle(&read_client, "bucket") .await .expect("get lifecycle rules"); assert_eq!(rules[0].expired_object_delete_marker, Some(true)); - - let set_response = http::Response::builder() - .status(200) - .body(SdkBody::empty()) - .expect("build set lifecycle response"); - let (write_client, write_request_receiver) = test_s3_client(Some(set_response)); + assert_eq!( + rules[0] + .del_marker_expiration + .as_ref() + .and_then(|expiration| expiration.days), + Some(2) + ); + get_server_handle + .join() + .expect("get server thread should finish"); + + let (put_endpoint, put_request_receiver, put_server_handle) = + start_replication_extension_test_server( + b"HTTP/1.1 200 OK\r\ncontent-length: 0\r\nconnection: close\r\n\r\n".to_vec(), + ); + let (write_client, _) = test_s3_client_with_endpoint(&put_endpoint, None); ObjectStore::set_bucket_lifecycle(&write_client, "bucket", rules) .await .expect("write lifecycle rules back"); - let request = write_request_receiver.expect_request(); - let body = request.body().bytes().expect("request body bytes"); - let body = std::str::from_utf8(body).expect("request body is utf8"); + let request = put_request_receiver + .recv_timeout(Duration::from_secs(5)) + .expect("put server should capture lifecycle request"); + let body = std::str::from_utf8(&request.body).expect("request body is utf8"); assert!(body.contains( "true" )); + assert!(body.contains("2")); + put_server_handle + .join() + .expect("put server thread should finish"); } #[tokio::test] @@ -7576,13 +7998,19 @@ mod tests { status: rc_core::LifecycleRuleStatus::Enabled, prefix: None, tags: None, + object_size_greater_than: None, + object_size_less_than: None, expiration: Some(rc_core::LifecycleExpiration { days: Some(30), date: None, + expired_object_all_versions: None, }), + del_marker_expiration: None, transition: None, + transitions: Vec::new(), noncurrent_version_expiration: None, noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), abort_incomplete_multipart_upload_days: None, expired_object_delete_marker: Some(true), }, @@ -7591,10 +8019,15 @@ mod tests { status: rc_core::LifecycleRuleStatus::Enabled, prefix: None, tags: Some(tags), + object_size_greater_than: None, + object_size_less_than: None, expiration: None, + del_marker_expiration: None, transition: None, + transitions: Vec::new(), noncurrent_version_expiration: None, noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), abort_incomplete_multipart_upload_days: None, expired_object_delete_marker: Some(true), }, @@ -7607,6 +8040,80 @@ mod tests { request_receiver.expect_no_request(); } + #[tokio::test] + async fn set_bucket_lifecycle_rejects_invalid_extension_values() { + let (client, request_receiver) = test_s3_client(None); + let rules = [ + LifecycleRule { + id: "all-versions-without-days".to_string(), + status: rc_core::LifecycleRuleStatus::Enabled, + prefix: None, + tags: None, + object_size_greater_than: None, + object_size_less_than: None, + expiration: Some(rc_core::LifecycleExpiration { + days: None, + date: None, + expired_object_all_versions: Some(true), + }), + del_marker_expiration: None, + transition: None, + transitions: Vec::new(), + noncurrent_version_expiration: None, + noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), + abort_incomplete_multipart_upload_days: None, + expired_object_delete_marker: None, + }, + LifecycleRule { + id: "all-versions-with-date".to_string(), + status: rc_core::LifecycleRuleStatus::Enabled, + prefix: None, + tags: None, + object_size_greater_than: None, + object_size_less_than: None, + expiration: Some(rc_core::LifecycleExpiration { + days: Some(1), + date: Some("2026-01-01T00:00:00Z".to_string()), + expired_object_all_versions: Some(false), + }), + del_marker_expiration: None, + transition: None, + transitions: Vec::new(), + noncurrent_version_expiration: None, + noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), + abort_incomplete_multipart_upload_days: None, + expired_object_delete_marker: None, + }, + LifecycleRule { + id: "delete-marker-without-days".to_string(), + status: rc_core::LifecycleRuleStatus::Enabled, + prefix: None, + tags: None, + object_size_greater_than: None, + object_size_less_than: None, + expiration: None, + del_marker_expiration: Some(rc_core::LifecycleDelMarkerExpiration { days: None }), + transition: None, + transitions: Vec::new(), + noncurrent_version_expiration: None, + noncurrent_version_transition: None, + noncurrent_version_transitions: Vec::new(), + abort_incomplete_multipart_upload_days: None, + expired_object_delete_marker: None, + }, + ]; + + for rule in rules { + assert!(matches!( + ObjectStore::set_bucket_lifecycle(&client, "bucket", vec![rule]).await, + Err(Error::InvalidPath(_)) + )); + } + request_receiver.expect_no_request(); + } + #[test] fn bucket_policy_error_kind_uses_error_code() { assert_eq!( @@ -8119,6 +8626,79 @@ mod tests { assert_eq!(url.as_str(), "https://example.com/bucket-name?cors="); } + #[test] + fn lifecycle_url_uses_dns_style_when_alias_requests_dns() { + let (mut client, _) = test_s3_client(None); + client.alias.bucket_lookup = "dns".to_string(); + + let url = client + .lifecycle_url("bucket-name") + .expect("build lifecycle url"); + + assert_eq!(url.host_str(), Some("bucket-name.example.com")); + assert_eq!(url.path(), "/"); + assert_eq!(url.query(), Some("lifecycle=")); + } + + #[test] + fn lifecycle_redirect_rejects_cross_host_locations() { + let current = reqwest::Url::parse("https://s3.example.com/bucket?lifecycle=") + .expect("valid current lifecycle URL"); + let mut headers = HeaderMap::new(); + headers.insert( + LOCATION, + HeaderValue::from_static("https://attacker.example/collect"), + ); + + let result = resolve_lifecycle_redirect(¤t, &headers); + assert!( + matches!(result, Err(Error::Network(message)) if message.contains("configured endpoint host")) + ); + } + + #[test] + fn lifecycle_redirect_allows_recognized_s3_provider_hosts() { + let current = reqwest::Url::parse("https://s3.amazonaws.com/bucket?lifecycle=") + .expect("valid current lifecycle URL"); + let mut headers = HeaderMap::new(); + headers.insert( + LOCATION, + HeaderValue::from_static("https://bucket.s3.eu-west-1.amazonaws.com/?lifecycle="), + ); + + let result = resolve_lifecycle_redirect(¤t, &headers) + .expect("AWS S3 endpoint redirects should be supported"); + assert_eq!(result.host_str(), Some("bucket.s3.eu-west-1.amazonaws.com")); + } + + #[test] + fn lifecycle_redirect_statuses_are_limited_to_s3_endpoint_redirects() { + assert!(is_lifecycle_redirect( + reqwest::StatusCode::MOVED_PERMANENTLY + )); + assert!(is_lifecycle_redirect( + reqwest::StatusCode::TEMPORARY_REDIRECT + )); + assert!(is_lifecycle_redirect( + reqwest::StatusCode::PERMANENT_REDIRECT + )); + assert!(!is_lifecycle_redirect(reqwest::StatusCode::FOUND)); + assert!(!is_lifecycle_redirect(reqwest::StatusCode::SEE_OTHER)); + } + + #[test] + fn lifecycle_retry_classification_uses_http_status_before_error_text() { + assert!(is_lifecycle_retryable_error(&Error::Network( + "HTTP 500: InternalError".to_string(), + ))); + assert!(is_lifecycle_retryable_error(&Error::Network( + "HTTP 429: throttled".to_string(), + ))); + assert!(!is_lifecycle_retryable_error(&Error::Network( + "HTTP 404: SlowDown".to_string(), + ))); + } + #[test] fn cors_url_rejects_endpoints_without_path_segments() { let (client, _) = test_s3_client_with_endpoint("mailto:test@example.com", None); diff --git a/crates/s3/src/lib.rs b/crates/s3/src/lib.rs index f122d0b6..021879b3 100644 --- a/crates/s3/src/lib.rs +++ b/crates/s3/src/lib.rs @@ -6,6 +6,7 @@ pub mod admin; pub mod client; +mod lifecycle_xml; pub mod multipart; mod ops; mod select; diff --git a/crates/s3/src/lifecycle_xml.rs b/crates/s3/src/lifecycle_xml.rs new file mode 100644 index 00000000..7ad38e1b --- /dev/null +++ b/crates/s3/src/lifecycle_xml.rs @@ -0,0 +1,735 @@ +use std::collections::HashMap; + +use quick_xml::Reader; +use quick_xml::de::from_str as from_xml_str; +use quick_xml::events::Event; +use rc_core::{ + Error, LifecycleDelMarkerExpiration, LifecycleExpiration, LifecycleRule, LifecycleRuleStatus, + LifecycleTransition, NoncurrentVersionExpiration, NoncurrentVersionTransition, Result, +}; +use serde::Deserialize; + +const S3_LIFECYCLE_XML_NAMESPACE: &str = "http://s3.amazonaws.com/doc/2006-03-01/"; + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +struct LifecycleConfigurationXml { + #[serde(rename = "Rule", default)] + rules: Vec, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +struct LifecycleRuleXml { + #[serde(rename = "ID")] + id: Option, + status: Option, + #[serde(rename = "Prefix")] + legacy_prefix: Option, + filter: Option, + expiration: Option, + #[serde(rename = "Transition", default)] + transitions: Vec, + noncurrent_version_expiration: Option, + #[serde(rename = "NoncurrentVersionTransition", default)] + noncurrent_version_transitions: Vec, + abort_incomplete_multipart_upload: Option, + del_marker_expiration: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +struct LifecycleFilterXml { + prefix: Option, + tag: Option, + object_size_greater_than: Option, + object_size_less_than: Option, + and: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +struct LifecycleAndXml { + prefix: Option, + #[serde(rename = "Tag", default)] + tags: Vec, + object_size_greater_than: Option, + object_size_less_than: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +struct LifecycleTagXml { + key: Option, + value: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +struct LifecycleExpirationXml { + date: Option, + days: Option, + expired_object_all_versions: Option, + expired_object_delete_marker: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +struct LifecycleTransitionXml { + date: Option, + days: Option, + storage_class: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +struct NoncurrentVersionExpirationXml { + noncurrent_days: Option, + newer_noncurrent_versions: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +struct NoncurrentVersionTransitionXml { + noncurrent_days: Option, + storage_class: Option, + newer_noncurrent_versions: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +struct AbortIncompleteMultipartUploadXml { + days_after_initiation: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +struct LifecycleDelMarkerExpirationXml { + days: Option, +} + +pub(crate) fn parse_lifecycle_configuration_xml(body: &str) -> Result> { + validate_lifecycle_root(body)?; + let config: LifecycleConfigurationXml = from_xml_str(body) + .map_err(|error| Error::General(format!("parse bucket lifecycle xml: {error}")))?; + + config + .rules + .into_iter() + .map(convert_lifecycle_rule) + .collect() +} + +pub(crate) fn validate_lifecycle_configuration_xml_response(body: &str) -> Result<()> { + if body.trim().is_empty() { + return Ok(()); + } + validate_lifecycle_root(body)?; + from_xml_str::(body) + .map(|_| ()) + .map_err(|error| Error::General(format!("parse bucket lifecycle xml: {error}"))) +} + +fn convert_lifecycle_rule(rule: LifecycleRuleXml) -> Result { + let prefix = rule + .filter + .as_ref() + .and_then(parse_filter_prefix) + .or(rule.legacy_prefix); + let tags = rule.filter.as_ref().and_then(parse_filter_tags); + let (object_size_greater_than, object_size_less_than) = rule + .filter + .as_ref() + .map(parse_filter_object_sizes) + .unwrap_or((None, None)); + + let (expiration, expired_object_delete_marker) = match rule.expiration { + Some(expiration) => ( + Some(LifecycleExpiration { + days: expiration.days, + date: expiration.date, + expired_object_all_versions: expiration.expired_object_all_versions, + }), + expiration + .expired_object_delete_marker + .filter(|value| *value), + ), + None => (None, None), + }; + + let transitions = rule + .transitions + .into_iter() + .map(|transition| LifecycleTransition { + days: transition.days, + date: transition.date, + storage_class: transition.storage_class.unwrap_or_default(), + }) + .collect::>(); + let transition = transitions.first().cloned(); + let additional_transitions = transitions.into_iter().skip(1).collect(); + + let noncurrent_version_expiration = + rule.noncurrent_version_expiration + .map(|expiration| NoncurrentVersionExpiration { + noncurrent_days: expiration.noncurrent_days.unwrap_or_default(), + newer_noncurrent_versions: expiration.newer_noncurrent_versions, + }); + + let noncurrent_version_transitions = rule + .noncurrent_version_transitions + .into_iter() + .map(|transition| NoncurrentVersionTransition { + noncurrent_days: transition.noncurrent_days.unwrap_or_default(), + storage_class: transition.storage_class.unwrap_or_default(), + newer_noncurrent_versions: transition.newer_noncurrent_versions, + }) + .collect::>(); + let noncurrent_version_transition = noncurrent_version_transitions.first().cloned(); + let additional_noncurrent_version_transitions = + noncurrent_version_transitions.into_iter().skip(1).collect(); + + Ok(LifecycleRule { + id: rule.id.unwrap_or_default(), + status: parse_rule_status(rule.status.as_deref()), + prefix, + tags, + object_size_greater_than, + object_size_less_than, + expiration, + del_marker_expiration: rule.del_marker_expiration.map(|expiration| { + LifecycleDelMarkerExpiration { + days: expiration.days, + } + }), + transition, + transitions: additional_transitions, + noncurrent_version_expiration, + noncurrent_version_transition, + noncurrent_version_transitions: additional_noncurrent_version_transitions, + abort_incomplete_multipart_upload_days: rule + .abort_incomplete_multipart_upload + .and_then(|upload| upload.days_after_initiation), + expired_object_delete_marker, + }) +} + +fn validate_lifecycle_root(body: &str) -> Result<()> { + let mut reader = Reader::from_str(body); + reader.config_mut().trim_text(true); + + loop { + match reader.read_event() { + Ok(Event::Start(element)) | Ok(Event::Empty(element)) => { + let root = element.local_name(); + if root.as_ref() == b"LifecycleConfiguration" { + return Ok(()); + } + return Err(Error::General(format!( + "unexpected lifecycle root '{}', expected 'LifecycleConfiguration'", + String::from_utf8_lossy(root.as_ref()) + ))); + } + Ok(Event::Eof) => { + return Err(Error::General( + "parse bucket lifecycle xml: missing root element".to_string(), + )); + } + Ok(_) => {} + Err(error) => { + return Err(Error::General(format!( + "parse bucket lifecycle xml: {error}" + ))); + } + } + } +} + +fn parse_rule_status(status: Option<&str>) -> LifecycleRuleStatus { + match status { + Some(value) if value.eq_ignore_ascii_case("enabled") => LifecycleRuleStatus::Enabled, + _ => LifecycleRuleStatus::Disabled, + } +} + +fn parse_filter_prefix(filter: &LifecycleFilterXml) -> Option { + filter + .prefix + .clone() + .or_else(|| filter.and.as_ref().and_then(|and| and.prefix.clone())) +} + +fn parse_filter_tags(filter: &LifecycleFilterXml) -> Option> { + let mut tags = HashMap::new(); + if let Some(tag) = &filter.tag + && let (Some(key), Some(value)) = (&tag.key, &tag.value) + { + tags.insert(key.clone(), value.clone()); + } + if let Some(and) = &filter.and { + for tag in &and.tags { + if let (Some(key), Some(value)) = (&tag.key, &tag.value) { + tags.insert(key.clone(), value.clone()); + } + } + } + (!tags.is_empty()).then_some(tags) +} + +fn parse_filter_object_sizes(filter: &LifecycleFilterXml) -> (Option, Option) { + ( + filter.object_size_greater_than.or_else(|| { + filter + .and + .as_ref() + .and_then(|and| and.object_size_greater_than) + }), + filter.object_size_less_than.or_else(|| { + filter + .and + .as_ref() + .and_then(|and| and.object_size_less_than) + }), + ) +} + +pub(crate) fn build_lifecycle_configuration_xml(rules: &[LifecycleRule]) -> String { + let mut xml = + String::from(r#""#); + + for rule in rules { + xml.push_str(""); + + if !rule.id.is_empty() { + append_xml_element(&mut xml, "ID", &rule.id); + } + append_xml_element( + &mut xml, + "Status", + match rule.status { + LifecycleRuleStatus::Enabled => "Enabled", + LifecycleRuleStatus::Disabled => "Disabled", + }, + ); + append_filter_xml( + &mut xml, + rule.prefix.as_deref(), + rule.tags.as_ref(), + rule.object_size_greater_than, + rule.object_size_less_than, + ); + append_expiration_xml( + &mut xml, + rule.expiration.as_ref(), + rule.expired_object_delete_marker, + ); + + append_transition_xml(&mut xml, rule.transition.as_ref()); + for transition in &rule.transitions { + append_transition_xml(&mut xml, Some(transition)); + } + + if let Some(expiration) = &rule.noncurrent_version_expiration { + xml.push_str(""); + append_xml_element( + &mut xml, + "NoncurrentDays", + &expiration.noncurrent_days.to_string(), + ); + append_optional_i32( + &mut xml, + "NewerNoncurrentVersions", + expiration.newer_noncurrent_versions, + ); + xml.push_str(""); + } + + append_noncurrent_version_transition_xml( + &mut xml, + rule.noncurrent_version_transition.as_ref(), + ); + for transition in &rule.noncurrent_version_transitions { + append_noncurrent_version_transition_xml(&mut xml, Some(transition)); + } + + if let Some(days) = rule.abort_incomplete_multipart_upload_days { + xml.push_str(""); + append_xml_element(&mut xml, "DaysAfterInitiation", &days.to_string()); + xml.push_str(""); + } + + if let Some(expiration) = &rule.del_marker_expiration { + xml.push_str(""); + append_optional_i32(&mut xml, "Days", expiration.days); + xml.push_str(""); + } + + xml.push_str(""); + } + + xml.push_str(""); + xml +} + +fn append_filter_xml( + xml: &mut String, + prefix: Option<&str>, + tags: Option<&HashMap>, + object_size_greater_than: Option, + object_size_less_than: Option, +) { + let tags = tags.filter(|tags| !tags.is_empty()); + let tag_count = tags.map_or(0, HashMap::len); + let predicate_count = usize::from(prefix.is_some()) + + tag_count + + usize::from(object_size_greater_than.is_some()) + + usize::from(object_size_less_than.is_some()); + if predicate_count == 0 { + return; + } + + xml.push_str(""); + if predicate_count == 1 { + if let Some(prefix) = prefix { + append_xml_element(xml, "Prefix", prefix); + } else if let Some(tags) = tags { + if let Some((key, value)) = sorted_tags(tags).into_iter().next() { + append_tag_xml(xml, key, value); + } + } else if let Some(value) = object_size_greater_than { + append_xml_element(xml, "ObjectSizeGreaterThan", &value.to_string()); + } else if let Some(value) = object_size_less_than { + append_xml_element(xml, "ObjectSizeLessThan", &value.to_string()); + } + } else { + xml.push_str(""); + if let Some(prefix) = prefix { + append_xml_element(xml, "Prefix", prefix); + } + if let Some(tags) = tags { + for (key, value) in sorted_tags(tags) { + append_tag_xml(xml, key, value); + } + } + if let Some(value) = object_size_greater_than { + append_xml_element(xml, "ObjectSizeGreaterThan", &value.to_string()); + } + if let Some(value) = object_size_less_than { + append_xml_element(xml, "ObjectSizeLessThan", &value.to_string()); + } + xml.push_str(""); + } + xml.push_str(""); +} + +fn append_transition_xml(xml: &mut String, transition: Option<&LifecycleTransition>) { + let Some(transition) = transition else { + return; + }; + xml.push_str(""); + append_optional_i32(xml, "Days", transition.days); + append_optional_string(xml, "Date", transition.date.as_deref()); + append_xml_element(xml, "StorageClass", &transition.storage_class); + xml.push_str(""); +} + +fn append_noncurrent_version_transition_xml( + xml: &mut String, + transition: Option<&NoncurrentVersionTransition>, +) { + let Some(transition) = transition else { + return; + }; + xml.push_str(""); + append_xml_element( + xml, + "NoncurrentDays", + &transition.noncurrent_days.to_string(), + ); + append_xml_element(xml, "StorageClass", &transition.storage_class); + append_optional_i32( + xml, + "NewerNoncurrentVersions", + transition.newer_noncurrent_versions, + ); + xml.push_str(""); +} + +fn append_expiration_xml( + xml: &mut String, + expiration: Option<&LifecycleExpiration>, + expired_object_delete_marker: Option, +) { + let marker_enabled = expired_object_delete_marker == Some(true); + let Some(expiration) = expiration else { + if marker_enabled { + xml.push_str(""); + append_xml_element(xml, "ExpiredObjectDeleteMarker", "true"); + xml.push_str(""); + } + return; + }; + + xml.push_str(""); + append_optional_i32(xml, "Days", expiration.days); + append_optional_string(xml, "Date", expiration.date.as_deref()); + append_optional_bool( + xml, + "ExpiredObjectAllVersions", + expiration.expired_object_all_versions, + ); + if marker_enabled { + append_xml_element(xml, "ExpiredObjectDeleteMarker", "true"); + } + xml.push_str(""); +} + +fn append_tag_xml(xml: &mut String, key: &str, value: &str) { + xml.push_str(""); + append_xml_element(xml, "Key", key); + append_xml_element(xml, "Value", value); + xml.push_str(""); +} + +fn append_optional_i32(xml: &mut String, tag: &str, value: Option) { + if let Some(value) = value { + append_xml_element(xml, tag, &value.to_string()); + } +} + +fn append_optional_string(xml: &mut String, tag: &str, value: Option<&str>) { + if let Some(value) = value { + append_xml_element(xml, tag, value); + } +} + +fn append_optional_bool(xml: &mut String, tag: &str, value: Option) { + if let Some(value) = value { + append_xml_element(xml, tag, if value { "true" } else { "false" }); + } +} + +fn append_xml_element(xml: &mut String, tag: &str, value: &str) { + xml.push('<'); + xml.push_str(tag); + xml.push('>'); + xml.push_str(&xml_escape(value)); + xml.push_str("'); +} + +fn sorted_tags(tags: &HashMap) -> Vec<(&str, &str)> { + let mut pairs: Vec<(&str, &str)> = tags + .iter() + .map(|(key, value)| (key.as_str(), value.as_str())) + .collect(); + pairs.sort_unstable(); + pairs +} + +fn xml_escape(value: &str) -> String { + value + .replace('&', "&") + .replace('<', "<") + .replace('>', ">") + .replace('"', """) + .replace('\'', "'") +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn lifecycle_xml_roundtrip_preserves_extension_fields_and_standard_actions() { + let mut tags = HashMap::new(); + tags.insert("env".to_string(), "prod&blue".to_string()); + let rules = vec![LifecycleRule { + id: "rule<&".to_string(), + status: LifecycleRuleStatus::Enabled, + prefix: Some("logs/".to_string()), + tags: Some(tags), + object_size_greater_than: None, + object_size_less_than: None, + expiration: Some(LifecycleExpiration { + days: Some(1), + date: None, + expired_object_all_versions: Some(true), + }), + del_marker_expiration: Some(LifecycleDelMarkerExpiration { days: Some(2) }), + transition: Some(LifecycleTransition { + days: Some(30), + date: None, + storage_class: "WARM<&".to_string(), + }), + noncurrent_version_expiration: Some(NoncurrentVersionExpiration { + noncurrent_days: 3, + newer_noncurrent_versions: Some(1), + }), + noncurrent_version_transition: None, + transitions: Vec::new(), + noncurrent_version_transitions: Vec::new(), + abort_incomplete_multipart_upload_days: Some(4), + expired_object_delete_marker: None, + }]; + + let xml = build_lifecycle_configuration_xml(&rules); + assert!(xml.contains("true")); + assert!(xml.contains("2")); + assert!(xml.contains("rule<&")); + assert!(xml.contains("prod&blue")); + + let parsed = parse_lifecycle_configuration_xml(&xml).expect("parse lifecycle XML"); + assert_eq!(parsed.len(), 1); + assert_eq!(parsed[0].id, "rule<&"); + assert_eq!( + parsed[0] + .expiration + .as_ref() + .and_then(|expiration| expiration.expired_object_all_versions), + Some(true) + ); + assert_eq!( + parsed[0] + .del_marker_expiration + .as_ref() + .and_then(|expiration| expiration.days), + Some(2) + ); + assert_eq!( + parsed[0].tags.as_ref().and_then(|tags| tags.get("env")), + Some(&"prod&blue".to_string()) + ); + } + + #[test] + fn lifecycle_xml_parser_accepts_legacy_prefix_and_marker_expiration() { + let xml = r#" + + + legacy + logs/ + Enabled + + 1 + true + + + + "#; + + let rules = parse_lifecycle_configuration_xml(xml).expect("parse legacy lifecycle XML"); + assert_eq!(rules[0].prefix.as_deref(), Some("logs/")); + assert_eq!(rules[0].expired_object_delete_marker, Some(true)); + } + + #[test] + fn lifecycle_xml_parser_rejects_error_and_unexpected_roots() { + for xml in [ + "InternalError", + "", + ] { + let result = parse_lifecycle_configuration_xml(xml); + assert!( + matches!(result, Err(Error::General(message)) if message.contains("unexpected lifecycle root")), + "unexpected root should fail with a root-validation error: {xml}" + ); + } + } + + #[test] + fn lifecycle_xml_response_validator_rejects_malformed_success_body() { + assert!( + validate_lifecycle_configuration_xml_response("") + .is_err() + ); + assert!( + validate_lifecycle_configuration_xml_response("").is_ok() + ); + assert!(validate_lifecycle_configuration_xml_response(" ").is_ok()); + } + + #[test] + fn lifecycle_xml_roundtrip_preserves_size_filters_and_all_actions() { + let rules = vec![LifecycleRule { + id: "full-rule".to_string(), + status: LifecycleRuleStatus::Enabled, + prefix: Some("logs/".to_string()), + tags: None, + object_size_greater_than: Some(500), + object_size_less_than: Some(64000), + expiration: None, + del_marker_expiration: None, + transition: Some(LifecycleTransition { + days: Some(30), + date: None, + storage_class: "WARM".to_string(), + }), + transitions: vec![LifecycleTransition { + days: Some(60), + date: None, + storage_class: "COLD".to_string(), + }], + noncurrent_version_expiration: Some(NoncurrentVersionExpiration { + noncurrent_days: 90, + newer_noncurrent_versions: Some(2), + }), + noncurrent_version_transition: Some(NoncurrentVersionTransition { + noncurrent_days: 90, + storage_class: "WARM".to_string(), + newer_noncurrent_versions: Some(2), + }), + noncurrent_version_transitions: vec![NoncurrentVersionTransition { + noncurrent_days: 180, + storage_class: "COLD".to_string(), + newer_noncurrent_versions: Some(1), + }], + abort_incomplete_multipart_upload_days: None, + expired_object_delete_marker: None, + }]; + + let xml = build_lifecycle_configuration_xml(&rules); + assert!(xml.contains("500")); + assert!(xml.contains("64000")); + assert_eq!(xml.matches("").count(), 2); + assert_eq!(xml.matches("").count(), 2); + assert!(xml.contains("2")); + assert!(xml.contains("1")); + + let parsed = parse_lifecycle_configuration_xml(&xml).expect("parse full lifecycle XML"); + let rule = &parsed[0]; + assert_eq!(rule.object_size_greater_than, Some(500)); + assert_eq!(rule.object_size_less_than, Some(64000)); + assert_eq!(rule.transitions.len(), 1); + assert_eq!(rule.noncurrent_version_transitions.len(), 1); + assert_eq!( + rule.noncurrent_version_transition + .as_ref() + .and_then(|transition| transition.newer_noncurrent_versions), + Some(2) + ); + assert_eq!( + rule.noncurrent_version_transitions[0].newer_noncurrent_versions, + Some(1) + ); + + let mut direct_filter = String::new(); + append_filter_xml(&mut direct_filter, None, None, Some(7), None); + assert_eq!( + direct_filter, + "7" + ); + let direct_xml = format!( + "direct-sizeEnabled{direct_filter}" + ); + let direct_rule = parse_lifecycle_configuration_xml(&direct_xml) + .expect("direct size filter should parse") + .into_iter() + .next() + .expect("direct size rule should be present"); + assert_eq!(direct_rule.object_size_greater_than, Some(7)); + } +}