use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use anyhow::Result;
use async_channel::Sender;
use async_trait::async_trait;
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::primitives::DateTime;
use aws_sdk_s3::types::{
DeletedObject, Object, ObjectIdentifier, ObjectVersion, ObjectVersionStorageClass, Tag,
};
use fancy_regex::Regex;
use proptest::prelude::*;
use tokio_util::sync::CancellationToken;
use super::*;
use crate::config::Config;
use crate::stage::Stage;
use crate::storage::StorageTrait;
use crate::test_utils::{
init_dummy_tracing_subscriber, make_s3_object, make_test_config, make_versioned_s3_object,
};
use crate::types::token::PipelineCancellationToken;
use crate::types::{DeletionStatistics, S3Object};
use single::extract_delete_object_error_code;
#[derive(Debug, Clone)]
#[allow(dead_code)]
struct DeleteObjectCall {
key: String,
version_id: Option<String>,
if_match: Option<String>,
}
#[derive(Debug, Clone)]
struct DeleteObjectsCall {
identifiers: Vec<ObjectIdentifier>,
}
#[derive(Debug, Clone)]
#[allow(dead_code)]
struct HeadObjectCall {
key: String,
version_id: Option<String>,
}
#[derive(Clone)]
struct MockStorage {
stats_sender: Sender<DeletionStatistics>,
delete_object_calls: Arc<Mutex<Vec<DeleteObjectCall>>>,
delete_objects_calls: Arc<Mutex<Vec<DeleteObjectsCall>>>,
head_object_calls: Arc<Mutex<Vec<HeadObjectCall>>>,
delete_object_error_keys: Arc<Mutex<HashMap<String, String>>>,
head_object_content_type: Arc<Mutex<Option<String>>>,
head_object_metadata: Arc<Mutex<Option<HashMap<String, String>>>>,
tagging_response_tags: Arc<Mutex<Option<Vec<Tag>>>>,
batch_error_keys: Arc<Mutex<HashMap<String, String>>>,
head_object_error: Arc<Mutex<Option<anyhow::Error>>>,
tagging_error: Arc<Mutex<Option<anyhow::Error>>>,
delete_object_sdk_error: Arc<Mutex<Option<anyhow::Error>>>,
}
impl MockStorage {
fn new(stats_sender: Sender<DeletionStatistics>) -> Self {
Self {
stats_sender,
delete_object_calls: Arc::new(Mutex::new(Vec::new())),
delete_objects_calls: Arc::new(Mutex::new(Vec::new())),
head_object_calls: Arc::new(Mutex::new(Vec::new())),
delete_object_error_keys: Arc::new(Mutex::new(HashMap::new())),
head_object_content_type: Arc::new(Mutex::new(None)),
head_object_metadata: Arc::new(Mutex::new(None)),
tagging_response_tags: Arc::new(Mutex::new(None)),
batch_error_keys: Arc::new(Mutex::new(HashMap::new())),
head_object_error: Arc::new(Mutex::new(None)),
tagging_error: Arc::new(Mutex::new(None)),
delete_object_sdk_error: Arc::new(Mutex::new(None)),
}
}
}
#[async_trait]
impl StorageTrait for MockStorage {
fn is_express_onezone_storage(&self) -> bool {
false
}
async fn list_objects(&self, _sender: &Sender<S3Object>, _max_keys: i32) -> Result<()> {
Ok(())
}
async fn list_object_versions(&self, _sender: &Sender<S3Object>, _max_keys: i32) -> Result<()> {
Ok(())
}
async fn head_object(
&self,
relative_key: &str,
version_id: Option<String>,
) -> Result<HeadObjectOutput> {
self.head_object_calls.lock().unwrap().push(HeadObjectCall {
key: relative_key.to_string(),
version_id,
});
if let Some(err) = self.head_object_error.lock().unwrap().take() {
return Err(err);
}
let mut builder = HeadObjectOutput::builder();
if let Some(ct) = self.head_object_content_type.lock().unwrap().as_ref() {
builder = builder.content_type(ct.clone());
}
if let Some(meta) = self.head_object_metadata.lock().unwrap().as_ref() {
for (k, v) in meta {
builder = builder.metadata(k.clone(), v.clone());
}
}
Ok(builder.build())
}
async fn get_object_tagging(
&self,
_relative_key: &str,
_version_id: Option<String>,
) -> Result<GetObjectTaggingOutput> {
if let Some(err) = self.tagging_error.lock().unwrap().take() {
return Err(err);
}
let mut builder = GetObjectTaggingOutput::builder();
let tags_guard = self.tagging_response_tags.lock().unwrap();
match tags_guard.as_ref() {
Some(tags) if !tags.is_empty() => {
for tag in tags {
builder = builder.tag_set(tag.clone());
}
}
_ => {
builder = builder.set_tag_set(Some(vec![]));
}
}
Ok(builder.build().unwrap())
}
async fn delete_object(
&self,
relative_key: &str,
version_id: Option<String>,
if_match: Option<String>,
) -> Result<DeleteObjectOutput> {
self.delete_object_calls
.lock()
.unwrap()
.push(DeleteObjectCall {
key: relative_key.to_string(),
version_id: version_id.clone(),
if_match: if_match.clone(),
});
if let Some(err) = self.delete_object_sdk_error.lock().unwrap().take() {
return Err(err);
}
let error_keys = self.delete_object_error_keys.lock().unwrap();
if let Some(msg) = error_keys.get(relative_key) {
return Err(anyhow::anyhow!("{}", msg));
}
Ok(DeleteObjectOutput::builder().build())
}
async fn delete_objects(&self, objects: Vec<ObjectIdentifier>) -> Result<DeleteObjectsOutput> {
let batch_error_keys = self.batch_error_keys.lock().unwrap().clone();
self.delete_objects_calls
.lock()
.unwrap()
.push(DeleteObjectsCall {
identifiers: objects.clone(),
});
let mut builder = DeleteObjectsOutput::builder();
for ident in &objects {
let key = ident.key();
let version_id = ident.version_id();
if let Some(error_code) = batch_error_keys.get(key) {
let mut err_builder = aws_sdk_s3::types::Error::builder()
.key(key)
.code(error_code.as_str())
.message(format!("{} error", error_code));
if let Some(vid) = version_id {
err_builder = err_builder.version_id(vid);
}
builder = builder.errors(err_builder.build());
} else {
let mut del_builder = DeletedObject::builder().key(key);
if let Some(vid) = version_id {
del_builder = del_builder.version_id(vid);
}
builder = builder.deleted(del_builder.build());
}
}
Ok(builder.build())
}
async fn is_versioning_enabled(&self) -> 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_mock_storage_boxed(
stats_sender: Sender<DeletionStatistics>,
) -> (Box<dyn StorageTrait + Send + Sync>, MockStorage) {
let mock = MockStorage::new(stats_sender.clone());
let boxed: Box<dyn StorageTrait + Send + Sync> = Box::new(mock.clone());
(boxed, mock)
}
fn make_stage_with_mock(
config: Config,
mock_storage: Box<dyn StorageTrait + Send + Sync>,
receiver: Option<async_channel::Receiver<S3Object>>,
sender: Option<Sender<S3Object>>,
) -> Stage {
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let has_warning = Arc::new(std::sync::atomic::AtomicBool::new(false));
Stage::new(
config,
mock_storage,
receiver,
sender,
cancellation_token,
has_warning,
)
}
fn make_stage_with_mock_and_token(
config: Config,
mock_storage: Box<dyn StorageTrait + Send + Sync>,
receiver: Option<async_channel::Receiver<S3Object>>,
sender: Option<Sender<S3Object>>,
cancellation_token: PipelineCancellationToken,
) -> Stage {
let has_warning = Arc::new(std::sync::atomic::AtomicBool::new(false));
Stage::new(
config,
mock_storage,
receiver,
sender,
cancellation_token,
has_warning,
)
}
#[test]
fn format_metadata_empty() {
init_dummy_tracing_subscriber();
let meta: HashMap<String, String> = HashMap::new();
assert_eq!(format_metadata(&meta), "");
}
#[test]
fn format_metadata_single_entry() {
let mut meta = HashMap::new();
meta.insert("key1".to_string(), "value1".to_string());
assert_eq!(format_metadata(&meta), "key1=value1");
}
#[test]
fn format_metadata_multiple_entries_sorted() {
let mut meta = HashMap::new();
meta.insert("zebra".to_string(), "z_val".to_string());
meta.insert("alpha".to_string(), "a_val".to_string());
meta.insert("middle".to_string(), "m_val".to_string());
let result = format_metadata(&meta);
assert_eq!(result, "alpha=a_val,middle=m_val,zebra=z_val");
}
#[test]
fn format_metadata_special_chars_encoded() {
let mut meta = HashMap::new();
meta.insert("key with spaces".to_string(), "val&ue".to_string());
let result = format_metadata(&meta);
assert!(result.contains("key with spaces=val%26ue"));
}
#[test]
fn format_tags_empty() {
init_dummy_tracing_subscriber();
let tags: Vec<Tag> = vec![];
assert_eq!(format_tags(&tags), "");
}
#[test]
fn format_tags_single_tag() {
let tags = vec![Tag::builder().key("env").value("prod").build().unwrap()];
assert_eq!(format_tags(&tags), "env=prod");
}
#[test]
fn format_tags_multiple_sorted() {
let tags = vec![
Tag::builder().key("z-tag").value("zval").build().unwrap(),
Tag::builder().key("a-tag").value("aval").build().unwrap(),
];
let result = format_tags(&tags);
assert_eq!(result, "a-tag=aval&z-tag=zval");
}
#[test]
fn format_tags_special_chars_encoded() {
let tags = vec![
Tag::builder()
.key("tag key")
.value("tag value&more")
.build()
.unwrap(),
];
let result = format_tags(&tags);
assert!(result.contains("tag%20key"));
assert!(result.contains("tag%20value%26more"));
}
#[test]
fn generate_tagging_string_none() {
assert!(generate_tagging_string(&None).is_none());
}
#[test]
fn generate_tagging_string_some() {
let output = GetObjectTaggingOutput::builder()
.tag_set(Tag::builder().key("env").value("dev").build().unwrap())
.tag_set(Tag::builder().key("app").value("test").build().unwrap())
.build()
.unwrap();
let result = generate_tagging_string(&Some(output));
assert!(result.is_some());
let s = result.unwrap();
assert!(s.contains("env=dev"));
assert!(s.contains("app=test"));
}
#[tokio::test]
async fn batch_deleter_empty_list() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let deleter = BatchDeleter::new(boxed);
let config = make_test_config();
let result = deleter.delete(&[], &config).await.unwrap();
assert_eq!(result.deleted.len(), 0);
}
#[tokio::test]
async fn batch_deleter_single_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = BatchDeleter::new(boxed);
let config = make_test_config();
let objects = vec![make_s3_object("test/key.txt", 1024)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 1);
let calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].identifiers.len(), 1);
assert_eq!(calls[0].identifiers[0].key(), "test/key.txt");
}
#[tokio::test]
async fn batch_deleter_respects_batch_size_config() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = BatchDeleter::new(boxed);
let mut config = make_test_config();
config.batch_size = 2;
let objects: Vec<S3Object> = (0..5)
.map(|i| make_s3_object(&format!("key/{i}"), 100))
.collect();
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 5);
let calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(calls.len(), 3);
assert_eq!(calls[0].identifiers.len(), 2);
assert_eq!(calls[1].identifiers.len(), 2);
assert_eq!(calls[2].identifiers.len(), 1);
}
#[tokio::test]
async fn batch_deleter_max_batch_size_enforced() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = BatchDeleter::new(boxed);
let mut config = make_test_config();
config.batch_size = u16::MAX;
let objects: Vec<S3Object> = (0..1500)
.map(|i| make_s3_object(&format!("key/{i}"), 10))
.collect();
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 1500);
let calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(calls.len(), 2);
assert_eq!(calls[0].identifiers.len(), 1000);
assert_eq!(calls[1].identifiers.len(), 500);
}
#[tokio::test]
async fn batch_deleter_with_version_ids() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = BatchDeleter::new(boxed);
let config = make_test_config();
let objects = vec![
make_versioned_s3_object("key/a", "v1", 100),
make_versioned_s3_object("key/b", "v2", 200),
];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 2);
let calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(calls.len(), 1);
let idents = &calls[0].identifiers;
assert_eq!(idents[0].key(), "key/a");
assert_eq!(idents[0].version_id(), Some("v1"));
assert_eq!(idents[1].key(), "key/b");
assert_eq!(idents[1].version_id(), Some("v2"));
}
#[tokio::test]
async fn batch_deleter_partial_failure() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/1".to_string(), "AccessDenied".to_string());
let deleter = BatchDeleter::new(boxed);
let config = make_test_config();
let objects = vec![
make_s3_object("key/0", 100),
make_s3_object("key/1", 200),
make_s3_object("key/2", 300),
];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 2);
assert_eq!(result.failed.len(), 1);
assert_eq!(result.failed[0].error_code, "AccessDenied");
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 0);
}
#[tokio::test]
async fn batch_deleter_includes_etag_when_if_match() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = BatchDeleter::new(boxed);
let mut config = make_test_config();
config.if_match = true;
let obj = S3Object::Versioning(
ObjectVersion::builder()
.key("key/with-etag")
.version_id("v1")
.size(100)
.is_latest(true)
.storage_class(ObjectVersionStorageClass::Standard)
.last_modified(DateTime::from_secs(1000))
.e_tag("\"abc123\"")
.build(),
);
let result = deleter.delete(&[obj], &config).await.unwrap();
assert_eq!(result.deleted.len(), 1);
let calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].identifiers.len(), 1);
assert_eq!(calls[0].identifiers[0].key(), "key/with-etag");
assert_eq!(calls[0].identifiers[0].e_tag(), Some("\"abc123\""));
}
#[tokio::test]
async fn batch_deleter_no_etag_when_if_match_disabled() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = BatchDeleter::new(boxed);
let config = make_test_config();
let obj = S3Object::Versioning(
ObjectVersion::builder()
.key("key/with-etag")
.version_id("v1")
.size(100)
.is_latest(true)
.storage_class(ObjectVersionStorageClass::Standard)
.last_modified(DateTime::from_secs(1000))
.e_tag("\"abc123\"")
.build(),
);
let result = deleter.delete(&[obj], &config).await.unwrap();
assert_eq!(result.deleted.len(), 1);
let calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].identifiers.len(), 1);
assert_eq!(calls[0].identifiers[0].key(), "key/with-etag");
assert_eq!(calls[0].identifiers[0].e_tag(), None);
}
#[tokio::test]
async fn single_deleter_single_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = SingleDeleter::new(boxed);
let config = make_test_config();
let objects = vec![make_s3_object("key/a", 100)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 1);
assert_eq!(result.deleted[0].key, "key/a");
let calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].key, "key/a");
}
#[tokio::test]
async fn single_deleter_with_version_id() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = SingleDeleter::new(boxed);
let config = make_test_config();
let objects = vec![make_versioned_s3_object("key/versioned", "v42", 500)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 1);
let calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(calls[0].key, "key/versioned");
assert_eq!(calls[0].version_id.as_deref(), Some("v42"));
}
#[tokio::test]
async fn single_deleter_returns_error_on_failure() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.delete_object_error_keys
.lock()
.unwrap()
.insert("key/fail".to_string(), "AccessDenied".to_string());
let deleter = SingleDeleter::new(boxed);
let config = make_test_config();
let objects = vec![make_s3_object("key/fail", 200)];
let result = deleter.delete(&objects, &config).await;
assert!(result.is_ok());
let delete_result = result.unwrap();
assert_eq!(delete_result.failed.len(), 1);
assert_eq!(delete_result.failed[0].key, "key/fail");
assert!(delete_result.deleted.is_empty());
let calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].key, "key/fail");
}
#[tokio::test]
async fn single_deleter_includes_etag_when_if_match() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = SingleDeleter::new(boxed);
let mut config = make_test_config();
config.if_match = true;
let objects = vec![S3Object::NotVersioning(
Object::builder()
.key("key/with-etag")
.size(100)
.last_modified(DateTime::from_secs(1000))
.e_tag("\"abc123\"")
.build(),
)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 1);
let calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].key, "key/with-etag");
assert_eq!(calls[0].if_match.as_deref(), Some("\"abc123\""));
}
#[tokio::test]
async fn single_deleter_no_etag_when_if_match_disabled() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = SingleDeleter::new(boxed);
let config = make_test_config();
let objects = vec![S3Object::NotVersioning(
Object::builder()
.key("key/with-etag")
.size(100)
.last_modified(DateTime::from_secs(1000))
.e_tag("\"abc123\"")
.build(),
)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 1);
let calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].if_match, None);
}
#[tokio::test]
async fn object_deleter_processes_objects() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, output_receiver) = async_channel::bounded::<S3Object>(10);
let config = make_test_config();
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter.clone());
input_sender
.send(make_s3_object("test/obj1.txt", 1024))
.await
.unwrap();
input_sender
.send(make_s3_object("test/obj2.txt", 2048))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1); assert_eq!(batch_calls[0].identifiers.len(), 2);
assert_eq!(batch_calls[0].identifiers[0].key(), "test/obj1.txt");
assert_eq!(batch_calls[0].identifiers[1].key(), "test/obj2.txt");
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 0);
let report = &*stats_report;
assert_eq!(report.stats_deleted_objects.load(Ordering::SeqCst), 2);
assert_eq!(report.stats_deleted_bytes.load(Ordering::SeqCst), 3072);
drop(output_receiver); drop(stats_receiver);
}
#[tokio::test]
async fn object_deleter_max_delete_threshold() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.max_delete = Some(2);
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let stage = make_stage_with_mock_and_token(
config,
boxed,
Some(input_receiver),
Some(output_sender),
cancellation_token.clone(),
);
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter.clone());
for i in 0..5 {
input_sender
.send(make_s3_object(&format!("key/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers.len(), 2);
assert_eq!(batch_calls[0].identifiers[0].key(), "key/0");
assert_eq!(batch_calls[0].identifiers[1].key(), "key/1");
let counter_val = delete_counter.load(Ordering::SeqCst);
assert!(
counter_val <= 3,
"counter {counter_val} should not exceed max_delete(2) + 1"
);
assert!(
cancellation_token.is_cancelled(),
"cancellation token must be set when max_delete is exceeded"
);
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(50))]
#[test]
fn property_20_max_delete_cancels_pipeline(
max_delete in 1u64..20,
total_objects in 5usize..50,
) {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let capacity = total_objects + 10;
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(capacity);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(capacity);
let mut config = make_test_config();
config.max_delete = Some(max_delete);
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let stage = make_stage_with_mock_and_token(
config,
boxed,
Some(input_receiver),
Some(output_sender),
cancellation_token.clone(),
);
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(
stage,
0,
stats_report.clone(),
delete_counter.clone(),
);
for i in 0..total_objects {
input_sender
.send(make_s3_object(&format!("max-del/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
let counter_val = delete_counter.load(Ordering::SeqCst);
if total_objects as u64 > max_delete {
prop_assert!(
counter_val <= max_delete + 1,
"counter {} should not exceed max_delete({}) + 1",
counter_val,
max_delete,
);
prop_assert!(
cancellation_token.is_cancelled(),
"cancellation token must be set when max_delete is exceeded"
);
}
Ok(())
})?;
}
}
#[tokio::test]
async fn object_deleter_cancellation() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let config = make_test_config();
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let has_warning = Arc::new(std::sync::atomic::AtomicBool::new(false));
let stage = Stage::new(
config,
boxed,
Some(input_receiver),
Some(output_sender),
cancellation_token.clone(),
has_warning,
);
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, delete_counter);
cancellation_token.cancel();
input_sender
.send(make_s3_object("key/1", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
}
#[tokio::test]
async fn object_deleter_content_type_include_filter() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("text/plain".to_string());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("application/json").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("doc.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 0);
let head_calls = mock.head_object_calls.lock().unwrap();
assert_eq!(head_calls.len(), 1);
}
#[tokio::test]
async fn object_deleter_content_type_exclude_filter() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("text/plain".to_string());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.exclude_content_type_regex = Some(Regex::new("text/.*").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("doc.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 0);
}
#[tokio::test]
async fn object_deleter_metadata_include_filter() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let mut meta = HashMap::new();
meta.insert("env".to_string(), "production".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_metadata_regex = Some(Regex::new("env=production").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("data.json", 500))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "data.json");
}
#[tokio::test]
async fn object_deleter_tag_include_filter() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("env").value("dev").build().unwrap(),
]);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("env=dev").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("tagged.txt", 256))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "tagged.txt");
}
#[tokio::test]
async fn object_deleter_tag_exclude_filter() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("retain").value("true").build().unwrap(),
]);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.exclude_tag_regex = Some(Regex::new("retain=true").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("important.txt", 999))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 0);
}
#[tokio::test]
async fn object_deleter_filter_combination_and_logic() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("application/json".to_string());
let mut meta = HashMap::new();
meta.insert("env".to_string(), "staging".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("application/json").unwrap());
config.filter_config.include_metadata_regex = Some(Regex::new("env=production").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("data.json", 1000))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 0);
}
#[tokio::test]
async fn object_deleter_no_head_without_filters() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let config = make_test_config(); let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("simple.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let head_calls = mock.head_object_calls.lock().unwrap();
assert_eq!(head_calls.len(), 0);
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "simple.txt");
}
#[tokio::test]
async fn object_deleter_batch_size_1_uses_single_deleter() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.batch_size = 1; let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("key/a", 100))
.await
.unwrap();
input_sender
.send(make_s3_object("key/b", 200))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 2);
assert_eq!(single_calls[0].key, "key/a");
assert_eq!(single_calls[1].key, "key/b");
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
let report = &*stats_report;
assert_eq!(report.stats_deleted_objects.load(Ordering::SeqCst), 2);
assert_eq!(report.stats_deleted_bytes.load(Ordering::SeqCst), 300);
}
#[tokio::test]
async fn object_deleter_if_match_uses_batch_with_etags() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.if_match = true; let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
let obj = S3Object::Versioning(
ObjectVersion::builder()
.key("if-match/obj.txt")
.version_id("v1")
.size(512)
.is_latest(true)
.storage_class(ObjectVersionStorageClass::Standard)
.last_modified(DateTime::from_secs(1000))
.e_tag("\"etag123\"")
.build(),
);
input_sender.send(obj).await.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "if-match/obj.txt");
assert_eq!(batch_calls[0].identifiers[0].e_tag(), Some("\"etag123\""));
let report = &*stats_report;
assert_eq!(report.stats_deleted_objects.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn object_deleter_flushes_at_batch_size_boundary() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(20);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(20);
let mut config = make_test_config();
config.batch_size = 3; let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..9 {
input_sender
.send(make_s3_object(&format!("key/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 3);
assert_eq!(batch_calls[0].identifiers.len(), 3);
assert_eq!(batch_calls[1].identifiers.len(), 3);
assert_eq!(batch_calls[2].identifiers.len(), 3);
let report = &*stats_report;
assert_eq!(report.stats_deleted_objects.load(Ordering::SeqCst), 9);
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(50))]
#[test]
fn prop_batch_never_exceeds_max_size(
obj_count in 1usize..3000,
batch_size in 1u16..2000,
) {
let objects: Vec<S3Object> = (0..obj_count)
.map(|i| make_s3_object(&format!("key/{i}"), 100))
.collect();
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = BatchDeleter::new(boxed);
let mut config = make_test_config();
config.batch_size = batch_size;
let result = deleter.delete(&objects, &config).await.unwrap();
prop_assert_eq!(result.deleted.len(), obj_count);
let calls = mock.delete_objects_calls.lock().unwrap();
let effective_batch = (batch_size as usize).min(batch::MAX_BATCH_SIZE);
for call in calls.iter() {
prop_assert!(call.identifiers.len() <= effective_batch);
}
let total_sent: usize = calls.iter().map(|c| c.identifiers.len()).sum();
prop_assert_eq!(total_sent, obj_count);
Ok(())
})?;
}
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(50))]
#[test]
fn prop_single_deleter_one_call_per_object(
obj_count in 1usize..200,
) {
let objects: Vec<S3Object> = (0..obj_count)
.map(|i| make_s3_object(&format!("key/{i}"), (i as i64 + 1) * 100))
.collect();
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let deleter = SingleDeleter::new(boxed);
let config = make_test_config();
for (i, obj) in objects.iter().enumerate() {
let result = deleter.delete(std::slice::from_ref(obj), &config).await.unwrap();
prop_assert_eq!(result.deleted.len(), 1);
prop_assert_eq!(&result.deleted[0].key, &format!("key/{i}"));
}
let calls = mock.delete_object_calls.lock().unwrap();
prop_assert_eq!(calls.len(), obj_count);
Ok(())
})?;
}
}
#[tokio::test]
async fn prop_concurrent_workers_process_all_objects() {
init_dummy_tracing_subscriber();
let total_objects = 2000;
let worker_count = 4;
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(total_objects);
for i in 0..total_objects {
input_sender
.send(make_s3_object(&format!("concurrent/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut handles = Vec::new();
for worker_idx in 0..worker_count {
let (boxed, _mock) = make_mock_storage_boxed(stats_sender.clone());
let (out_sender, _out_receiver) = async_channel::bounded::<S3Object>(total_objects);
let config = make_test_config();
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let has_warning = Arc::new(std::sync::atomic::AtomicBool::new(false));
let stage = Stage::new(
config,
boxed,
Some(input_receiver.clone()),
Some(out_sender),
cancellation_token,
has_warning,
);
let report = stats_report.clone();
let counter = delete_counter.clone();
let handle = tokio::spawn(async move {
let mut deleter = ObjectDeleter::new(stage, worker_idx as u16, report, counter);
deleter.delete().await.unwrap();
});
handles.push(handle);
}
for handle in handles {
handle.await.unwrap();
}
let report = &*stats_report;
assert_eq!(
report.stats_deleted_objects.load(Ordering::SeqCst),
total_objects as u64
);
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(30))]
#[test]
fn prop_partial_batch_failure_counts_successes(
total in 5usize..100,
fail_pct in 1usize..50,
) {
let fail_count = (total * fail_pct / 100).max(1).min(total - 1);
let fail_keys: HashMap<String, String> = (0..fail_count)
.map(|i| (format!("key/{i}"), "AccessDenied".to_string()))
.collect();
let objects: Vec<S3Object> = (0..total)
.map(|i| make_s3_object(&format!("key/{i}"), 100))
.collect();
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.batch_error_keys.lock().unwrap() = fail_keys.clone();
let deleter = BatchDeleter::new(boxed);
let config = make_test_config();
let result = deleter.delete(&objects, &config).await.unwrap();
let expected_successes = total - fail_count;
prop_assert_eq!(result.deleted.len(), expected_successes);
Ok(())
})?;
}
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(50))]
#[test]
fn prop_format_metadata_always_sorted(
keys in prop::collection::vec("[a-z]{1,5}", 1..10),
values in prop::collection::vec("[a-z0-9]{1,5}", 1..10),
) {
let len = keys.len().min(values.len());
let mut meta: HashMap<String, String> = HashMap::new();
for i in 0..len {
meta.insert(keys[i].clone(), values[i].clone());
}
let result = format_metadata(&meta);
let parts: Vec<&str> = result.split(',').filter(|s| !s.is_empty()).collect();
for window in parts.windows(2) {
prop_assert!(window[0] <= window[1], "Not sorted: {} > {}", window[0], window[1]);
}
}
}
#[tokio::test]
async fn object_deleter_dry_run_skips_api_calls() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::bounded(100);
let (mock_storage, mock) = make_mock_storage_boxed(stats_sender.clone());
let (input_sender, input_receiver) = async_channel::bounded(100);
let (output_sender, output_receiver) = async_channel::bounded::<S3Object>(100);
let mut config = make_test_config();
config.dry_run = true;
config.batch_size = 1000;
let stage = make_stage_with_mock(
config,
mock_storage,
Some(input_receiver),
Some(output_sender),
);
let stats_report = Arc::new(crate::types::DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter.clone());
let obj1 = make_s3_object("dry-run/file1.txt", 100);
let obj2 = make_s3_object("dry-run/file2.txt", 200);
let obj3 = make_s3_object("dry-run/file3.txt", 300);
input_sender.send(obj1).await.unwrap();
input_sender.send(obj2).await.unwrap();
input_sender.send(obj3).await.unwrap();
input_sender.close();
deleter.delete().await.unwrap();
assert_eq!(
mock.delete_object_calls.lock().unwrap().len(),
0,
"dry-run should not call delete_object"
);
assert_eq!(
mock.delete_objects_calls.lock().unwrap().len(),
0,
"dry-run should not call delete_objects"
);
let report = &*stats_report;
let snapshot = report.snapshot();
assert_eq!(
snapshot.stats_deleted_objects, 3,
"all 3 objects should be counted as deleted"
);
assert_eq!(
snapshot.stats_deleted_bytes, 600,
"total bytes should be 100+200+300=600"
);
assert_eq!(snapshot.stats_failed_objects, 0, "no failures in dry-run");
let mut forwarded = Vec::new();
while let Ok(obj) = output_receiver.try_recv() {
forwarded.push(obj.key().to_string());
}
assert_eq!(forwarded.len(), 3);
assert!(forwarded.contains(&"dry-run/file1.txt".to_string()));
assert!(forwarded.contains(&"dry-run/file2.txt".to_string()));
assert!(forwarded.contains(&"dry-run/file3.txt".to_string()));
}
#[tokio::test]
async fn object_deleter_dry_run_single_mode() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::bounded(100);
let (mock_storage, mock) = make_mock_storage_boxed(stats_sender.clone());
let (input_sender, input_receiver) = async_channel::bounded(100);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(100);
let mut config = make_test_config();
config.dry_run = true;
config.batch_size = 1;
let stage = make_stage_with_mock(
config,
mock_storage,
Some(input_receiver),
Some(output_sender),
);
let stats_report = Arc::new(crate::types::DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
let obj = make_s3_object("dry-run/single.txt", 500);
input_sender.send(obj).await.unwrap();
input_sender.close();
deleter.delete().await.unwrap();
assert_eq!(mock.delete_object_calls.lock().unwrap().len(), 0);
assert_eq!(mock.delete_objects_calls.lock().unwrap().len(), 0);
let report = &*stats_report;
let snapshot = report.snapshot();
assert_eq!(snapshot.stats_deleted_objects, 1);
assert_eq!(snapshot.stats_deleted_bytes, 500);
}
#[tokio::test]
async fn object_deleter_dry_run_versioned_objects() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::bounded(100);
let (mock_storage, mock) = make_mock_storage_boxed(stats_sender.clone());
let (input_sender, input_receiver) = async_channel::bounded(100);
let (output_sender, output_receiver) = async_channel::bounded::<S3Object>(100);
let mut config = make_test_config();
config.dry_run = true;
config.batch_size = 1000;
let stage = make_stage_with_mock(
config,
mock_storage,
Some(input_receiver),
Some(output_sender),
);
let stats_report = Arc::new(crate::types::DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
let obj = make_versioned_s3_object("versioned/file.txt", "ver-001", 1024);
input_sender.send(obj).await.unwrap();
input_sender.close();
deleter.delete().await.unwrap();
assert_eq!(mock.delete_object_calls.lock().unwrap().len(), 0);
assert_eq!(mock.delete_objects_calls.lock().unwrap().len(), 0);
let report = &*stats_report;
let snapshot = report.snapshot();
assert_eq!(snapshot.stats_deleted_objects, 1);
assert_eq!(snapshot.stats_deleted_bytes, 1024);
let forwarded = output_receiver.try_recv().unwrap();
assert_eq!(forwarded.key(), "versioned/file.txt");
assert_eq!(forwarded.version_id(), Some("ver-001"));
}
#[test]
fn test_is_retryable_error_code() {
use crate::deleter::batch::is_retryable_error_code;
assert!(is_retryable_error_code("InternalError"));
assert!(is_retryable_error_code("SlowDown"));
assert!(is_retryable_error_code("ServiceUnavailable"));
assert!(is_retryable_error_code("RequestTimeout"));
assert!(is_retryable_error_code("unknown"));
assert!(!is_retryable_error_code("AccessDenied"));
assert!(!is_retryable_error_code("NoSuchKey"));
assert!(!is_retryable_error_code("InvalidArgument"));
assert!(!is_retryable_error_code(""));
}
#[tokio::test]
async fn batch_deleter_retryable_error_falls_back_to_single() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/1".to_string(), "InternalError".to_string());
let deleter = BatchDeleter::new(boxed);
let config = make_test_config();
let objects = vec![
make_s3_object("key/0", 100),
make_s3_object("key/1", 200),
make_s3_object("key/2", 300),
];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 3);
assert_eq!(result.failed.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 1);
assert_eq!(single_calls[0].key, "key/1");
}
#[tokio::test]
async fn batch_deleter_non_retryable_skips_single_fallback() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/1".to_string(), "AccessDenied".to_string());
let deleter = BatchDeleter::new(boxed);
let config = make_test_config();
let objects = vec![make_s3_object("key/0", 100), make_s3_object("key/1", 200)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 1);
assert_eq!(result.failed.len(), 1);
assert_eq!(result.failed[0].error_code, "AccessDenied");
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 0);
}
#[tokio::test]
async fn batch_deleter_retryable_fallback_exhausted() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/0".to_string(), "SlowDown".to_string());
mock.delete_object_error_keys
.lock()
.unwrap()
.insert("key/0".to_string(), "SlowDown error".to_string());
let deleter = BatchDeleter::new(boxed);
let mut config = make_test_config();
config.force_retry_config.force_retry_count = 2;
let objects = vec![make_s3_object("key/0", 100)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 0);
assert_eq!(result.failed.len(), 1);
assert_eq!(result.failed[0].key, "key/0");
assert_eq!(result.failed[0].error_code, "SlowDown");
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 3);
}
#[tokio::test]
async fn batch_deleter_mixed_retryable_non_retryable() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
{
let mut errors = mock.batch_error_keys.lock().unwrap();
errors.insert("key/0".to_string(), "InternalError".to_string());
errors.insert("key/2".to_string(), "NoSuchKey".to_string());
}
let deleter = BatchDeleter::new(boxed);
let config = make_test_config();
let objects = vec![
make_s3_object("key/0", 100),
make_s3_object("key/1", 200),
make_s3_object("key/2", 300),
];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 2); assert_eq!(result.failed.len(), 1);
assert_eq!(result.failed[0].key, "key/2");
assert_eq!(result.failed[0].error_code, "NoSuchKey");
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 1);
assert_eq!(single_calls[0].key, "key/0");
}
#[tokio::test]
async fn batch_deleter_retryable_fallback_with_if_match() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/0".to_string(), "ServiceUnavailable".to_string());
let deleter = BatchDeleter::new(boxed);
let mut config = make_test_config();
config.if_match = true;
let objects = vec![S3Object::NotVersioning(
Object::builder()
.key("key/0")
.size(100)
.e_tag("\"abc123\"")
.build(),
)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 1);
assert_eq!(result.failed.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 1);
assert_eq!(single_calls[0].key, "key/0");
assert_eq!(single_calls[0].if_match, Some("\"abc123\"".to_string()));
}
#[tokio::test]
async fn batch_deleter_retryable_fallback_passes_version_id() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/0".to_string(), "RequestTimeout".to_string());
let deleter = BatchDeleter::new(boxed);
let config = make_test_config();
let objects = vec![make_versioned_s3_object("key/0", "ver-abc", 100)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 1);
assert_eq!(result.failed.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 1);
assert_eq!(single_calls[0].key, "key/0");
assert_eq!(single_calls[0].version_id, Some("ver-abc".to_string()));
}
#[tokio::test]
async fn batch_deleter_retryable_fallback_succeeds_on_second_attempt() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let mock = MockStorage::new(stats_sender.clone());
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/0".to_string(), "InternalError".to_string());
let fail_counter = Arc::new(AtomicU64::new(0));
let fail_counter_clone = fail_counter.clone();
#[derive(Clone)]
struct FailOnceMock {
inner: MockStorage,
fail_counter: Arc<AtomicU64>,
}
#[async_trait]
impl StorageTrait for FailOnceMock {
fn is_express_onezone_storage(&self) -> bool {
self.inner.is_express_onezone_storage()
}
async fn list_objects(&self, sender: &Sender<S3Object>, max_keys: i32) -> Result<()> {
self.inner.list_objects(sender, max_keys).await
}
async fn list_object_versions(
&self,
sender: &Sender<S3Object>,
max_keys: i32,
) -> Result<()> {
self.inner.list_object_versions(sender, max_keys).await
}
async fn head_object(
&self,
key: &str,
version_id: Option<String>,
) -> Result<HeadObjectOutput> {
self.inner.head_object(key, version_id).await
}
async fn get_object_tagging(
&self,
key: &str,
version_id: Option<String>,
) -> Result<GetObjectTaggingOutput> {
self.inner.get_object_tagging(key, version_id).await
}
async fn delete_object(
&self,
key: &str,
version_id: Option<String>,
if_match: Option<String>,
) -> Result<DeleteObjectOutput> {
self.inner
.delete_object_calls
.lock()
.unwrap()
.push(DeleteObjectCall {
key: key.to_string(),
version_id: version_id.clone(),
if_match: if_match.clone(),
});
let count = self.fail_counter.fetch_add(1, Ordering::SeqCst);
if count == 0 {
Err(anyhow::anyhow!("transient error"))
} else {
Ok(DeleteObjectOutput::builder().build())
}
}
async fn delete_objects(
&self,
objects: Vec<ObjectIdentifier>,
) -> Result<DeleteObjectsOutput> {
self.inner.delete_objects(objects).await
}
async fn is_versioning_enabled(&self) -> Result<bool> {
self.inner.is_versioning_enabled().await
}
fn get_client(&self) -> Option<Arc<aws_sdk_s3::Client>> {
None
}
fn get_stats_sender(&self) -> Sender<DeletionStatistics> {
self.inner.get_stats_sender()
}
async fn send_stats(&self, stats: DeletionStatistics) {
self.inner.send_stats(stats).await;
}
fn set_warning(&self) {}
}
let fail_once_mock = FailOnceMock {
inner: mock.clone(),
fail_counter: fail_counter_clone,
};
let boxed: Box<dyn StorageTrait + Send + Sync> = Box::new(fail_once_mock);
let deleter = BatchDeleter::new(boxed);
let mut config = make_test_config();
config.force_retry_config.force_retry_count = 2;
let objects = vec![make_s3_object("key/0", 100)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 1);
assert_eq!(result.failed.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 2);
assert_eq!(fail_counter.load(Ordering::SeqCst), 2);
}
fn make_stage_with_observables(
config: Config,
mock_storage: Box<dyn StorageTrait + Send + Sync>,
receiver: Option<async_channel::Receiver<S3Object>>,
sender: Option<Sender<S3Object>>,
) -> (
Stage,
PipelineCancellationToken,
Arc<std::sync::atomic::AtomicBool>,
) {
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let has_warning = Arc::new(std::sync::atomic::AtomicBool::new(false));
let stage = Stage::new(
config,
mock_storage,
receiver,
sender,
cancellation_token.clone(),
has_warning.clone(),
);
(stage, cancellation_token, has_warning)
}
#[tokio::test]
async fn warn_as_error_false_failure_sets_warning_but_continues() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/1".to_string(), "AccessDenied".to_string());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.warn_as_error = false; config.batch_size = 1000;
let (stage, cancellation_token, has_warning) =
make_stage_with_observables(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..3 {
input_sender
.send(make_s3_object(&format!("key/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
assert!(
has_warning.load(Ordering::SeqCst),
"Warning flag should be set when a batch deletion partially fails"
);
assert!(
!cancellation_token.is_cancelled(),
"Pipeline should continue when warn_as_error=false"
);
let snapshot = stats_report.snapshot();
assert_eq!(snapshot.stats_deleted_objects, 2);
assert_eq!(snapshot.stats_failed_objects, 1);
let mut forwarded = Vec::new();
while let Ok(obj) = output_receiver.try_recv() {
forwarded.push(obj.key().to_string());
}
assert_eq!(
forwarded.len(),
3,
"All objects should be forwarded to next stage even with failures"
);
}
#[tokio::test]
async fn warn_as_error_true_failure_cancels_pipeline() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/1".to_string(), "AccessDenied".to_string());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.warn_as_error = true;
config.batch_size = 1000;
let (stage, cancellation_token, has_warning) =
make_stage_with_observables(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..3 {
input_sender
.send(make_s3_object(&format!("key/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
assert!(
has_warning.load(Ordering::SeqCst),
"Warning flag should be set on batch failure"
);
assert!(
cancellation_token.is_cancelled(),
"Pipeline must be cancelled when warn_as_error=true and failures occur"
);
let snapshot = stats_report.snapshot();
assert_eq!(snapshot.stats_deleted_objects, 2);
assert_eq!(snapshot.stats_failed_objects, 1);
let forwarded_count = output_receiver.len();
assert_eq!(
forwarded_count, 0,
"No objects should be forwarded when warn_as_error cancels the pipeline"
);
}
#[tokio::test]
async fn warn_as_error_true_no_failures_completes_normally() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.warn_as_error = true;
config.batch_size = 1000;
let (stage, cancellation_token, has_warning) =
make_stage_with_observables(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..3 {
input_sender
.send(make_s3_object(&format!("key/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
assert!(
!has_warning.load(Ordering::SeqCst),
"Warning flag should not be set when there are no failures"
);
assert!(
!cancellation_token.is_cancelled(),
"Pipeline should not be cancelled when there are no failures"
);
let snapshot = stats_report.snapshot();
assert_eq!(snapshot.stats_deleted_objects, 3);
assert_eq!(snapshot.stats_failed_objects, 0);
let mut forwarded = Vec::new();
while let Ok(obj) = output_receiver.try_recv() {
forwarded.push(obj.key().to_string());
}
assert_eq!(forwarded.len(), 3);
}
#[tokio::test]
async fn warn_as_error_false_no_failures_baseline() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, output_receiver) = async_channel::bounded::<S3Object>(10);
let config = make_test_config();
let (stage, cancellation_token, has_warning) =
make_stage_with_observables(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..3 {
input_sender
.send(make_s3_object(&format!("key/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
assert!(!has_warning.load(Ordering::SeqCst));
assert!(!cancellation_token.is_cancelled());
let snapshot = stats_report.snapshot();
assert_eq!(snapshot.stats_deleted_objects, 3);
assert_eq!(snapshot.stats_failed_objects, 0);
let mut forwarded = Vec::new();
while let Ok(obj) = output_receiver.try_recv() {
forwarded.push(obj.key().to_string());
}
assert_eq!(forwarded.len(), 3);
}
#[tokio::test]
async fn warn_as_error_true_cancels_on_first_failing_batch() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/0".to_string(), "AccessDenied".to_string());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(20);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(20);
let mut config = make_test_config();
config.warn_as_error = true;
config.batch_size = 3;
let (stage, cancellation_token, has_warning) =
make_stage_with_observables(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..9 {
input_sender
.send(make_s3_object(&format!("key/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
assert!(
cancellation_token.is_cancelled(),
"Pipeline must cancel on first batch with failures when warn_as_error=true"
);
assert!(has_warning.load(Ordering::SeqCst));
let snapshot = stats_report.snapshot();
assert_eq!(
snapshot.stats_deleted_objects, 2,
"Only objects from the first batch should be deleted"
);
assert_eq!(
snapshot.stats_failed_objects, 1,
"Only the failure from the first batch should be counted"
);
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(
batch_calls.len(),
1,
"Only one batch should be sent to S3 before cancellation"
);
assert_eq!(batch_calls[0].identifiers.len(), 3);
}
#[tokio::test]
async fn warn_as_error_true_single_deleter_failure_cancels() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/0".to_string(), "AccessDenied".to_string());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.warn_as_error = true;
config.batch_size = 2;
let (stage, cancellation_token, has_warning) =
make_stage_with_observables(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..4 {
input_sender
.send(make_s3_object(&format!("key/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
assert!(cancellation_token.is_cancelled());
assert!(has_warning.load(Ordering::SeqCst));
let snapshot = stats_report.snapshot();
assert_eq!(snapshot.stats_deleted_objects, 1); assert_eq!(snapshot.stats_failed_objects, 1); }
#[tokio::test]
async fn delete_error_stats_emitted_for_each_failed_key() {
init_dummy_tracing_subscriber();
for warn_as_error in [false, true] {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender.clone());
{
let mut errors = mock.batch_error_keys.lock().unwrap();
errors.insert("key/1".to_string(), "AccessDenied".to_string());
errors.insert("key/3".to_string(), "NoSuchKey".to_string());
}
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.warn_as_error = warn_as_error;
config.batch_size = 1000;
let (stage, _cancellation_token, _has_warning) =
make_stage_with_observables(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..5 {
input_sender
.send(make_s3_object(&format!("key/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
stats_sender.close();
let mut error_keys = Vec::new();
while let Ok(stat) = stats_receiver.try_recv() {
if let DeletionStatistics::DeleteError { key } = stat {
error_keys.push(key);
}
}
assert!(
error_keys.contains(&"key/1".to_string()),
"warn_as_error={warn_as_error}: DeleteError should be emitted for key/1"
);
assert!(
error_keys.contains(&"key/3".to_string()),
"warn_as_error={warn_as_error}: DeleteError should be emitted for key/3"
);
}
}
#[tokio::test]
async fn warn_as_error_true_retryable_failures_recovered_no_cancellation() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/1".to_string(), "InternalError".to_string());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.warn_as_error = true;
config.batch_size = 1000;
let (stage, cancellation_token, has_warning) =
make_stage_with_observables(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..3 {
input_sender
.send(make_s3_object(&format!("key/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
assert!(
!cancellation_token.is_cancelled(),
"Pipeline should not cancel when all retryable failures are recovered"
);
assert!(
!has_warning.load(Ordering::SeqCst),
"Warning should not be set when all failures are recovered via fallback"
);
let snapshot = stats_report.snapshot();
assert_eq!(snapshot.stats_deleted_objects, 3);
assert_eq!(snapshot.stats_failed_objects, 0);
let mut forwarded = Vec::new();
while let Ok(obj) = output_receiver.try_recv() {
forwarded.push(obj.key().to_string());
}
assert_eq!(forwarded.len(), 3);
}
#[tokio::test]
async fn warn_as_error_true_mixed_failures_non_retryable_triggers_cancel() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
{
let mut errors = mock.batch_error_keys.lock().unwrap();
errors.insert("key/0".to_string(), "InternalError".to_string());
errors.insert("key/2".to_string(), "NoSuchKey".to_string());
}
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.warn_as_error = true;
config.batch_size = 1000;
let (stage, cancellation_token, has_warning) =
make_stage_with_observables(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..3 {
input_sender
.send(make_s3_object(&format!("key/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
assert!(
cancellation_token.is_cancelled(),
"Pipeline must cancel when non-retryable failure remains with warn_as_error=true"
);
assert!(has_warning.load(Ordering::SeqCst));
let snapshot = stats_report.snapshot();
assert_eq!(snapshot.stats_deleted_objects, 2);
assert_eq!(snapshot.stats_failed_objects, 1);
}
#[tokio::test]
async fn object_deleter_tag_include_filter_multiple_tags() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("team").value("backend").build().unwrap(),
Tag::builder()
.key("env")
.value("production")
.build()
.unwrap(),
Tag::builder().key("retain").value("false").build().unwrap(),
]);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_tag_regex =
Some(Regex::new("env=production&retain=false&team=backend").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("multi-tagged.txt", 256))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "multi-tagged.txt");
let head_calls = mock.head_object_calls.lock().unwrap();
assert_eq!(head_calls.len(), 0);
}
#[tokio::test]
async fn object_deleter_metadata_include_filter_multiple_entries() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let mut meta = HashMap::new();
meta.insert("version".to_string(), "v2".to_string());
meta.insert("env".to_string(), "staging".to_string());
meta.insert("team".to_string(), "data".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_metadata_regex =
Some(Regex::new("env=staging,team=data,version=v2").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("metadata-multi.json", 500))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "metadata-multi.json");
}
#[tokio::test]
async fn object_deleter_tag_exclude_filter_alternation_pattern() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("team").value("backend").build().unwrap(),
Tag::builder()
.key("env")
.value("production")
.build()
.unwrap(),
Tag::builder().key("retain").value("true").build().unwrap(),
]);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.exclude_tag_regex =
Some(Regex::new("env=(production|staging)&retain=true").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("excluded-by-tag.txt", 256))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(
batch_calls.len(),
0,
"Excluded object should not be deleted"
);
}
#[tokio::test]
async fn object_deleter_metadata_exclude_filter_alternation_pattern() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let mut meta = HashMap::new();
meta.insert("version".to_string(), "v2".to_string());
meta.insert("env".to_string(), "staging".to_string());
meta.insert("team".to_string(), "backend".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.exclude_metadata_regex =
Some(Regex::new("env=(staging|production),team=backend").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("excluded-by-meta.json", 500))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(
batch_calls.len(),
0,
"Excluded object should not be deleted"
);
}
#[tokio::test]
async fn batch_deleter_retryable_error_falls_back_to_single_delete() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/1".to_string(), "InternalError".to_string());
let deleter = BatchDeleter::new(boxed);
let mut config = make_test_config();
config.force_retry_config.force_retry_count = 1;
let objects = vec![
make_s3_object("key/0", 100),
make_s3_object("key/1", 200),
make_s3_object("key/2", 300),
];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.deleted.len(), 3);
assert_eq!(result.failed.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
assert!(
single_calls.iter().any(|c| c.key == "key/1"),
"single-delete fallback must be attempted for key/1"
);
}
#[tokio::test]
async fn object_deleter_all_include_filters_pass_through() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("application/json".to_string());
let mut meta = HashMap::new();
meta.insert("env".to_string(), "staging".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("team").value("data").build().unwrap(),
]);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("application/json").unwrap());
config.filter_config.include_metadata_regex = Some(Regex::new("env=staging").unwrap());
config.filter_config.include_tag_regex = Some(Regex::new("team=data").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("match-all.json", 1024))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "match-all.json");
let head_calls = mock.head_object_calls.lock().unwrap();
assert_eq!(head_calls.len(), 1);
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(30))]
#[test]
fn prop_content_type_include_filter(
content_type in "(text|image|audio)/[a-z]{3,8}",
num_objects in 1usize..10,
) {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some(content_type.clone());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(num_objects + 1);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(num_objects + 1);
let mut config = make_test_config();
config.filter_config.include_content_type_regex =
Some(Regex::new("application/octet-stream").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..num_objects {
input_sender
.send(make_s3_object(&format!("ct-obj/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
prop_assert_eq!(batch_calls.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
prop_assert_eq!(single_calls.len(), 0);
let head_calls = mock.head_object_calls.lock().unwrap();
prop_assert_eq!(head_calls.len(), num_objects);
Ok(())
})?;
}
#[test]
fn prop_metadata_include_filter(
meta_key in "[a-z]{3,8}",
meta_value in "[a-z0-9]{3,8}",
num_objects in 1usize..10,
) {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let mut meta = HashMap::new();
meta.insert(meta_key.clone(), meta_value.clone());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(num_objects + 1);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(num_objects + 1);
let mut config = make_test_config();
config.filter_config.include_metadata_regex =
Some(Regex::new("ZZZZZ_NO_MATCH_ZZZZZ").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..num_objects {
input_sender
.send(make_s3_object(&format!("meta-obj/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
prop_assert_eq!(batch_calls.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
prop_assert_eq!(single_calls.len(), 0);
Ok(())
})?;
}
#[test]
fn prop_tag_include_filter(
tag_key in "[a-z]{3,8}",
tag_value in "[a-z0-9]{3,8}",
num_objects in 1usize..10,
) {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let tag = Tag::builder()
.key(tag_key.clone())
.value(tag_value.clone())
.build()
.unwrap();
*mock.tagging_response_tags.lock().unwrap() = Some(vec![tag]);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(num_objects + 1);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(num_objects + 1);
let mut config = make_test_config();
config.filter_config.include_tag_regex =
Some(Regex::new("ZZZZZ_NO_MATCH_ZZZZZ").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
for i in 0..num_objects {
input_sender
.send(make_s3_object(&format!("tag-obj/{i}"), 100))
.await
.unwrap();
}
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
prop_assert_eq!(batch_calls.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
prop_assert_eq!(single_calls.len(), 0);
Ok(())
})?;
}
}
#[tokio::test]
async fn object_deleter_include_content_type_none_filters_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("application/json").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("no-ct.bin", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 0);
}
#[tokio::test]
async fn object_deleter_exclude_content_type_none_passes_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.exclude_content_type_regex = Some(Regex::new("text/.*").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("no-ct.bin", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "no-ct.bin");
}
#[tokio::test]
async fn object_deleter_include_metadata_none_filters_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_metadata_regex = Some(Regex::new("env=prod").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("no-meta.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
}
#[tokio::test]
async fn object_deleter_exclude_metadata_none_passes_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.exclude_metadata_regex = Some(Regex::new("env=dev").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("no-meta.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "no-meta.txt");
}
#[tokio::test]
async fn object_deleter_include_metadata_empty_filters_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_metadata.lock().unwrap() = Some(HashMap::new());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_metadata_regex = Some(Regex::new("env=prod").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("empty-meta.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
}
#[tokio::test]
async fn object_deleter_exclude_metadata_empty_passes_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_metadata.lock().unwrap() = Some(HashMap::new());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.exclude_metadata_regex = Some(Regex::new("env=dev").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("empty-meta.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
}
#[tokio::test]
async fn object_deleter_warn_as_error_cancels_on_failure() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.batch_error_keys
.lock()
.unwrap()
.insert("key/fail".to_string(), "AccessDenied".to_string());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.warn_as_error = true;
config.batch_size = 10; let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let stage = make_stage_with_mock_and_token(
config,
boxed,
Some(input_receiver),
Some(output_sender),
cancellation_token.clone(),
);
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("key/ok", 100))
.await
.unwrap();
input_sender
.send(make_s3_object("key/fail", 200))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
assert!(
cancellation_token.is_cancelled(),
"warn_as_error should cancel the pipeline on batch failure"
);
}
#[tokio::test]
async fn object_deleter_emits_filter_skip_event_with_callback() {
use crate::callback::event_manager::EventManager;
use crate::types::event_callback::{EventCallback, EventData, EventType};
struct CollectingCallback {
events: Arc<tokio::sync::Mutex<Vec<EventData>>>,
}
#[async_trait]
impl EventCallback for CollectingCallback {
async fn on_event(&mut self, event_data: EventData) {
self.events.lock().await.push(event_data);
}
}
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let events = Arc::new(tokio::sync::Mutex::new(Vec::new()));
let callback = CollectingCallback {
events: events.clone(),
};
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("application/json").unwrap());
let mut event_manager = EventManager::new();
event_manager.register_callback(EventType::ALL_EVENTS, callback, false);
config.event_manager = event_manager;
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, delete_counter);
input_sender
.send(make_s3_object("filtered.txt", 256))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let collected = events.lock().await;
assert_eq!(collected.len(), 1);
assert_eq!(collected[0].event_type, EventType::DELETE_FILTERED);
assert_eq!(collected[0].key.as_deref(), Some("filtered.txt"));
assert!(
collected[0]
.message
.as_ref()
.unwrap()
.contains("include_content_type_regex_filter")
);
}
#[tokio::test]
async fn object_deleter_include_tag_none_response_filters_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![]);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("NONEXISTENT_TAG_PATTERN").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, delete_counter);
input_sender
.send(make_s3_object("no-tags.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
}
#[tokio::test]
async fn object_deleter_exclude_tag_empty_tags_passes_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![]);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.exclude_tag_regex = Some(Regex::new("retain=true").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, delete_counter);
input_sender
.send(make_s3_object("no-tags.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
}
#[tokio::test]
async fn object_deleter_head_object_not_found_skips_object() {
use aws_sdk_s3::operation::head_object::HeadObjectError;
use aws_smithy_runtime_api::http::StatusCode;
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let service_err = HeadObjectError::NotFound(
aws_sdk_s3::types::error::NotFound::builder()
.message("Not Found")
.build(),
);
let raw_response = aws_smithy_runtime_api::http::Response::new(
StatusCode::try_from(404).unwrap(),
SdkBody::from(""),
);
let sdk_error: SdkError<HeadObjectError, Response<SdkBody>> =
SdkError::service_error(service_err, raw_response);
*mock.head_object_error.lock().unwrap() = Some(sdk_error.into());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("application/json").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, delete_counter);
input_sender
.send(make_s3_object("missing.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
let single_calls = mock.delete_object_calls.lock().unwrap();
assert_eq!(single_calls.len(), 0);
let mut found_skip = false;
while let Ok(stat) = stats_receiver.try_recv() {
if matches!(stat, DeletionStatistics::DeleteSkip { .. }) {
found_skip = true;
}
}
assert!(
found_skip,
"DeleteSkip stat should be emitted for not-found"
);
}
#[tokio::test]
async fn object_deleter_head_object_error_cancels_pipeline() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_error.lock().unwrap() = Some(anyhow::anyhow!("connection refused"));
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("application/json").unwrap());
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let stage = make_stage_with_mock_and_token(
config,
boxed,
Some(input_receiver),
Some(output_sender),
cancellation_token.clone(),
);
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, delete_counter);
input_sender
.send(make_s3_object("fail.txt", 100))
.await
.unwrap();
drop(input_sender);
let result = deleter.delete().await;
assert!(
result.is_err(),
"head_object failure should propagate as error"
);
assert!(
cancellation_token.is_cancelled(),
"pipeline should be cancelled on head_object error"
);
}
#[tokio::test]
async fn object_deleter_tagging_not_found_skips_object() {
use aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingError;
use aws_smithy_runtime_api::http::StatusCode;
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let service_err = GetObjectTaggingError::generic(
aws_smithy_types::error::ErrorMetadata::builder()
.code("NoSuchKey")
.message("The specified key does not exist.")
.build(),
);
let raw_response = aws_smithy_runtime_api::http::Response::new(
StatusCode::try_from(404).unwrap(),
SdkBody::from(""),
);
let sdk_error: SdkError<GetObjectTaggingError, Response<SdkBody>> =
SdkError::service_error(service_err, raw_response);
*mock.tagging_error.lock().unwrap() = Some(sdk_error.into());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("env=prod").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, delete_counter);
input_sender
.send(make_s3_object("vanished.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
let mut found_skip = false;
while let Ok(stat) = stats_receiver.try_recv() {
if matches!(stat, DeletionStatistics::DeleteSkip { .. }) {
found_skip = true;
}
}
assert!(
found_skip,
"DeleteSkip stat should be emitted for tagging not-found"
);
}
#[tokio::test]
async fn object_deleter_tagging_error_cancels_pipeline() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_error.lock().unwrap() = Some(anyhow::anyhow!("service unavailable"));
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("env=prod").unwrap());
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let stage = make_stage_with_mock_and_token(
config,
boxed,
Some(input_receiver),
Some(output_sender),
cancellation_token.clone(),
);
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, delete_counter);
input_sender
.send(make_s3_object("fail-tag.txt", 100))
.await
.unwrap();
drop(input_sender);
let result = deleter.delete().await;
assert!(result.is_err(), "tagging failure should propagate as error");
assert!(
cancellation_token.is_cancelled(),
"pipeline should be cancelled on get_object_tagging error"
);
}
#[test]
fn is_not_found_error_head_object_not_found() {
use aws_sdk_s3::operation::head_object::HeadObjectError;
use aws_smithy_runtime_api::http::StatusCode;
let service_err = HeadObjectError::NotFound(
aws_sdk_s3::types::error::NotFound::builder()
.message("Not Found")
.build(),
);
let raw_response = aws_smithy_runtime_api::http::Response::new(
StatusCode::try_from(404).unwrap(),
SdkBody::from(""),
);
let sdk_error: SdkError<HeadObjectError, Response<SdkBody>> =
SdkError::service_error(service_err, raw_response);
let anyhow_err: anyhow::Error = sdk_error.into();
assert!(is_not_found_error(&anyhow_err));
}
#[test]
fn is_not_found_error_head_object_other_error() {
let err = anyhow::anyhow!("connection refused");
assert!(!is_not_found_error(&err));
}
#[test]
fn is_not_found_error_tagging_no_such_key() {
use aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingError;
use aws_smithy_runtime_api::http::StatusCode;
let service_err = GetObjectTaggingError::generic(
aws_smithy_types::error::ErrorMetadata::builder()
.code("NoSuchKey")
.message("The specified key does not exist.")
.build(),
);
let raw_response = aws_smithy_runtime_api::http::Response::new(
StatusCode::try_from(404).unwrap(),
SdkBody::from(""),
);
let sdk_error: SdkError<GetObjectTaggingError, Response<SdkBody>> =
SdkError::service_error(service_err, raw_response);
let anyhow_err: anyhow::Error = sdk_error.into();
assert!(is_not_found_error(&anyhow_err));
}
#[test]
fn is_not_found_error_tagging_404_status() {
use aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingError;
use aws_smithy_runtime_api::http::StatusCode;
let service_err = GetObjectTaggingError::generic(
aws_smithy_types::error::ErrorMetadata::builder()
.code("NotFound")
.message("Not Found")
.build(),
);
let raw_response = aws_smithy_runtime_api::http::Response::new(
StatusCode::try_from(404).unwrap(),
SdkBody::from(""),
);
let sdk_error: SdkError<GetObjectTaggingError, Response<SdkBody>> =
SdkError::service_error(service_err, raw_response);
let anyhow_err: anyhow::Error = sdk_error.into();
assert!(is_not_found_error(&anyhow_err));
}
#[test]
fn is_not_found_error_tagging_404_status_unknown_code() {
use aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingError;
use aws_smithy_runtime_api::http::StatusCode;
let service_err = GetObjectTaggingError::generic(
aws_smithy_types::error::ErrorMetadata::builder()
.code("SomeUnknownCode")
.message("something went wrong")
.build(),
);
let raw_response = aws_smithy_runtime_api::http::Response::new(
StatusCode::try_from(404).unwrap(),
SdkBody::from(""),
);
let sdk_error: SdkError<GetObjectTaggingError, Response<SdkBody>> =
SdkError::service_error(service_err, raw_response);
let anyhow_err: anyhow::Error = sdk_error.into();
assert!(is_not_found_error(&anyhow_err));
}
#[test]
fn is_not_found_error_tagging_non_404_status_unknown_code() {
use aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingError;
use aws_smithy_runtime_api::http::StatusCode;
let service_err = GetObjectTaggingError::generic(
aws_smithy_types::error::ErrorMetadata::builder()
.code("SomeUnknownCode")
.message("something went wrong")
.build(),
);
let raw_response = aws_smithy_runtime_api::http::Response::new(
StatusCode::try_from(500).unwrap(),
SdkBody::from(""),
);
let sdk_error: SdkError<GetObjectTaggingError, Response<SdkBody>> =
SdkError::service_error(service_err, raw_response);
let anyhow_err: anyhow::Error = sdk_error.into();
assert!(!is_not_found_error(&anyhow_err));
}
#[tokio::test]
async fn object_deleter_emits_filter_skip_event_for_tag_filter() {
use crate::callback::event_manager::EventManager;
use crate::types::event_callback::{EventCallback, EventData, EventType};
struct CollectingCallback {
events: Arc<tokio::sync::Mutex<Vec<EventData>>>,
}
#[async_trait]
impl EventCallback for CollectingCallback {
async fn on_event(&mut self, event_data: EventData) {
self.events.lock().await.push(event_data);
}
}
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("env").value("staging").build().unwrap(),
]);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let events = Arc::new(tokio::sync::Mutex::new(Vec::new()));
let callback = CollectingCallback {
events: events.clone(),
};
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("env=production").unwrap());
let mut event_manager = EventManager::new();
event_manager.register_callback(EventType::ALL_EVENTS, callback, false);
config.event_manager = event_manager;
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, delete_counter);
input_sender
.send(make_s3_object("tagged-wrong.txt", 512))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let collected = events.lock().await;
assert_eq!(collected.len(), 1);
assert_eq!(collected[0].event_type, EventType::DELETE_FILTERED);
assert_eq!(collected[0].key.as_deref(), Some("tagged-wrong.txt"));
assert_eq!(collected[0].size, Some(512));
assert!(
collected[0]
.message
.as_ref()
.unwrap()
.contains("include_tag_regex_filter")
);
}
#[tokio::test]
async fn object_deleter_content_type_include_passes_matching() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("application/json".to_string());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("application/json").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("data.json", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "data.json");
}
#[tokio::test]
async fn object_deleter_content_type_exclude_passes_non_matching() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("application/json".to_string());
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.exclude_content_type_regex = Some(Regex::new("text/.*").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("data.json", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "data.json");
}
#[tokio::test]
async fn object_deleter_metadata_exclude_filter_skips_matching() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let mut meta = HashMap::new();
meta.insert("env".to_string(), "staging".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.filter_config.exclude_metadata_regex = Some(Regex::new("env=staging").unwrap());
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, delete_counter);
input_sender
.send(make_s3_object("staging-data.txt", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 0);
}
#[test]
fn mock_storage_is_express_onezone_returns_false() {
let (stats_sender, _) = async_channel::unbounded();
let mock = MockStorage::new(stats_sender);
assert!(!mock.is_express_onezone_storage());
}
#[tokio::test]
async fn mock_storage_list_objects_returns_ok() {
let (stats_sender, _) = async_channel::unbounded();
let mock = MockStorage::new(stats_sender);
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_objects(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn mock_storage_list_object_versions_returns_ok() {
let (stats_sender, _) = async_channel::unbounded();
let mock = MockStorage::new(stats_sender);
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_object_versions(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn mock_storage_is_versioning_enabled_returns_false() {
let (stats_sender, _) = async_channel::unbounded();
let mock = MockStorage::new(stats_sender);
assert!(!mock.is_versioning_enabled().await.unwrap());
}
#[test]
fn mock_storage_get_client_returns_none() {
let (stats_sender, _) = async_channel::unbounded();
let mock = MockStorage::new(stats_sender);
assert!(mock.get_client().is_none());
}
#[tokio::test]
async fn mock_storage_get_stats_sender_works() {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let mock = MockStorage::new(stats_sender);
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 mock_storage_send_stats_delivers_stat() {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let mock = MockStorage::new(stats_sender);
mock.send_stats(DeletionStatistics::DeleteComplete {
key: "k".to_string(),
})
.await;
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(
received,
DeletionStatistics::DeleteComplete { .. }
));
}
#[test]
fn mock_storage_set_warning_is_noop() {
let (stats_sender, _) = async_channel::unbounded();
let mock = MockStorage::new(stats_sender);
mock.set_warning(); }
fn make_object_deleter_for_unit_test(
config: Config,
mock_storage: Box<dyn StorageTrait + Send + Sync>,
delete_counter: Arc<AtomicU64>,
) -> ObjectDeleter {
let stage = make_stage_with_mock(config, mock_storage, None, None);
let stats_report = Arc::new(DeletionStatsReport::new());
ObjectDeleter::new(stage, 0, stats_report, delete_counter)
}
fn make_object_deleter_with_token(
config: Config,
mock_storage: Box<dyn StorageTrait + Send + Sync>,
delete_counter: Arc<AtomicU64>,
cancellation_token: PipelineCancellationToken,
) -> ObjectDeleter {
let stage =
make_stage_with_mock_and_token(config, mock_storage, None, None, cancellation_token);
let stats_report = Arc::new(DeletionStatsReport::new());
ObjectDeleter::new(stage, 0, stats_report, delete_counter)
}
#[tokio::test]
async fn check_max_delete_no_limit_returns_false() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config(); let counter = Arc::new(AtomicU64::new(0));
let mut deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.check_max_delete(&obj).await.unwrap());
}
#[tokio::test]
async fn check_max_delete_under_limit_returns_false() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.max_delete = Some(5);
let counter = Arc::new(AtomicU64::new(0));
let mut deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.check_max_delete(&obj).await.unwrap());
}
#[tokio::test]
async fn check_max_delete_at_exact_limit_returns_false() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.max_delete = Some(3);
let counter = Arc::new(AtomicU64::new(2));
let mut deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.check_max_delete(&obj).await.unwrap());
}
#[tokio::test]
async fn check_max_delete_exceeds_limit_returns_true_and_cancels() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.max_delete = Some(2);
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let counter = Arc::new(AtomicU64::new(2));
let mut deleter =
make_object_deleter_with_token(config, boxed, counter, cancellation_token.clone());
let obj = make_s3_object("key/1", 100);
assert!(deleter.check_max_delete(&obj).await.unwrap());
assert!(cancellation_token.is_cancelled());
}
#[tokio::test]
async fn check_max_delete_increments_counter() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config(); let counter = Arc::new(AtomicU64::new(5));
let mut deleter = make_object_deleter_for_unit_test(config, boxed, counter.clone());
let obj = make_s3_object("key/1", 100);
deleter.check_max_delete(&obj).await.unwrap();
assert_eq!(counter.load(Ordering::SeqCst), 6);
}
#[tokio::test]
async fn check_max_delete_sends_stats_on_exceed() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.max_delete = Some(1);
let counter = Arc::new(AtomicU64::new(1)); let mut deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/exceeded", 100);
deleter.check_max_delete(&obj).await.unwrap();
let stat = stats_receiver.recv().await.unwrap();
assert!(matches!(stat, DeletionStatistics::DeleteError { key } if key == "key/exceeded"));
}
#[tokio::test]
async fn check_max_delete_flushes_buffer_before_cancel() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.max_delete = Some(2);
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let counter = Arc::new(AtomicU64::new(2));
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let stage = make_stage_with_mock_and_token(
config,
boxed,
None,
Some(output_sender),
cancellation_token.clone(),
);
let stats_report = Arc::new(DeletionStatsReport::new());
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), counter);
deleter.buffer.push(make_s3_object("key/allowed", 500));
let obj = make_s3_object("key/exceeded", 100);
deleter.check_max_delete(&obj).await.unwrap();
assert_eq!(deleter.buffer.len(), 0);
let batch_calls = mock.delete_objects_calls.lock().unwrap();
assert_eq!(batch_calls.len(), 1);
assert_eq!(batch_calls[0].identifiers.len(), 1);
assert_eq!(batch_calls[0].identifiers[0].key(), "key/allowed");
assert!(cancellation_token.is_cancelled());
}
#[tokio::test]
async fn apply_head_object_filters_no_filters_returns_false() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config(); let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_include_content_type_match_passes() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("text/plain".to_string());
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("text/.*").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_include_content_type_no_match_filters() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("image/png".to_string());
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("text/.*").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_exclude_content_type_match_filters() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("text/html".to_string());
let mut config = make_test_config();
config.filter_config.exclude_content_type_regex = Some(Regex::new("text/html").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_exclude_content_type_no_match_passes() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("application/json".to_string());
let mut config = make_test_config();
config.filter_config.exclude_content_type_regex = Some(Regex::new("text/html").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_include_metadata_match_passes() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let mut meta = HashMap::new();
meta.insert("env".to_string(), "production".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let mut config = make_test_config();
config.filter_config.include_metadata_regex = Some(Regex::new("env=production").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_include_metadata_no_match_filters() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let mut meta = HashMap::new();
meta.insert("env".to_string(), "staging".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let mut config = make_test_config();
config.filter_config.include_metadata_regex = Some(Regex::new("env=production").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_exclude_metadata_match_filters() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let mut meta = HashMap::new();
meta.insert("env".to_string(), "staging".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let mut config = make_test_config();
config.filter_config.exclude_metadata_regex = Some(Regex::new("env=staging").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_exclude_metadata_no_match_passes() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let mut meta = HashMap::new();
meta.insert("env".to_string(), "production".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let mut config = make_test_config();
config.filter_config.exclude_metadata_regex = Some(Regex::new("env=staging").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_no_content_type_include_filter_rejects() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("text/.*").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_no_content_type_exclude_filter_passes() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.filter_config.exclude_content_type_regex = Some(Regex::new("text/.*").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_no_metadata_include_filter_rejects() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.filter_config.include_metadata_regex = Some(Regex::new("env=prod").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_no_metadata_exclude_filter_passes() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.filter_config.exclude_metadata_regex = Some(Regex::new("env=prod").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_head_object_error_cancels_pipeline() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_error.lock().unwrap() = Some(anyhow::anyhow!("access denied"));
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new(".*").unwrap());
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let counter = Arc::new(AtomicU64::new(0));
let deleter =
make_object_deleter_with_token(config, boxed, counter, cancellation_token.clone());
let obj = make_s3_object("key/1", 100);
let result = deleter.apply_head_object_filters(&obj).await;
assert!(result.is_err());
assert!(cancellation_token.is_cancelled());
}
#[tokio::test]
async fn apply_head_object_filters_combined_content_type_and_metadata() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("text/plain".to_string());
let mut meta = HashMap::new();
meta.insert("env".to_string(), "production".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("text/.*").unwrap());
config.filter_config.include_metadata_regex = Some(Regex::new("env=production").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_head_object_filters_content_type_passes_metadata_rejects() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("text/plain".to_string());
let mut meta = HashMap::new();
meta.insert("env".to_string(), "staging".to_string());
*mock.head_object_metadata.lock().unwrap() = Some(meta);
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("text/.*").unwrap());
config.filter_config.include_metadata_regex = Some(Regex::new("env=production").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(deleter.apply_head_object_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_tag_filters_no_filters_returns_false() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config();
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_tag_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_tag_filters_include_match_passes() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("env").value("prod").build().unwrap(),
]);
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("env=prod").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_tag_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_tag_filters_include_no_match_filters() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("env").value("staging").build().unwrap(),
]);
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("env=prod").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(deleter.apply_tag_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_tag_filters_exclude_match_filters() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("env").value("dev").build().unwrap(),
]);
let mut config = make_test_config();
config.filter_config.exclude_tag_regex = Some(Regex::new("env=dev").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(deleter.apply_tag_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_tag_filters_exclude_no_match_passes() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("env").value("prod").build().unwrap(),
]);
let mut config = make_test_config();
config.filter_config.exclude_tag_regex = Some(Regex::new("env=dev").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_tag_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_tag_filters_tagging_error_cancels_pipeline() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_error.lock().unwrap() = Some(anyhow::anyhow!("access denied"));
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new(".*").unwrap());
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let counter = Arc::new(AtomicU64::new(0));
let deleter =
make_object_deleter_with_token(config, boxed, counter, cancellation_token.clone());
let obj = make_s3_object("key/1", 100);
let result = deleter.apply_tag_filters(&obj).await;
assert!(result.is_err());
assert!(cancellation_token.is_cancelled());
}
#[tokio::test]
async fn apply_tag_filters_empty_tags_include_rejects() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![]);
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("env=prod").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(deleter.apply_tag_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_tag_filters_multiple_tags_sorted() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("z-tag").value("zval").build().unwrap(),
Tag::builder().key("a-tag").value("aval").build().unwrap(),
]);
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("a-tag=aval&z-tag=zval").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_tag_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_tag_filters_combined_include_and_exclude() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("env").value("prod").build().unwrap(),
Tag::builder().key("team").value("backend").build().unwrap(),
]);
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("env=prod").unwrap());
config.filter_config.exclude_tag_regex = Some(Regex::new("team=frontend").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(!deleter.apply_tag_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn apply_tag_filters_include_passes_exclude_rejects() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("env").value("prod").build().unwrap(),
Tag::builder().key("team").value("backend").build().unwrap(),
]);
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("env=prod").unwrap());
config.filter_config.exclude_tag_regex = Some(Regex::new("team=backend").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let obj = make_s3_object("key/1", 100);
assert!(deleter.apply_tag_filters(&obj).await.unwrap());
}
#[tokio::test]
async fn buffer_and_delete_adds_to_buffer() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.batch_size = 1000; let counter = Arc::new(AtomicU64::new(0));
let mut deleter = make_object_deleter_for_unit_test(config, boxed, counter);
assert_eq!(deleter.buffer.len(), 0);
deleter
.buffer_and_delete(make_s3_object("key/1", 100))
.await
.unwrap();
assert_eq!(deleter.buffer.len(), 1);
deleter
.buffer_and_delete(make_s3_object("key/2", 200))
.await
.unwrap();
assert_eq!(deleter.buffer.len(), 2);
}
#[tokio::test]
async fn buffer_and_delete_flushes_at_batch_size() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.batch_size = 2; let counter = Arc::new(AtomicU64::new(0));
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let stage = make_stage_with_mock(config, boxed, None, Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, counter);
deleter
.buffer_and_delete(make_s3_object("key/1", 100))
.await
.unwrap();
assert_eq!(deleter.buffer.len(), 1);
assert_eq!(mock.delete_objects_calls.lock().unwrap().len(), 0);
deleter
.buffer_and_delete(make_s3_object("key/2", 200))
.await
.unwrap();
assert_eq!(deleter.buffer.len(), 0); assert_eq!(mock.delete_objects_calls.lock().unwrap().len(), 1);
}
#[tokio::test]
async fn handle_api_error_non_404_returns_err_and_cancels() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config();
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let counter = Arc::new(AtomicU64::new(0));
let deleter =
make_object_deleter_with_token(config, boxed, counter, cancellation_token.clone());
let err = anyhow::anyhow!("access denied");
let result = deleter.handle_api_error(&err, "key/1", "head_object").await;
assert!(result.is_err());
let msg = result.unwrap_err().to_string();
assert!(msg.contains("head_object failed for key: key/1"));
assert!(cancellation_token.is_cancelled());
}
#[tokio::test]
async fn handle_api_error_with_different_operation_name() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config();
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let err = anyhow::anyhow!("throttled");
let result = deleter
.handle_api_error(&err, "key/abc", "get_object_tagging")
.await;
assert!(result.is_err());
let msg = result.unwrap_err().to_string();
assert!(msg.contains("get_object_tagging failed for key: key/abc"));
}
#[tokio::test]
async fn process_object_no_filters_buffers_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config();
let counter = Arc::new(AtomicU64::new(0));
let mut deleter = make_object_deleter_for_unit_test(config, boxed, counter);
deleter
.process_object(make_s3_object("key/1", 100))
.await
.unwrap();
assert_eq!(deleter.buffer.len(), 1);
assert_eq!(deleter.buffer[0].key(), "key/1");
}
#[tokio::test]
async fn process_object_max_delete_exceeded_does_not_buffer() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let mut config = make_test_config();
config.max_delete = Some(1);
let counter = Arc::new(AtomicU64::new(1)); let mut deleter = make_object_deleter_for_unit_test(config, boxed, counter);
deleter
.process_object(make_s3_object("key/1", 100))
.await
.unwrap();
assert_eq!(deleter.buffer.len(), 0); }
#[tokio::test]
async fn process_object_head_filter_rejects_does_not_buffer() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("image/png".to_string());
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("text/.*").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let mut deleter = make_object_deleter_for_unit_test(config, boxed, counter);
deleter
.process_object(make_s3_object("key/1", 100))
.await
.unwrap();
assert_eq!(deleter.buffer.len(), 0);
}
#[tokio::test]
async fn process_object_tag_filter_rejects_does_not_buffer() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("env").value("dev").build().unwrap(),
]);
let mut config = make_test_config();
config.filter_config.include_tag_regex = Some(Regex::new("env=prod").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let mut deleter = make_object_deleter_for_unit_test(config, boxed, counter);
deleter
.process_object(make_s3_object("key/1", 100))
.await
.unwrap();
assert_eq!(deleter.buffer.len(), 0);
}
#[tokio::test]
async fn process_object_all_filters_pass_buffers_object() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.head_object_content_type.lock().unwrap() = Some("text/plain".to_string());
*mock.tagging_response_tags.lock().unwrap() = Some(vec![
Tag::builder().key("env").value("prod").build().unwrap(),
]);
let mut config = make_test_config();
config.filter_config.include_content_type_regex = Some(Regex::new("text/.*").unwrap());
config.filter_config.include_tag_regex = Some(Regex::new("env=prod").unwrap());
let counter = Arc::new(AtomicU64::new(0));
let mut deleter = make_object_deleter_for_unit_test(config, boxed, counter);
deleter
.process_object(make_s3_object("key/1", 100))
.await
.unwrap();
assert_eq!(deleter.buffer.len(), 1);
}
fn make_delete_object_sdk_error(code: &str, message: &str) -> anyhow::Error {
use aws_sdk_s3::operation::delete_object::DeleteObjectError;
use aws_smithy_runtime_api::http::StatusCode;
let service_err = DeleteObjectError::generic(
aws_smithy_types::error::ErrorMetadata::builder()
.code(code)
.message(message)
.build(),
);
let raw_response = aws_smithy_runtime_api::http::Response::new(
StatusCode::try_from(403).unwrap(),
SdkBody::from(""),
);
let sdk_error: SdkError<DeleteObjectError, Response<SdkBody>> =
SdkError::service_error(service_err, raw_response);
anyhow::anyhow!(sdk_error).context("aws_sdk_s3::client::delete_object() failed.")
}
#[test]
fn extract_error_code_from_sdk_error_access_denied() {
let err = make_delete_object_sdk_error("AccessDenied", "Access Denied");
assert_eq!(extract_delete_object_error_code(&err), "AccessDenied");
}
#[test]
fn extract_error_code_from_sdk_error_no_such_key() {
let err = make_delete_object_sdk_error("NoSuchKey", "The specified key does not exist.");
assert_eq!(extract_delete_object_error_code(&err), "NoSuchKey");
}
#[test]
fn extract_error_code_from_sdk_error_internal_error() {
let err = make_delete_object_sdk_error("InternalError", "We encountered an internal error.");
assert_eq!(extract_delete_object_error_code(&err), "InternalError");
}
#[test]
fn extract_error_code_from_sdk_error_slow_down() {
let err = make_delete_object_sdk_error("SlowDown", "Please reduce your request rate.");
assert_eq!(extract_delete_object_error_code(&err), "SlowDown");
}
#[test]
fn extract_error_code_from_sdk_error_service_unavailable() {
let err =
make_delete_object_sdk_error("ServiceUnavailable", "Service is temporarily unavailable.");
assert_eq!(extract_delete_object_error_code(&err), "ServiceUnavailable");
}
#[test]
fn extract_error_code_from_plain_anyhow_error_returns_unknown() {
let err = anyhow::anyhow!("connection refused");
assert_eq!(extract_delete_object_error_code(&err), "unknown");
}
#[test]
fn extract_error_code_from_non_service_sdk_error_returns_unknown() {
use aws_sdk_s3::operation::delete_object::DeleteObjectError;
let sdk_error: SdkError<DeleteObjectError, Response<SdkBody>> =
SdkError::timeout_error("connection timed out");
let err = anyhow::anyhow!(sdk_error).context("delete_object() failed.");
assert_eq!(extract_delete_object_error_code(&err), "unknown");
}
#[test]
fn extract_error_code_from_nested_context_chain() {
let base_err = make_delete_object_sdk_error("RequestTimeout", "Request timed out.");
let wrapped = base_err
.context("additional context")
.context("more context");
assert_eq!(extract_delete_object_error_code(&wrapped), "RequestTimeout");
}
#[tokio::test]
async fn single_deleter_error_code_extracted_from_sdk_error() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.delete_object_sdk_error.lock().unwrap() = Some(make_delete_object_sdk_error(
"AccessDenied",
"Access Denied",
));
let deleter = SingleDeleter::new(boxed);
let config = make_test_config();
let objects = vec![make_s3_object("key/denied", 200)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.failed.len(), 1);
assert_eq!(result.failed[0].key, "key/denied");
assert_eq!(
result.failed[0].error_code, "AccessDenied",
"Should extract specific S3 error code, not generic 'DeleteObjectError'"
);
}
#[tokio::test]
async fn single_deleter_error_code_unknown_for_non_sdk_error() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
mock.delete_object_error_keys
.lock()
.unwrap()
.insert("key/fail".to_string(), "connection refused".to_string());
let deleter = SingleDeleter::new(boxed);
let config = make_test_config();
let objects = vec![make_s3_object("key/fail", 200)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.failed.len(), 1);
assert_eq!(
result.failed[0].error_code, "unknown",
"Non-SDK errors should have 'unknown' error code"
);
}
#[tokio::test]
async fn single_deleter_error_code_internal_error() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.delete_object_sdk_error.lock().unwrap() = Some(make_delete_object_sdk_error(
"InternalError",
"Internal server error",
));
let deleter = SingleDeleter::new(boxed);
let config = make_test_config();
let objects = vec![make_s3_object("key/internal", 200)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.failed.len(), 1);
assert_eq!(result.failed[0].error_code, "InternalError");
}
#[tokio::test]
async fn single_deleter_error_message_preserved_with_error_code() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.delete_object_sdk_error.lock().unwrap() = Some(make_delete_object_sdk_error(
"AccessDenied",
"Access Denied",
));
let deleter = SingleDeleter::new(boxed);
let config = make_test_config();
let objects = vec![make_s3_object("key/denied", 200)];
let result = deleter.delete(&objects, &config).await.unwrap();
assert_eq!(result.failed.len(), 1);
assert!(
!result.failed[0].error_message.is_empty(),
"error_message should not be empty"
);
}
#[tokio::test]
async fn single_deleter_sends_delete_error_stats_on_sdk_error() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.delete_object_sdk_error.lock().unwrap() = Some(make_delete_object_sdk_error(
"AccessDenied",
"Access Denied",
));
let deleter = SingleDeleter::new(boxed);
let config = make_test_config();
let objects = vec![make_s3_object("key/denied", 200)];
deleter.delete(&objects, &config).await.unwrap();
let stat = stats_receiver.recv().await.unwrap();
assert!(
matches!(stat, DeletionStatistics::DeleteError { key } if key == "key/denied"),
"Should send DeleteError stat for the failed key"
);
}
struct FailingDeleter {
error_message: String,
}
#[async_trait]
impl Deleter for FailingDeleter {
async fn delete(&self, _objects: &[S3Object], _config: &Config) -> Result<DeleteResult> {
Err(anyhow::anyhow!("{}", self.error_message))
}
}
#[tokio::test]
async fn delete_buffered_objects_error_chain_preserves_original_error() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config();
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let counter = Arc::new(AtomicU64::new(0));
let stage =
make_stage_with_mock_and_token(config, boxed, None, None, cancellation_token.clone());
let stats_report = Arc::new(DeletionStatsReport::new());
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, counter);
deleter.deleter = Box::new(FailingDeleter {
error_message: "S3 InternalError: something went wrong".to_string(),
});
deleter.buffer.push(make_s3_object("key/test", 100));
let result = deleter.delete_buffered_objects().await;
assert!(result.is_err(), "Should return error when deleter fails");
let err = result.unwrap_err();
let err_string = format!("{:#}", err);
assert!(
err_string.contains("delete worker has been cancelled with error"),
"Error should contain the context message, got: {err_string}"
);
assert!(
err_string.contains("S3 InternalError: something went wrong"),
"Error chain should preserve the original error message, got: {err_string}"
);
assert!(cancellation_token.is_cancelled());
}
#[tokio::test]
async fn delete_buffered_objects_error_chain_preserves_sdk_error() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config();
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let counter = Arc::new(AtomicU64::new(0));
let stage =
make_stage_with_mock_and_token(config, boxed, None, None, cancellation_token.clone());
let stats_report = Arc::new(DeletionStatsReport::new());
let mut deleter = ObjectDeleter::new(stage, 0, stats_report, counter);
deleter.deleter = Box::new(FailingDeleter {
error_message: "AccessDenied: User is not authorized".to_string(),
});
deleter.buffer.push(make_s3_object("key/denied", 100));
let result = deleter.delete_buffered_objects().await;
assert!(result.is_err());
let err = result.unwrap_err();
let err_string = format!("{:#}", err);
assert!(
err_string.contains("AccessDenied"),
"Error chain should contain the original S3 error code, got: {err_string}"
);
}
#[tokio::test]
async fn handle_api_error_preserves_original_error_in_chain() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config();
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let counter = Arc::new(AtomicU64::new(0));
let deleter =
make_object_deleter_with_token(config, boxed, counter, cancellation_token.clone());
let original_error = anyhow::anyhow!("connection refused by remote host");
let result = deleter
.handle_api_error(&original_error, "test-key", "head_object")
.await;
assert!(result.is_err(), "Non-404 error should return Err");
let err = result.unwrap_err();
let err_string = format!("{:#}", err);
assert!(
err_string.contains("head_object failed for key: test-key"),
"Error should contain operation context, got: {err_string}"
);
assert!(
err_string.contains("connection refused by remote host"),
"Error chain should preserve original error message, got: {err_string}"
);
assert!(cancellation_token.is_cancelled());
}
#[tokio::test]
async fn handle_api_error_preserves_different_operation_names() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config();
let counter = Arc::new(AtomicU64::new(0));
let deleter = make_object_deleter_for_unit_test(config, boxed, counter);
let original_error = anyhow::anyhow!("timeout reached");
let result = deleter
.handle_api_error(&original_error, "some/key.txt", "get_object_tagging")
.await;
assert!(result.is_err());
let err = result.unwrap_err();
let err_string = format!("{:#}", err);
assert!(
err_string.contains("get_object_tagging failed for key: some/key.txt"),
"Should include operation name and key, got: {err_string}"
);
assert!(
err_string.contains("timeout reached"),
"Should preserve original error, got: {err_string}"
);
}
#[tokio::test]
async fn handle_api_error_not_found_skips_without_error() {
use aws_sdk_s3::operation::head_object::HeadObjectError;
use aws_smithy_runtime_api::http::StatusCode;
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let (boxed, _mock) = make_mock_storage_boxed(stats_sender);
let config = make_test_config();
let cancellation_token: PipelineCancellationToken = CancellationToken::new();
let counter = Arc::new(AtomicU64::new(0));
let deleter =
make_object_deleter_with_token(config, boxed, counter, cancellation_token.clone());
let service_err = HeadObjectError::NotFound(
aws_sdk_s3::types::error::NotFound::builder()
.message("Not Found")
.build(),
);
let raw_response = aws_smithy_runtime_api::http::Response::new(
StatusCode::try_from(404).unwrap(),
SdkBody::from(""),
);
let sdk_error: SdkError<HeadObjectError, Response<SdkBody>> =
SdkError::service_error(service_err, raw_response);
let not_found_error: anyhow::Error = sdk_error.into();
let result = deleter
.handle_api_error(¬_found_error, "missing-key", "head_object")
.await;
assert!(result.is_ok());
assert!(result.unwrap(), "Should return true to skip the object");
assert!(!cancellation_token.is_cancelled());
let stat = stats_receiver.recv().await.unwrap();
assert!(matches!(stat, DeletionStatistics::DeleteSkip { key } if key == "missing-key"));
}
#[tokio::test]
async fn object_deleter_single_mode_sdk_error_preserves_error_code_in_event() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let (boxed, mock) = make_mock_storage_boxed(stats_sender);
*mock.delete_object_sdk_error.lock().unwrap() = Some(make_delete_object_sdk_error(
"AccessDenied",
"Access Denied",
));
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, output_receiver) = async_channel::bounded::<S3Object>(10);
let mut config = make_test_config();
config.batch_size = 1;
let stage = make_stage_with_mock(config, boxed, Some(input_receiver), Some(output_sender));
let stats_report = Arc::new(DeletionStatsReport::new());
let delete_counter = Arc::new(AtomicU64::new(0));
let mut deleter = ObjectDeleter::new(stage, 0, stats_report.clone(), delete_counter);
input_sender
.send(make_s3_object("key/denied", 100))
.await
.unwrap();
drop(input_sender);
deleter.delete().await.unwrap();
assert_eq!(
stats_report.stats_failed_objects.load(Ordering::SeqCst),
1,
"Should record one failed deletion"
);
let forwarded = output_receiver.try_recv();
assert!(
forwarded.is_ok(),
"Object should still be forwarded to next stage"
);
drop(output_receiver);
}
mod extract_error_code_properties {
use super::*;
proptest! {
#![proptest_config(ProptestConfig {
cases: 50,
timeout: 5000,
..ProptestConfig::default()
})]
#[test]
fn never_panics_on_arbitrary_anyhow_error(msg in "[a-zA-Z0-9 _/.-]{1,100}") {
let err = anyhow::anyhow!("{}", msg);
let code = extract_delete_object_error_code(&err);
prop_assert_eq!(code, "unknown");
}
#[test]
fn preserves_known_error_codes(
code in prop::sample::select(vec![
"AccessDenied",
"InternalError",
"NoSuchKey",
"SlowDown",
"ServiceUnavailable",
"RequestTimeout",
"InvalidObjectState",
"NoSuchBucket",
])
) {
let err = make_delete_object_sdk_error(code, "test message");
let extracted = extract_delete_object_error_code(&err);
prop_assert_eq!(extracted, code);
}
}
}