use std::collections::HashMap;
use std::fmt;
use std::fmt::{Debug, Display, Formatter};
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use aws_sdk_s3::primitives::DateTime;
use aws_sdk_s3::types::{DeleteMarkerEntry, Object, ObjectVersion};
use zeroize_derive::{Zeroize, ZeroizeOnDrop};
pub mod error;
pub mod event_callback;
pub mod filter_callback;
pub mod token;
#[derive(Debug, Clone, PartialEq)]
pub enum S3Object {
NotVersioning(Object),
Versioning(ObjectVersion),
DeleteMarker(DeleteMarkerEntry),
}
impl S3Object {
pub fn new(key: &str, size: i64) -> Self {
S3Object::NotVersioning(
Object::builder()
.key(key)
.size(size)
.last_modified(DateTime::from_secs(0))
.build(),
)
}
pub fn new_versioned(key: &str, version_id: &str, size: i64) -> Self {
use aws_sdk_s3::types::ObjectVersionStorageClass;
S3Object::Versioning(
ObjectVersion::builder()
.key(key)
.version_id(version_id)
.size(size)
.is_latest(true)
.storage_class(ObjectVersionStorageClass::Standard)
.last_modified(DateTime::from_secs(0))
.build(),
)
}
pub fn key(&self) -> &str {
match &self {
Self::Versioning(object) => object.key().expect("S3 ObjectVersion missing key"),
Self::NotVersioning(object) => object.key().expect("S3 Object missing key"),
Self::DeleteMarker(marker) => marker.key().expect("S3 DeleteMarker missing key"),
}
}
pub fn last_modified(&self) -> &DateTime {
match &self {
Self::Versioning(object) => object
.last_modified()
.expect("S3 ObjectVersion missing last_modified"),
Self::NotVersioning(object) => object
.last_modified()
.expect("S3 Object missing last_modified"),
Self::DeleteMarker(marker) => marker
.last_modified()
.expect("S3 DeleteMarker missing last_modified"),
}
}
pub fn size(&self) -> i64 {
match &self {
Self::Versioning(object) => object.size().expect("S3 ObjectVersion missing size"),
Self::NotVersioning(object) => object.size().expect("S3 Object missing size"),
Self::DeleteMarker(_) => 0,
}
}
pub fn version_id(&self) -> Option<&str> {
match &self {
Self::Versioning(object) => object.version_id(),
Self::NotVersioning(_) => None,
Self::DeleteMarker(object) => object.version_id(),
}
}
pub fn e_tag(&self) -> Option<&str> {
match &self {
Self::Versioning(object) => object.e_tag(),
Self::NotVersioning(object) => object.e_tag(),
Self::DeleteMarker(_) => None,
}
}
pub fn is_latest(&self) -> bool {
match &self {
Self::Versioning(object) => object.is_latest().unwrap_or(true),
Self::NotVersioning(_) => true,
Self::DeleteMarker(marker) => marker.is_latest().unwrap_or(true),
}
}
pub fn is_delete_marker(&self) -> bool {
matches!(self, Self::DeleteMarker(_))
}
}
pub type ObjectKeyMap = Arc<Mutex<HashMap<String, S3Object>>>;
#[derive(Debug, Clone, PartialEq)]
pub enum DeletionStatistics {
DeleteBytes(u64),
DeleteComplete { key: String },
DeleteSkip { key: String },
DeleteError { key: String },
}
#[derive(Debug)]
pub struct DeletionStatsReport {
pub stats_deleted_objects: AtomicU64,
pub stats_deleted_bytes: AtomicU64,
pub stats_failed_objects: AtomicU64,
}
impl DeletionStatsReport {
pub fn new() -> Self {
Self {
stats_deleted_objects: AtomicU64::new(0),
stats_deleted_bytes: AtomicU64::new(0),
stats_failed_objects: AtomicU64::new(0),
}
}
pub fn increment_deleted(&self, bytes: u64) {
self.stats_deleted_objects.fetch_add(1, Ordering::Relaxed);
self.stats_deleted_bytes.fetch_add(bytes, Ordering::Relaxed);
}
pub fn increment_failed(&self) {
self.stats_failed_objects.fetch_add(1, Ordering::Relaxed);
}
pub fn snapshot(&self) -> DeletionStats {
DeletionStats {
stats_deleted_objects: self.stats_deleted_objects.load(Ordering::Relaxed),
stats_deleted_bytes: self.stats_deleted_bytes.load(Ordering::Relaxed),
stats_failed_objects: self.stats_failed_objects.load(Ordering::Relaxed),
duration: Duration::default(),
}
}
}
impl Default for DeletionStatsReport {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct DeletionStats {
pub stats_deleted_objects: u64,
pub stats_deleted_bytes: u64,
pub stats_failed_objects: u64,
pub duration: Duration,
}
#[derive(Debug, Clone, PartialEq)]
pub enum DeletionOutcome {
Success {
key: String,
version_id: Option<String>,
},
Failed {
key: String,
version_id: Option<String>,
error: DeletionError,
retry_count: u32,
},
}
impl DeletionOutcome {
pub fn is_success(&self) -> bool {
matches!(self, DeletionOutcome::Success { .. })
}
pub fn key(&self) -> &str {
match self {
DeletionOutcome::Success { key, .. } => key,
DeletionOutcome::Failed { key, .. } => key,
}
}
pub fn version_id(&self) -> Option<&str> {
match self {
DeletionOutcome::Success { version_id, .. } => version_id.as_deref(),
DeletionOutcome::Failed { version_id, .. } => version_id.as_deref(),
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum DeletionError {
NotFound,
AccessDenied,
PreconditionFailed,
Throttled,
NetworkError(String),
ServiceError(String),
}
impl DeletionError {
pub fn is_retryable(&self) -> bool {
matches!(
self,
DeletionError::Throttled
| DeletionError::NetworkError(_)
| DeletionError::ServiceError(_)
)
}
}
impl Display for DeletionError {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
match self {
DeletionError::NotFound => write!(f, "Object not found"),
DeletionError::AccessDenied => write!(f, "Access denied"),
DeletionError::PreconditionFailed => {
write!(f, "Precondition failed (ETag mismatch)")
}
DeletionError::Throttled => write!(f, "Request throttled"),
DeletionError::NetworkError(msg) => write!(f, "Network error: {msg}"),
DeletionError::ServiceError(msg) => write!(f, "Service error: {msg}"),
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum DeletionEvent {
PipelineStart,
ObjectDeleted {
key: String,
version_id: Option<String>,
size: u64,
},
ObjectFailed {
key: String,
version_id: Option<String>,
error: DeletionError,
},
PipelineEnd,
PipelineError { message: String },
}
#[derive(Debug, Clone, PartialEq)]
pub struct S3Target {
pub bucket: String,
pub prefix: Option<String>,
pub endpoint: Option<String>,
pub region: Option<String>,
}
impl S3Target {
pub fn parse(s3_uri: &str) -> anyhow::Result<Self> {
if !s3_uri.starts_with("s3://") {
return Err(anyhow::anyhow!(error::S3rmError::InvalidUri(format!(
"Target URI must start with 's3://': {s3_uri}"
))));
}
let without_scheme = &s3_uri[5..];
if without_scheme.is_empty() {
return Err(anyhow::anyhow!(error::S3rmError::InvalidUri(format!(
"Bucket name cannot be empty: {s3_uri}"
))));
}
let (bucket, prefix) = match without_scheme.find('/') {
Some(idx) => {
let bucket = &without_scheme[..idx];
let prefix = &without_scheme[idx + 1..];
(
bucket.to_string(),
if prefix.is_empty() {
None
} else {
Some(prefix.to_string())
},
)
}
None => (without_scheme.to_string(), None),
};
if bucket.is_empty() {
return Err(anyhow::anyhow!(error::S3rmError::InvalidUri(format!(
"Bucket name cannot be empty: {s3_uri}"
))));
}
Ok(S3Target {
bucket,
prefix,
endpoint: None,
region: None,
})
}
}
impl Display for S3Target {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
match &self.prefix {
Some(prefix) => write!(f, "s3://{}/{}", self.bucket, prefix),
None => write!(f, "s3://{}", self.bucket),
}
}
}
#[derive(Debug, Clone)]
pub enum StoragePath {
S3 { bucket: String, prefix: String },
}
#[derive(Debug, Clone)]
pub struct ClientConfigLocation {
pub aws_config_file: Option<PathBuf>,
pub aws_shared_credentials_file: Option<PathBuf>,
}
#[derive(Debug, Clone)]
pub enum S3Credentials {
Profile(String),
Credentials { access_keys: AccessKeys },
FromEnvironment,
}
#[derive(Clone, Zeroize, ZeroizeOnDrop)]
pub struct AccessKeys {
pub access_key: String,
pub secret_access_key: String,
pub session_token: Option<String>,
}
impl Debug for AccessKeys {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
let mut keys = f.debug_struct("AccessKeys");
let session_token = self
.session_token
.as_ref()
.map_or("None", |_| "** redacted **");
let masked_access_key = mask_access_key(&self.access_key);
keys.field("access_key", &masked_access_key)
.field("secret_access_key", &"** redacted **")
.field("session_token", &session_token);
keys.finish()
}
}
pub(crate) fn mask_access_key(key: &str) -> String {
const PREFIX_LEN: usize = 4;
const SUFFIX_LEN: usize = 4;
if key.len() < PREFIX_LEN + SUFFIX_LEN + 1 {
"*".repeat(key.len())
} else {
let masked_len = key.len() - PREFIX_LEN - SUFFIX_LEN;
format!(
"{}{}{}",
&key[..PREFIX_LEN],
"*".repeat(masked_len),
&key[key.len() - SUFFIX_LEN..]
)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test_utils::init_dummy_tracing_subscriber;
use aws_sdk_s3::types::{ObjectStorageClass, ObjectVersionStorageClass, Owner};
#[test]
fn non_versioning_object_getters() {
init_dummy_tracing_subscriber();
let object = Object::builder()
.key("test/key.txt")
.size(1024)
.e_tag("my-etag")
.storage_class(ObjectStorageClass::Standard)
.owner(
Owner::builder()
.id("test_id")
.display_name("test_name")
.build(),
)
.last_modified(DateTime::from_secs(777))
.build();
let s3_object = S3Object::NotVersioning(object);
assert_eq!(s3_object.key(), "test/key.txt");
assert_eq!(s3_object.size(), 1024);
assert_eq!(s3_object.e_tag().unwrap(), "my-etag");
assert_eq!(*s3_object.last_modified(), DateTime::from_secs(777));
assert!(s3_object.version_id().is_none());
assert!(s3_object.is_latest());
assert!(!s3_object.is_delete_marker());
}
#[test]
fn versioning_object_getters() {
init_dummy_tracing_subscriber();
let object = ObjectVersion::builder()
.key("test/key.txt")
.version_id("version1")
.is_latest(true)
.size(2048)
.e_tag("my-etag-v1")
.storage_class(ObjectVersionStorageClass::Standard)
.last_modified(DateTime::from_secs(888))
.build();
let s3_object = S3Object::Versioning(object);
assert_eq!(s3_object.key(), "test/key.txt");
assert_eq!(s3_object.size(), 2048);
assert_eq!(s3_object.e_tag().unwrap(), "my-etag-v1");
assert_eq!(*s3_object.last_modified(), DateTime::from_secs(888));
assert_eq!(s3_object.version_id().unwrap(), "version1");
assert!(s3_object.is_latest());
assert!(!s3_object.is_delete_marker());
}
#[test]
fn delete_marker_getters() {
init_dummy_tracing_subscriber();
let marker = DeleteMarkerEntry::builder()
.key("test/deleted.txt")
.version_id("dm-version1")
.is_latest(true)
.last_modified(DateTime::from_secs(999))
.build();
let s3_object = S3Object::DeleteMarker(marker);
assert_eq!(s3_object.key(), "test/deleted.txt");
assert_eq!(s3_object.size(), 0);
assert!(s3_object.e_tag().is_none());
assert_eq!(*s3_object.last_modified(), DateTime::from_secs(999));
assert_eq!(s3_object.version_id().unwrap(), "dm-version1");
assert!(s3_object.is_latest());
assert!(s3_object.is_delete_marker());
}
#[test]
fn s3_object_new_sets_key_and_size() {
let obj = S3Object::new("photos/cat.jpg", 2048);
assert_eq!(obj.key(), "photos/cat.jpg");
assert_eq!(obj.size(), 2048);
}
#[test]
fn s3_object_new_is_not_versioning() {
let obj = S3Object::new("key.txt", 100);
assert!(obj.version_id().is_none());
assert!(obj.is_latest());
assert!(!obj.is_delete_marker());
assert!(matches!(obj, S3Object::NotVersioning(_)));
}
#[test]
fn s3_object_new_defaults_last_modified_to_epoch() {
let obj = S3Object::new("key.txt", 0);
assert_eq!(*obj.last_modified(), DateTime::from_secs(0));
}
#[test]
fn s3_object_new_zero_size() {
let obj = S3Object::new("empty.txt", 0);
assert_eq!(obj.size(), 0);
}
#[test]
fn s3_object_new_versioned_sets_key_version_size() {
let obj = S3Object::new_versioned("logs/app.log", "v1", 512);
assert_eq!(obj.key(), "logs/app.log");
assert_eq!(obj.version_id(), Some("v1"));
assert_eq!(obj.size(), 512);
}
#[test]
fn s3_object_new_versioned_is_latest() {
let obj = S3Object::new_versioned("key.txt", "ver-abc", 100);
assert!(obj.is_latest());
assert!(!obj.is_delete_marker());
assert!(matches!(obj, S3Object::Versioning(_)));
}
#[test]
fn s3_object_new_versioned_defaults_last_modified_to_epoch() {
let obj = S3Object::new_versioned("key.txt", "v1", 0);
assert_eq!(*obj.last_modified(), DateTime::from_secs(0));
}
#[test]
fn s3_object_new_versioned_has_no_etag() {
let obj = S3Object::new_versioned("key.txt", "v1", 100);
assert!(obj.e_tag().is_none());
}
#[test]
fn s3_object_new_has_no_etag() {
let obj = S3Object::new("key.txt", 100);
assert!(obj.e_tag().is_none());
}
#[test]
fn debug_print_access_keys_redacts_secrets() {
let access_keys = AccessKeys {
access_key: "AKIAIOSFODNN7EXAMPLE".to_string(),
secret_access_key: "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY".to_string(),
session_token: Some("session_token_value".to_string()),
};
let debug_string = format!("{access_keys:?}");
assert!(debug_string.contains("access_key: \"AKIA************MPLE\""));
assert!(debug_string.contains("secret_access_key: \"** redacted **\""));
assert!(debug_string.contains("session_token: \"** redacted **\""));
assert!(!debug_string.contains("wJalrXUtnFEMI"));
assert!(!debug_string.contains("AKIAIOSFODNN7EXAMPLE"));
}
#[test]
fn mask_access_key_empty_string() {
assert_eq!(super::mask_access_key(""), "");
}
#[test]
fn mask_access_key_single_char() {
assert_eq!(super::mask_access_key("A"), "*");
}
#[test]
fn mask_access_key_short_keys_fully_masked() {
assert_eq!(super::mask_access_key("AB"), "**");
assert_eq!(super::mask_access_key("ABCD"), "****");
assert_eq!(super::mask_access_key("ABCDEFG"), "*******");
}
#[test]
fn mask_access_key_eight_chars_fully_masked() {
assert_eq!(super::mask_access_key("ABCDEFGH"), "********");
}
#[test]
fn mask_access_key_nine_chars_threshold() {
assert_eq!(super::mask_access_key("ABCDEFGHI"), "ABCD*FGHI");
}
#[test]
fn mask_access_key_ten_chars() {
assert_eq!(super::mask_access_key("ABCDEFGHIJ"), "ABCD**GHIJ");
}
#[test]
fn mask_access_key_standard_aws_key() {
let masked = super::mask_access_key("AKIAIOSFODNN7EXAMPLE");
assert_eq!(masked, "AKIA************MPLE");
assert_eq!(masked.len(), 20);
}
#[test]
fn mask_access_key_preserves_length() {
for len in 0..=30 {
let input: String = "X".repeat(len);
let masked = super::mask_access_key(&input);
assert_eq!(
masked.len(),
input.len(),
"Masked output length must equal input length for {len}-char key"
);
}
}
#[test]
fn mask_access_key_prefix_and_suffix_preserved() {
let masked = super::mask_access_key("AKIAIOSFODNN7EXAMPLE");
assert!(masked.starts_with("AKIA"));
assert!(masked.ends_with("MPLE"));
}
#[test]
fn mask_access_key_middle_fully_masked() {
let masked = super::mask_access_key("AKIAIOSFODNN7EXAMPLE");
let middle = &masked[4..16];
assert!(
middle.chars().all(|c| c == '*'),
"Middle portion should be all asterisks, got: {middle}"
);
}
#[test]
fn mask_access_key_no_original_middle_leaked() {
let original = "AKIAIOSFODNN7EXAMPLE";
let masked = super::mask_access_key(original);
let original_middle = &original[4..16];
assert!(
!masked.contains(original_middle),
"Masked output must not contain the original middle portion"
);
}
#[test]
fn mask_access_key_short_key_no_chars_leaked() {
let original = "ABCDEFGH";
let masked = super::mask_access_key(original);
assert!(
masked.chars().all(|c| c == '*'),
"Fully masked key should contain only asterisks"
);
}
#[test]
fn debug_access_keys_without_session_token() {
let access_keys = AccessKeys {
access_key: "AKIAIOSFODNN7EXAMPLE".to_string(),
secret_access_key: "secret".to_string(),
session_token: None,
};
let debug_string = format!("{access_keys:?}");
assert!(debug_string.contains("access_key: \"AKIA************MPLE\""));
assert!(debug_string.contains("session_token: \"None\""));
assert!(!debug_string.contains("AKIAIOSFODNN7EXAMPLE"));
}
#[test]
fn debug_access_keys_with_short_key() {
let access_keys = AccessKeys {
access_key: "SHORT".to_string(),
secret_access_key: "secret".to_string(),
session_token: None,
};
let debug_string = format!("{access_keys:?}");
assert!(debug_string.contains("access_key: \"*****\""));
assert!(!debug_string.contains("SHORT"));
}
#[test]
fn object_key_map_insert_and_retrieve() {
let map: ObjectKeyMap = Arc::new(Mutex::new(HashMap::new()));
let object = S3Object::NotVersioning(
Object::builder()
.key("test/key.txt")
.size(100)
.last_modified(DateTime::from_secs(1000))
.build(),
);
map.lock()
.unwrap()
.insert("test/key.txt".to_string(), object.clone());
let retrieved = map.lock().unwrap().get("test/key.txt").cloned();
assert_eq!(retrieved, Some(object));
}
#[test]
fn object_key_map_concurrent_access() {
let map: ObjectKeyMap = Arc::new(Mutex::new(HashMap::new()));
let map_clone = Arc::clone(&map);
map.lock().unwrap().insert(
"key1".to_string(),
S3Object::NotVersioning(
Object::builder()
.key("key1")
.size(10)
.last_modified(DateTime::from_secs(1))
.build(),
),
);
map_clone.lock().unwrap().insert(
"key2".to_string(),
S3Object::NotVersioning(
Object::builder()
.key("key2")
.size(20)
.last_modified(DateTime::from_secs(2))
.build(),
),
);
assert_eq!(map.lock().unwrap().len(), 2);
}
#[test]
fn deletion_stats_report_new() {
let report = DeletionStatsReport::new();
assert_eq!(report.stats_deleted_objects.load(Ordering::SeqCst), 0);
assert_eq!(report.stats_deleted_bytes.load(Ordering::SeqCst), 0);
assert_eq!(report.stats_failed_objects.load(Ordering::SeqCst), 0);
}
#[test]
fn deletion_stats_report_default() {
let report = DeletionStatsReport::default();
assert_eq!(report.stats_deleted_objects.load(Ordering::SeqCst), 0);
}
#[test]
fn deletion_stats_report_increment_deleted() {
let report = DeletionStatsReport::new();
report.increment_deleted(1024);
report.increment_deleted(2048);
assert_eq!(report.stats_deleted_objects.load(Ordering::SeqCst), 2);
assert_eq!(report.stats_deleted_bytes.load(Ordering::SeqCst), 3072);
assert_eq!(report.stats_failed_objects.load(Ordering::SeqCst), 0);
}
#[test]
fn deletion_stats_report_increment_failed() {
let report = DeletionStatsReport::new();
report.increment_failed();
report.increment_failed();
report.increment_failed();
assert_eq!(report.stats_deleted_objects.load(Ordering::SeqCst), 0);
assert_eq!(report.stats_failed_objects.load(Ordering::SeqCst), 3);
}
#[test]
fn deletion_stats_report_snapshot() {
let report = DeletionStatsReport::new();
report.increment_deleted(500);
report.increment_deleted(300);
report.increment_failed();
let stats = report.snapshot();
assert_eq!(stats.stats_deleted_objects, 2);
assert_eq!(stats.stats_deleted_bytes, 800);
assert_eq!(stats.stats_failed_objects, 1);
assert_eq!(stats.duration, Duration::default());
}
#[test]
fn deletion_stats_clone() {
let stats = DeletionStats {
stats_deleted_objects: 100,
stats_deleted_bytes: 50_000,
stats_failed_objects: 5,
duration: Duration::from_secs(10),
};
let cloned = stats.clone();
assert_eq!(stats, cloned);
}
#[test]
fn deletion_outcome_success() {
let outcome = DeletionOutcome::Success {
key: "test/key.txt".to_string(),
version_id: Some("v1".to_string()),
};
assert!(outcome.is_success());
assert_eq!(outcome.key(), "test/key.txt");
assert_eq!(outcome.version_id(), Some("v1"));
}
#[test]
fn deletion_outcome_success_no_version() {
let outcome = DeletionOutcome::Success {
key: "test/key.txt".to_string(),
version_id: None,
};
assert!(outcome.is_success());
assert!(outcome.version_id().is_none());
}
#[test]
fn deletion_outcome_failed() {
let outcome = DeletionOutcome::Failed {
key: "test/key.txt".to_string(),
version_id: None,
error: DeletionError::AccessDenied,
retry_count: 3,
};
assert!(!outcome.is_success());
assert_eq!(outcome.key(), "test/key.txt");
}
#[test]
fn deletion_error_is_retryable() {
assert!(!DeletionError::NotFound.is_retryable());
assert!(!DeletionError::AccessDenied.is_retryable());
assert!(!DeletionError::PreconditionFailed.is_retryable());
assert!(DeletionError::Throttled.is_retryable());
assert!(DeletionError::NetworkError("timeout".to_string()).is_retryable());
assert!(DeletionError::ServiceError("500".to_string()).is_retryable());
}
#[test]
fn deletion_error_display() {
assert_eq!(DeletionError::NotFound.to_string(), "Object not found");
assert_eq!(DeletionError::AccessDenied.to_string(), "Access denied");
assert_eq!(
DeletionError::PreconditionFailed.to_string(),
"Precondition failed (ETag mismatch)"
);
assert_eq!(DeletionError::Throttled.to_string(), "Request throttled");
assert_eq!(
DeletionError::NetworkError("conn reset".to_string()).to_string(),
"Network error: conn reset"
);
assert_eq!(
DeletionError::ServiceError("Internal".to_string()).to_string(),
"Service error: Internal"
);
}
#[test]
fn deletion_event_pipeline_start() {
let event = DeletionEvent::PipelineStart;
assert_eq!(event, DeletionEvent::PipelineStart);
}
#[test]
fn deletion_event_object_deleted() {
let event = DeletionEvent::ObjectDeleted {
key: "test/key.txt".to_string(),
version_id: Some("v1".to_string()),
size: 1024,
};
if let DeletionEvent::ObjectDeleted {
key,
version_id,
size,
} = &event
{
assert_eq!(key, "test/key.txt");
assert_eq!(version_id.as_deref(), Some("v1"));
assert_eq!(*size, 1024);
} else {
panic!("Expected ObjectDeleted event");
}
}
#[test]
fn deletion_event_object_failed() {
let event = DeletionEvent::ObjectFailed {
key: "test/key.txt".to_string(),
version_id: None,
error: DeletionError::AccessDenied,
};
if let DeletionEvent::ObjectFailed { key, error, .. } = &event {
assert_eq!(key, "test/key.txt");
assert_eq!(*error, DeletionError::AccessDenied);
} else {
panic!("Expected ObjectFailed event");
}
}
#[test]
fn deletion_event_pipeline_end() {
let event = DeletionEvent::PipelineEnd;
assert_eq!(event, DeletionEvent::PipelineEnd);
}
#[test]
fn deletion_event_pipeline_error() {
let event = DeletionEvent::PipelineError {
message: "something went wrong".to_string(),
};
if let DeletionEvent::PipelineError { message } = &event {
assert_eq!(message, "something went wrong");
} else {
panic!("Expected PipelineError event");
}
}
#[test]
fn deletion_event_clone() {
let event = DeletionEvent::ObjectDeleted {
key: "key".to_string(),
version_id: None,
size: 42,
};
let cloned = event.clone();
assert_eq!(event, cloned);
}
#[test]
fn s3_target_parse_bucket_only() {
let target = S3Target::parse("s3://my-bucket").unwrap();
assert_eq!(target.bucket, "my-bucket");
assert!(target.prefix.is_none());
assert!(target.endpoint.is_none());
assert!(target.region.is_none());
}
#[test]
fn s3_target_parse_bucket_with_trailing_slash() {
let target = S3Target::parse("s3://my-bucket/").unwrap();
assert_eq!(target.bucket, "my-bucket");
assert!(target.prefix.is_none());
}
#[test]
fn s3_target_parse_bucket_with_prefix() {
let target = S3Target::parse("s3://my-bucket/logs/2023/").unwrap();
assert_eq!(target.bucket, "my-bucket");
assert_eq!(target.prefix.as_deref(), Some("logs/2023/"));
}
#[test]
fn s3_target_parse_bucket_with_simple_prefix() {
let target = S3Target::parse("s3://my-bucket/prefix").unwrap();
assert_eq!(target.bucket, "my-bucket");
assert_eq!(target.prefix.as_deref(), Some("prefix"));
}
#[test]
fn s3_target_parse_bucket_with_deep_prefix() {
let target = S3Target::parse("s3://my-bucket/a/b/c/d/e").unwrap();
assert_eq!(target.bucket, "my-bucket");
assert_eq!(target.prefix.as_deref(), Some("a/b/c/d/e"));
}
#[test]
fn s3_target_parse_invalid_no_scheme() {
let result = S3Target::parse("my-bucket/prefix");
assert!(result.is_err());
let err_msg = result.unwrap_err().to_string();
assert!(err_msg.contains("Target URI must start with 's3://'"));
}
#[test]
fn s3_target_parse_invalid_wrong_scheme() {
let result = S3Target::parse("http://my-bucket/prefix");
assert!(result.is_err());
}
#[test]
fn s3_target_parse_invalid_empty_bucket() {
let result = S3Target::parse("s3://");
assert!(result.is_err());
let err_msg = result.unwrap_err().to_string();
assert!(err_msg.contains("Bucket name cannot be empty"));
}
#[test]
fn s3_target_parse_invalid_empty_bucket_with_prefix() {
let result = S3Target::parse("s3:///prefix");
assert!(result.is_err());
let err_msg = result.unwrap_err().to_string();
assert!(err_msg.contains("Bucket name cannot be empty"));
}
#[test]
fn s3_target_display_bucket_only() {
let target = S3Target {
bucket: "my-bucket".to_string(),
prefix: None,
endpoint: None,
region: None,
};
assert_eq!(target.to_string(), "s3://my-bucket");
}
#[test]
fn s3_target_display_with_prefix() {
let target = S3Target {
bucket: "my-bucket".to_string(),
prefix: Some("logs/2023/".to_string()),
endpoint: None,
region: None,
};
assert_eq!(target.to_string(), "s3://my-bucket/logs/2023/");
}
#[test]
fn s3_target_roundtrip() {
let uri = "s3://my-bucket/some/prefix/";
let target = S3Target::parse(uri).unwrap();
assert_eq!(target.to_string(), uri);
}
#[test]
fn s3_target_clone_and_eq() {
let target = S3Target::parse("s3://bucket/key").unwrap();
let cloned = target.clone();
assert_eq!(target, cloned);
}
}