#[cfg(test)]
mod tests {
use crate::config::args::parse_from_args;
use crate::config::{Config, FilterConfig, ForceRetryConfig};
use crate::types::{DeletionError, StoragePath};
use proptest::prelude::*;
fn make_config(if_match: bool) -> Config {
Config {
target: StoragePath::S3 {
bucket: "test-bucket".to_string(),
prefix: "prefix/".to_string(),
},
show_no_progress: false,
target_client_config: None,
force_retry_config: ForceRetryConfig {
force_retry_count: 0,
force_retry_interval_milliseconds: 0,
},
tracing_config: None,
worker_size: 4,
warn_as_error: false,
dry_run: false,
rate_limit_objects: None,
max_parallel_listings: 1,
object_listing_queue_size: 1000,
max_parallel_listing_max_depth: 0,
allow_parallel_listings_in_express_one_zone: false,
filter_config: FilterConfig::default(),
max_keys: 1000,
auto_complete_shell: None,
event_callback_lua_script: None,
filter_callback_lua_script: None,
allow_lua_os_library: false,
allow_lua_unsafe_vm: false,
lua_vm_memory_limit: 0,
lua_callback_timeout_milliseconds: 0,
if_match,
max_delete: None,
filter_manager: crate::callback::filter_manager::FilterManager::new(),
event_manager: crate::callback::event_manager::EventManager::new(),
batch_size: 1000,
delete_all_versions: false,
force: false,
test_user_defined_callback: false,
}
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(100))]
#[test]
fn prop_single_deleter_if_match_etag_inclusion(
if_match_enabled in proptest::bool::ANY,
has_etag in proptest::bool::ANY,
etag_value in "[a-f0-9]{32}",
) {
use aws_sdk_s3::primitives::DateTime;
use aws_sdk_s3::types::Object;
use crate::types::S3Object;
use crate::deleter::single::SingleDeleter;
use crate::deleter::Deleter;
let config = make_config(if_match_enabled);
let mut builder = Object::builder()
.key("test/key")
.size(100)
.last_modified(DateTime::from_secs(1000));
if has_etag {
builder = builder.e_tag(format!("\"{}\"", etag_value));
}
let obj = S3Object::NotVersioning(builder.build());
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
use std::sync::{Arc, Mutex};
use crate::types::DeletionStatistics;
let (stats_sender, _stats_rx) = async_channel::unbounded::<DeletionStatistics>();
let recorded_if_match: Arc<Mutex<Option<Option<String>>>> =
Arc::new(Mutex::new(None));
let recorded_clone = recorded_if_match.clone();
let mock = RecordingMockStorage {
stats_sender: stats_sender.clone(),
recorded_if_match: recorded_clone,
};
let boxed: Box<dyn crate::storage::StorageTrait + Send + Sync> = Box::new(mock);
let deleter = SingleDeleter::new(boxed);
let result = deleter.delete(&[obj], &config).await;
prop_assert!(result.is_ok(), "Deletion should succeed");
let captured = recorded_if_match.lock().unwrap();
let captured_if_match = captured.as_ref().expect("delete_object should have been called");
if if_match_enabled && has_etag {
prop_assert!(
captured_if_match.is_some(),
"When if_match=true and object has ETag, If-Match must be set"
);
let expected_etag = format!("\"{}\"", etag_value);
prop_assert_eq!(
captured_if_match.as_deref().unwrap(),
expected_etag.as_str(),
"ETag value must match the object's ETag"
);
} else {
prop_assert!(
captured_if_match.is_none(),
"When if_match=false or object has no ETag, If-Match must be None"
);
}
Ok(())
})?;
}
}
#[test]
fn precondition_failed_is_not_retryable() {
assert!(
!DeletionError::PreconditionFailed.is_retryable(),
"PreconditionFailed must not be retryable — object should be skipped"
);
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(100))]
#[test]
fn prop_if_match_flag_parsing(include_flag in proptest::bool::ANY) {
let _ = tracing_subscriber::fmt()
.with_env_filter("dummy=trace")
.try_init();
let mut args = vec!["s3rm".to_string(), "s3://bucket/prefix/".to_string()];
if include_flag {
args.push("--if-match".to_string());
}
let args_ref: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
let parsed = parse_from_args(args_ref).unwrap();
let config = Config::try_from(parsed).unwrap();
prop_assert_eq!(
config.if_match, include_flag,
"Config.if_match must match whether --if-match flag was provided"
);
}
#[test]
fn prop_if_match_default_is_false(
bucket in "[a-z][a-z0-9-]{2,10}",
) {
let _ = tracing_subscriber::fmt()
.with_env_filter("dummy=trace")
.try_init();
let target = format!("s3://{}/", bucket);
let args = vec!["s3rm", target.as_str()];
let parsed = parse_from_args(args).unwrap();
let config = Config::try_from(parsed).unwrap();
prop_assert!(!config.if_match, "Default if_match must be false when flag is absent");
}
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(50))]
#[test]
fn prop_batch_deleter_if_match_per_object_etags(
if_match_enabled in proptest::bool::ANY,
n_objects in 1usize..10,
etag_base in "[a-f0-9]{8}",
) {
use aws_sdk_s3::primitives::DateTime;
use aws_sdk_s3::types::{ObjectVersion, ObjectVersionStorageClass};
use crate::types::S3Object;
use crate::deleter::batch::BatchDeleter;
use crate::deleter::Deleter;
let config = make_config(if_match_enabled);
let objects: Vec<S3Object> = (0..n_objects)
.map(|i| {
S3Object::Versioning(
ObjectVersion::builder()
.key(format!("key/{}", i))
.version_id(format!("v{}", i))
.size(100)
.is_latest(true)
.storage_class(ObjectVersionStorageClass::Standard)
.last_modified(DateTime::from_secs(1000))
.e_tag(format!("\"{}{}\"", etag_base, i))
.build(),
)
})
.collect();
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
use std::sync::{Arc, Mutex};
use aws_sdk_s3::types::ObjectIdentifier;
use crate::types::DeletionStatistics;
let (stats_sender, _stats_rx) = async_channel::unbounded::<DeletionStatistics>();
let recorded_identifiers: Arc<Mutex<Vec<Vec<ObjectIdentifier>>>> =
Arc::new(Mutex::new(Vec::new()));
let recorded_clone = recorded_identifiers.clone();
let mock = BatchRecordingMockStorage {
stats_sender: stats_sender.clone(),
recorded_identifiers: recorded_clone,
};
let boxed: Box<dyn crate::storage::StorageTrait + Send + Sync> = Box::new(mock);
let deleter = BatchDeleter::new(boxed);
let result = deleter.delete(&objects, &config).await;
prop_assert!(result.is_ok(), "Batch deletion should succeed");
let batches = recorded_identifiers.lock().unwrap();
prop_assert!(!batches.is_empty(), "At least one batch should be sent");
let all_idents: Vec<&ObjectIdentifier> =
batches.iter().flat_map(|b| b.iter()).collect();
prop_assert_eq!(all_idents.len(), n_objects, "All objects should be in batches");
for (i, ident) in all_idents.iter().enumerate() {
let expected_key = format!("key/{}", i);
prop_assert_eq!(ident.key(), expected_key.as_str());
if if_match_enabled {
let expected_etag = format!("\"{}{}\"", etag_base, i);
prop_assert_eq!(
ident.e_tag(),
Some(expected_etag.as_str()),
"When if_match=true, ETag must be included for object {}",
i
);
} else {
prop_assert_eq!(
ident.e_tag(),
None,
"When if_match=false, ETag must be None for object {}",
i
);
}
}
Ok(())
})?;
}
#[test]
fn prop_batch_if_match_respects_batch_size(
batch_size in 2u16..10,
extra_objects in 1usize..8,
) {
use aws_sdk_s3::primitives::DateTime;
use aws_sdk_s3::types::{ObjectVersion, ObjectVersionStorageClass};
use crate::types::S3Object;
use crate::deleter::batch::BatchDeleter;
use crate::deleter::Deleter;
let mut config = make_config(true); config.batch_size = batch_size;
let n_objects = batch_size as usize + extra_objects;
let objects: Vec<S3Object> = (0..n_objects)
.map(|i| {
S3Object::Versioning(
ObjectVersion::builder()
.key(format!("key/{}", i))
.version_id(format!("v{}", i))
.size(50)
.is_latest(true)
.storage_class(ObjectVersionStorageClass::Standard)
.last_modified(DateTime::from_secs(2000))
.e_tag(format!("\"etag{}\"", i))
.build(),
)
})
.collect();
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
use std::sync::{Arc, Mutex};
use aws_sdk_s3::types::ObjectIdentifier;
use crate::types::DeletionStatistics;
let (stats_sender, _stats_rx) = async_channel::unbounded::<DeletionStatistics>();
let recorded_identifiers: Arc<Mutex<Vec<Vec<ObjectIdentifier>>>> =
Arc::new(Mutex::new(Vec::new()));
let recorded_clone = recorded_identifiers.clone();
let mock = BatchRecordingMockStorage {
stats_sender: stats_sender.clone(),
recorded_identifiers: recorded_clone,
};
let boxed: Box<dyn crate::storage::StorageTrait + Send + Sync> = Box::new(mock);
let deleter = BatchDeleter::new(boxed);
let result = deleter.delete(&objects, &config).await;
prop_assert!(result.is_ok(), "Batch deletion should succeed");
let batches = recorded_identifiers.lock().unwrap();
let expected_batches = n_objects.div_ceil(batch_size as usize);
prop_assert_eq!(
batches.len(),
expected_batches,
"Objects should be split into correct number of batches"
);
for batch in batches.iter() {
prop_assert!(
batch.len() <= batch_size as usize,
"Batch size {} exceeds configured max {}",
batch.len(),
batch_size
);
}
let all_idents: Vec<&ObjectIdentifier> =
batches.iter().flat_map(|b| b.iter()).collect();
prop_assert_eq!(all_idents.len(), n_objects);
for ident in &all_idents {
prop_assert!(
ident.e_tag().is_some(),
"Every identifier must include ETag when if_match=true"
);
}
Ok(())
})?;
}
}
use async_channel::Sender;
use async_trait::async_trait;
use std::sync::{Arc, Mutex};
use aws_sdk_s3::operation::delete_object::DeleteObjectOutput;
use aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput;
use aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput;
use aws_sdk_s3::operation::head_object::HeadObjectOutput;
use aws_sdk_s3::types::{DeletedObject, ObjectIdentifier};
use crate::storage::StorageTrait;
use crate::types::{DeletionStatistics, S3Object};
#[derive(Clone)]
struct RecordingMockStorage {
stats_sender: Sender<DeletionStatistics>,
recorded_if_match: Arc<Mutex<Option<Option<String>>>>,
}
#[async_trait]
impl StorageTrait for RecordingMockStorage {
fn is_express_onezone_storage(&self) -> bool {
false
}
async fn list_objects(
&self,
_sender: &Sender<S3Object>,
_max_keys: i32,
) -> anyhow::Result<()> {
Ok(())
}
async fn list_object_versions(
&self,
_sender: &Sender<S3Object>,
_max_keys: i32,
) -> anyhow::Result<()> {
Ok(())
}
async fn head_object(
&self,
_key: &str,
_version_id: Option<String>,
) -> anyhow::Result<HeadObjectOutput> {
Ok(HeadObjectOutput::builder().build())
}
async fn get_object_tagging(
&self,
_key: &str,
_version_id: Option<String>,
) -> anyhow::Result<GetObjectTaggingOutput> {
Ok(GetObjectTaggingOutput::builder()
.set_tag_set(Some(vec![]))
.build()
.unwrap())
}
async fn delete_object(
&self,
_key: &str,
_version_id: Option<String>,
if_match: Option<String>,
) -> anyhow::Result<DeleteObjectOutput> {
*self.recorded_if_match.lock().unwrap() = Some(if_match);
Ok(DeleteObjectOutput::builder().build())
}
async fn delete_objects(
&self,
objects: Vec<ObjectIdentifier>,
) -> anyhow::Result<DeleteObjectsOutput> {
let mut builder = DeleteObjectsOutput::builder();
for ident in &objects {
builder = builder.deleted(DeletedObject::builder().key(ident.key()).build());
}
Ok(builder.build())
}
async fn is_versioning_enabled(&self) -> anyhow::Result<bool> {
Ok(false)
}
fn get_client(&self) -> Option<Arc<aws_sdk_s3::Client>> {
None
}
fn get_stats_sender(&self) -> Sender<DeletionStatistics> {
self.stats_sender.clone()
}
async fn send_stats(&self, stats: DeletionStatistics) {
let _ = self.stats_sender.send(stats).await;
}
fn set_warning(&self) {}
}
#[derive(Clone)]
struct BatchRecordingMockStorage {
stats_sender: Sender<DeletionStatistics>,
recorded_identifiers: Arc<Mutex<Vec<Vec<ObjectIdentifier>>>>,
}
#[async_trait]
impl StorageTrait for BatchRecordingMockStorage {
fn is_express_onezone_storage(&self) -> bool {
false
}
async fn list_objects(
&self,
_sender: &Sender<S3Object>,
_max_keys: i32,
) -> anyhow::Result<()> {
Ok(())
}
async fn list_object_versions(
&self,
_sender: &Sender<S3Object>,
_max_keys: i32,
) -> anyhow::Result<()> {
Ok(())
}
async fn head_object(
&self,
_key: &str,
_version_id: Option<String>,
) -> anyhow::Result<HeadObjectOutput> {
Ok(HeadObjectOutput::builder().build())
}
async fn get_object_tagging(
&self,
_key: &str,
_version_id: Option<String>,
) -> anyhow::Result<GetObjectTaggingOutput> {
Ok(GetObjectTaggingOutput::builder()
.set_tag_set(Some(vec![]))
.build()
.unwrap())
}
async fn delete_object(
&self,
_key: &str,
_version_id: Option<String>,
_if_match: Option<String>,
) -> anyhow::Result<DeleteObjectOutput> {
Ok(DeleteObjectOutput::builder().build())
}
async fn delete_objects(
&self,
objects: Vec<ObjectIdentifier>,
) -> anyhow::Result<DeleteObjectsOutput> {
self.recorded_identifiers
.lock()
.unwrap()
.push(objects.clone());
let mut builder = DeleteObjectsOutput::builder();
for ident in &objects {
builder = builder.deleted(DeletedObject::builder().key(ident.key()).build());
}
Ok(builder.build())
}
async fn is_versioning_enabled(&self) -> anyhow::Result<bool> {
Ok(false)
}
fn get_client(&self) -> Option<Arc<aws_sdk_s3::Client>> {
None
}
fn get_stats_sender(&self) -> Sender<DeletionStatistics> {
self.stats_sender.clone()
}
async fn send_stats(&self, stats: DeletionStatistics) {
let _ = self.stats_sender.send(stats).await;
}
fn set_warning(&self) {}
}
fn make_recording_mock() -> (
RecordingMockStorage,
async_channel::Receiver<DeletionStatistics>,
) {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let mock = RecordingMockStorage {
stats_sender,
recorded_if_match: Arc::new(Mutex::new(None)),
};
(mock, stats_receiver)
}
#[test]
fn recording_mock_is_express_onezone_returns_false() {
let (mock, _) = make_recording_mock();
assert!(!mock.is_express_onezone_storage());
}
#[tokio::test]
async fn recording_mock_list_objects_returns_ok() {
let (mock, _) = make_recording_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_objects(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn recording_mock_list_object_versions_returns_ok() {
let (mock, _) = make_recording_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_object_versions(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn recording_mock_head_object_returns_ok() {
let (mock, _) = make_recording_mock();
assert!(mock.head_object("key", None).await.is_ok());
}
#[tokio::test]
async fn recording_mock_get_object_tagging_returns_ok() {
let (mock, _) = make_recording_mock();
let result = mock.get_object_tagging("key", None).await;
assert!(result.is_ok());
assert!(result.unwrap().tag_set().is_empty());
}
#[tokio::test]
async fn recording_mock_delete_object_records_if_match() {
let (mock, _) = make_recording_mock();
let result = mock
.delete_object("key", None, Some("etag-123".to_string()))
.await;
assert!(result.is_ok());
let recorded = mock.recorded_if_match.lock().unwrap();
assert_eq!(*recorded, Some(Some("etag-123".to_string())));
}
#[tokio::test]
async fn recording_mock_delete_object_records_none_if_match() {
let (mock, _) = make_recording_mock();
let result = mock.delete_object("key", None, None).await;
assert!(result.is_ok());
let recorded = mock.recorded_if_match.lock().unwrap();
assert_eq!(*recorded, Some(None));
}
#[tokio::test]
async fn recording_mock_delete_objects_returns_deleted_keys() {
let (mock, _) = make_recording_mock();
let idents = vec![
ObjectIdentifier::builder().key("a.txt").build().unwrap(),
ObjectIdentifier::builder().key("b.txt").build().unwrap(),
];
let result = mock.delete_objects(idents).await;
assert!(result.is_ok());
let output = result.unwrap();
let deleted: Vec<&str> = output.deleted().iter().filter_map(|d| d.key()).collect();
assert_eq!(deleted, vec!["a.txt", "b.txt"]);
}
#[tokio::test]
async fn recording_mock_is_versioning_enabled_returns_false() {
let (mock, _) = make_recording_mock();
assert!(!mock.is_versioning_enabled().await.unwrap());
}
#[test]
fn recording_mock_get_client_returns_none() {
let (mock, _) = make_recording_mock();
assert!(mock.get_client().is_none());
}
#[tokio::test]
async fn recording_mock_get_stats_sender_works() {
let (mock, stats_receiver) = make_recording_mock();
let sender = mock.get_stats_sender();
sender
.send(DeletionStatistics::DeleteBytes(42))
.await
.unwrap();
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(received, DeletionStatistics::DeleteBytes(42)));
}
#[tokio::test]
async fn recording_mock_send_stats_delivers_stat() {
let (mock, stats_receiver) = make_recording_mock();
mock.send_stats(DeletionStatistics::DeleteComplete {
key: "k".to_string(),
})
.await;
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(
received,
DeletionStatistics::DeleteComplete { .. }
));
}
#[test]
fn recording_mock_set_warning_is_callable() {
let (mock, _) = make_recording_mock();
mock.set_warning(); }
fn make_batch_recording_mock() -> (
BatchRecordingMockStorage,
async_channel::Receiver<DeletionStatistics>,
) {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let mock = BatchRecordingMockStorage {
stats_sender,
recorded_identifiers: Arc::new(Mutex::new(Vec::new())),
};
(mock, stats_receiver)
}
#[test]
fn batch_recording_mock_is_express_onezone_returns_false() {
let (mock, _) = make_batch_recording_mock();
assert!(!mock.is_express_onezone_storage());
}
#[tokio::test]
async fn batch_recording_mock_list_objects_returns_ok() {
let (mock, _) = make_batch_recording_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_objects(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn batch_recording_mock_list_object_versions_returns_ok() {
let (mock, _) = make_batch_recording_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_object_versions(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn batch_recording_mock_head_object_returns_ok() {
let (mock, _) = make_batch_recording_mock();
assert!(mock.head_object("key", None).await.is_ok());
}
#[tokio::test]
async fn batch_recording_mock_get_object_tagging_returns_ok() {
let (mock, _) = make_batch_recording_mock();
let result = mock.get_object_tagging("key", None).await;
assert!(result.is_ok());
assert!(result.unwrap().tag_set().is_empty());
}
#[tokio::test]
async fn batch_recording_mock_delete_object_returns_ok() {
let (mock, _) = make_batch_recording_mock();
assert!(mock.delete_object("key", None, None).await.is_ok());
}
#[tokio::test]
async fn batch_recording_mock_delete_objects_records_and_returns() {
let (mock, _) = make_batch_recording_mock();
let idents = vec![
ObjectIdentifier::builder().key("x.txt").build().unwrap(),
ObjectIdentifier::builder().key("y.txt").build().unwrap(),
ObjectIdentifier::builder().key("z.txt").build().unwrap(),
];
let result = mock.delete_objects(idents).await;
assert!(result.is_ok());
let recorded = mock.recorded_identifiers.lock().unwrap();
assert_eq!(recorded.len(), 1);
let keys: Vec<&str> = recorded[0].iter().map(|i| i.key()).collect();
assert_eq!(keys, vec!["x.txt", "y.txt", "z.txt"]);
let output = result.unwrap();
let deleted: Vec<&str> = output.deleted().iter().filter_map(|d| d.key()).collect();
assert_eq!(deleted, vec!["x.txt", "y.txt", "z.txt"]);
}
#[tokio::test]
async fn batch_recording_mock_is_versioning_enabled_returns_false() {
let (mock, _) = make_batch_recording_mock();
assert!(!mock.is_versioning_enabled().await.unwrap());
}
#[test]
fn batch_recording_mock_get_client_returns_none() {
let (mock, _) = make_batch_recording_mock();
assert!(mock.get_client().is_none());
}
#[tokio::test]
async fn batch_recording_mock_get_stats_sender_works() {
let (mock, stats_receiver) = make_batch_recording_mock();
let sender = mock.get_stats_sender();
sender
.send(DeletionStatistics::DeleteBytes(99))
.await
.unwrap();
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(received, DeletionStatistics::DeleteBytes(99)));
}
#[tokio::test]
async fn batch_recording_mock_send_stats_delivers_stat() {
let (mock, stats_receiver) = make_batch_recording_mock();
mock.send_stats(DeletionStatistics::DeleteComplete {
key: "k".to_string(),
})
.await;
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(
received,
DeletionStatistics::DeleteComplete { .. }
));
}
#[test]
fn batch_recording_mock_set_warning_is_callable() {
let (mock, _) = make_batch_recording_mock();
mock.set_warning(); }
}