use std::collections::VecDeque;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use anyhow::Result;
use async_channel::Receiver;
use tokio::task::JoinHandle;
use tracing::{debug, error, info};
use crate::config::Config;
use crate::deleter::ObjectDeleter;
use crate::filters::user_defined::UserDefinedFilter;
use crate::filters::{
DeleteMarkerOnlyFilter, ExcludeRegexFilter, IncludeRegexFilter, KeepLatestOnlyFilter,
LargerSizeFilter, MtimeAfterFilter, MtimeBeforeFilter, ObjectFilter, SmallerSizeFilter,
};
use crate::lister::ObjectLister;
use crate::safety::SafetyChecker;
use crate::stage::Stage;
use crate::storage::{self, Storage};
use crate::terminator::Terminator;
use crate::types::event_callback::{EventData, EventType};
use crate::types::token::PipelineCancellationToken;
use crate::types::{DeletionStatistics, DeletionStats, DeletionStatsReport, S3Object};
pub struct DeletionPipeline {
config: Config,
target: Storage,
cancellation_token: PipelineCancellationToken,
stats_receiver: Receiver<DeletionStatistics>,
has_error: Arc<AtomicBool>,
has_panic: Arc<AtomicBool>,
has_warning: Arc<AtomicBool>,
errors: Arc<Mutex<VecDeque<anyhow::Error>>>,
ready: bool,
prerequisites_checked: bool,
deletion_stats_report: Arc<DeletionStatsReport>,
}
impl DeletionPipeline {
pub async fn new(config: Config, cancellation_token: PipelineCancellationToken) -> Self {
let has_warning = Arc::new(AtomicBool::new(false));
let (stats_sender, stats_receiver) = async_channel::unbounded();
let target = storage::create_storage(
config.clone(),
cancellation_token.clone(),
stats_sender,
has_warning.clone(),
)
.await;
Self {
config,
target,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
}
}
pub async fn run(&mut self) {
assert!(self.ready, "DeletionPipeline::run() called more than once");
self.ready = false;
self.config
.event_manager
.trigger_event(EventData::new(EventType::PIPELINE_START))
.await;
if !self.prerequisites_checked {
if let Err(e) = self.check_prerequisites().await {
self.record_error(e);
self.shutdown();
self.fire_completion_events().await;
return;
}
}
self.execute_pipeline().await;
self.shutdown();
self.fire_completion_events().await;
}
pub fn has_error(&self) -> bool {
self.has_error.load(Ordering::SeqCst)
}
pub fn has_panic(&self) -> bool {
self.has_panic.load(Ordering::SeqCst)
}
pub fn has_warning(&self) -> bool {
self.has_warning.load(Ordering::SeqCst)
}
pub fn get_errors_and_consume(&self) -> Option<Vec<anyhow::Error>> {
if !self.has_error() {
return None;
}
let mut error_list = self
.errors
.lock()
.expect("error list mutex poisoned: a pipeline task panicked while holding the lock");
let mut errors = Vec::with_capacity(error_list.len());
while let Some(e) = error_list.pop_front() {
errors.push(e);
}
Some(errors)
}
pub fn get_error_messages(&self) -> Option<Vec<String>> {
if !self.has_error() {
return None;
}
let error_list = self
.errors
.lock()
.expect("error list mutex poisoned: a pipeline task panicked while holding the lock");
Some(error_list.iter().map(|e| e.to_string()).collect())
}
pub fn get_stats_receiver(&self) -> Receiver<DeletionStatistics> {
self.stats_receiver.clone()
}
pub fn get_deletion_stats(&self) -> DeletionStats {
self.deletion_stats_report.snapshot()
}
pub fn close_stats_sender(&self) {
self.target.get_stats_sender().close();
}
pub async fn check_prerequisites(&mut self) -> Result<()> {
let checker = SafetyChecker::new(&self.config);
checker.check_before_deletion()?;
if self.config.filter_config.keep_latest_only && !self.config.delete_all_versions {
return Err(anyhow::anyhow!(
"--keep-latest-only requires --delete-all-versions."
));
}
if self.config.filter_config.delete_marker_only && !self.config.delete_all_versions {
return Err(anyhow::anyhow!(
"--filter-delete-marker-only requires --delete-all-versions."
));
}
if self.config.if_match && self.config.delete_all_versions {
return Err(anyhow::anyhow!(
"--if-match cannot be used with --delete-all-versions. \
S3 does not support If-Match conditional headers when deleting by version ID."
));
}
if self.config.delete_all_versions {
let versioning_enabled = self.target.is_versioning_enabled().await?;
if !versioning_enabled {
if self.config.filter_config.keep_latest_only {
return Err(anyhow::anyhow!(
"--keep-latest-only cannot be used on a non-versioned bucket."
));
}
if self.config.filter_config.delete_marker_only {
return Err(anyhow::anyhow!(
"--filter-delete-marker-only cannot be used on a non-versioned bucket."
));
}
self.config.delete_all_versions = false;
}
}
self.prerequisites_checked = true;
Ok(())
}
async fn execute_pipeline(&self) {
let listed_objects = self.list_target();
let filtered_objects = self.filter_objects(listed_objects);
let deleted_objects = self.delete_objects(filtered_objects);
let terminator_handle = self.terminate(deleted_objects);
if let Err(e) = terminator_handle.await {
self.has_panic.store(true, Ordering::SeqCst);
error!("terminator task panicked: {}", e);
self.record_error(anyhow::anyhow!("terminator task panicked: {}", e));
}
if self.config.warn_as_error && self.has_warning.load(Ordering::SeqCst) {
self.record_error(anyhow::anyhow!(
"warnings promoted to errors (--warn-as-error)"
));
}
}
fn record_error(&self, error: anyhow::Error) {
self.has_error.store(true, Ordering::SeqCst);
self.errors
.lock()
.expect("error list mutex poisoned: a pipeline task panicked while holding the lock")
.push_back(error);
}
fn shutdown(&self) {
self.close_stats_sender();
}
async fn fire_completion_events(&self) {
if self.has_error() {
let unknown = "Unknown error".to_string();
let mut event_data = EventData::new(EventType::PIPELINE_ERROR);
event_data.message = Some(
self.get_error_messages()
.unwrap_or_default()
.first()
.unwrap_or(&unknown)
.to_string(),
);
self.config.event_manager.trigger_event(event_data).await;
}
if self.cancellation_token.is_cancelled() {
let mut event_data = EventData::new(EventType::DELETE_CANCEL);
event_data.message = Some("Pipeline was cancelled".to_string());
self.config.event_manager.trigger_event(event_data).await;
}
self.config
.event_manager
.trigger_event(EventData::new(EventType::PIPELINE_END))
.await;
}
fn create_spsc_stage(
&self,
previous_stage_receiver: Option<Receiver<S3Object>>,
has_warning: Arc<AtomicBool>,
) -> (Stage, Receiver<S3Object>) {
let (sender, next_stage_receiver) =
async_channel::bounded::<S3Object>(self.config.object_listing_queue_size as usize);
let stage = Stage::new(
self.config.clone(),
dyn_clone::clone_box(&*self.target),
previous_stage_receiver,
Some(sender),
self.cancellation_token.clone(),
has_warning,
);
(stage, next_stage_receiver)
}
fn create_mpmc_stage(
&self,
sender: async_channel::Sender<S3Object>,
receiver: Receiver<S3Object>,
has_warning: Arc<AtomicBool>,
) -> Stage {
Stage::new(
self.config.clone(),
dyn_clone::clone_box(&*self.target),
Some(receiver),
Some(sender),
self.cancellation_token.clone(),
has_warning,
)
}
fn spawn_filter(&self, filter: Box<dyn ObjectFilter + Send + Sync>) {
let has_error = self.has_error.clone();
let has_panic = self.has_panic.clone();
let error_list = self.errors.clone();
let cancellation_token = self.cancellation_token.clone();
tokio::spawn(async move {
let join_result = tokio::spawn(async move { filter.filter().await }).await;
match join_result {
Ok(Ok(())) => {}
Ok(Err(e)) => {
cancellation_token.cancel();
has_error.store(true, Ordering::SeqCst);
error!("filter stage failed: {}", e);
error_list
.lock()
.expect("error list mutex poisoned")
.push_back(e);
}
Err(e) => {
cancellation_token.cancel();
has_error.store(true, Ordering::SeqCst);
has_panic.store(true, Ordering::SeqCst);
error!("filter task panicked: {}", e);
error_list
.lock()
.expect("error list mutex poisoned")
.push_back(anyhow::anyhow!("filter task panicked: {}", e));
}
}
});
}
fn list_target(&self) -> Receiver<S3Object> {
let (stage, receiver) = self.create_spsc_stage(None, self.has_warning.clone());
let max_keys = self.config.max_keys;
let has_error = self.has_error.clone();
let has_panic = self.has_panic.clone();
let error_list = self.errors.clone();
let cancellation_token = self.cancellation_token.clone();
tokio::spawn(async move {
let lister = ObjectLister::new(stage);
let join_result = tokio::spawn(async move { lister.list_target(max_keys).await }).await;
match join_result {
Ok(Ok(())) => {
debug!("object lister completed successfully.");
}
Ok(Err(e)) => {
cancellation_token.cancel();
has_error.store(true, Ordering::SeqCst);
error!("object lister failed: {}", e);
error_list
.lock()
.expect("error list mutex poisoned")
.push_back(e);
}
Err(e) => {
cancellation_token.cancel();
has_error.store(true, Ordering::SeqCst);
has_panic.store(true, Ordering::SeqCst);
error!("object lister task panicked: {}", e);
error_list
.lock()
.expect("error list mutex poisoned")
.push_back(anyhow::anyhow!("object lister task panicked: {}", e));
}
}
});
receiver
}
fn filter_objects(&self, objects_list: Receiver<S3Object>) -> Receiver<S3Object> {
let mut previous_stage_receiver = objects_list;
if self.config.filter_config.delete_marker_only {
let (stage, new_receiver) =
self.create_spsc_stage(Some(previous_stage_receiver), self.has_warning.clone());
self.spawn_filter(Box::new(DeleteMarkerOnlyFilter::new(stage)));
previous_stage_receiver = new_receiver;
}
if self.config.filter_config.before_time.is_some() {
let (stage, new_receiver) =
self.create_spsc_stage(Some(previous_stage_receiver), self.has_warning.clone());
self.spawn_filter(Box::new(MtimeBeforeFilter::new(stage)));
previous_stage_receiver = new_receiver;
}
if self.config.filter_config.after_time.is_some() {
let (stage, new_receiver) =
self.create_spsc_stage(Some(previous_stage_receiver), self.has_warning.clone());
self.spawn_filter(Box::new(MtimeAfterFilter::new(stage)));
previous_stage_receiver = new_receiver;
}
if self.config.filter_config.smaller_size.is_some() {
let (stage, new_receiver) =
self.create_spsc_stage(Some(previous_stage_receiver), self.has_warning.clone());
self.spawn_filter(Box::new(SmallerSizeFilter::new(stage)));
previous_stage_receiver = new_receiver;
}
if self.config.filter_config.larger_size.is_some() {
let (stage, new_receiver) =
self.create_spsc_stage(Some(previous_stage_receiver), self.has_warning.clone());
self.spawn_filter(Box::new(LargerSizeFilter::new(stage)));
previous_stage_receiver = new_receiver;
}
if self.config.filter_config.include_regex.is_some() {
let (stage, new_receiver) =
self.create_spsc_stage(Some(previous_stage_receiver), self.has_warning.clone());
self.spawn_filter(Box::new(IncludeRegexFilter::new(stage)));
previous_stage_receiver = new_receiver;
}
if self.config.filter_config.exclude_regex.is_some() {
let (stage, new_receiver) =
self.create_spsc_stage(Some(previous_stage_receiver), self.has_warning.clone());
self.spawn_filter(Box::new(ExcludeRegexFilter::new(stage)));
previous_stage_receiver = new_receiver;
}
if self.config.filter_config.keep_latest_only {
let (stage, new_receiver) =
self.create_spsc_stage(Some(previous_stage_receiver), self.has_warning.clone());
self.spawn_filter(Box::new(KeepLatestOnlyFilter::new(stage)));
previous_stage_receiver = new_receiver;
}
if self.config.filter_manager.is_callback_registered() {
let (stage, new_receiver) =
self.create_spsc_stage(Some(previous_stage_receiver), self.has_warning.clone());
let has_error = self.has_error.clone();
let has_panic = self.has_panic.clone();
let error_list = self.errors.clone();
let cancellation_token = self.cancellation_token.clone();
tokio::spawn(async move {
let filter = UserDefinedFilter::new(stage);
let join_result = tokio::spawn(async move { filter.filter().await }).await;
match join_result {
Ok(Ok(())) => {}
Ok(Err(e)) => {
cancellation_token.cancel();
has_error.store(true, Ordering::SeqCst);
error!("user defined filter failed: {}", e);
error_list
.lock()
.expect("error list mutex poisoned")
.push_back(e);
}
Err(e) => {
cancellation_token.cancel();
has_error.store(true, Ordering::SeqCst);
has_panic.store(true, Ordering::SeqCst);
error!("user defined filter task panicked: {}", e);
error_list
.lock()
.unwrap()
.push_back(anyhow::anyhow!("user defined filter panicked: {}", e));
}
}
});
previous_stage_receiver = new_receiver;
}
previous_stage_receiver
}
fn delete_objects(&self, objects_to_be_deleted: Receiver<S3Object>) -> Receiver<S3Object> {
let (sender, next_stage_receiver) =
async_channel::bounded::<S3Object>(self.config.object_listing_queue_size as usize);
let delete_counter = Arc::new(AtomicU64::new(0));
for worker_index in 0..self.config.worker_size {
let stage = self.create_mpmc_stage(
sender.clone(),
objects_to_be_deleted.clone(),
self.has_warning.clone(),
);
let mut object_deleter = ObjectDeleter::new(
stage,
worker_index,
self.deletion_stats_report.clone(),
delete_counter.clone(),
);
let has_error = self.has_error.clone();
let has_panic = self.has_panic.clone();
let error_list = self.errors.clone();
let cancellation_token = self.cancellation_token.clone();
tokio::spawn(async move {
let join_result = tokio::spawn(async move { object_deleter.delete().await }).await;
match join_result {
Ok(Ok(())) => {
debug!(worker_index, "delete worker completed successfully.");
}
Ok(Err(e)) => {
if crate::types::error::is_cancelled_error(&e) {
info!(worker_index, "delete worker cancelled.");
} else {
cancellation_token.cancel();
has_error.store(true, Ordering::SeqCst);
error!(worker_index, "delete worker failed: {}", e);
error_list
.lock()
.expect("error list mutex poisoned")
.push_back(e);
}
}
Err(e) => {
cancellation_token.cancel();
has_error.store(true, Ordering::SeqCst);
has_panic.store(true, Ordering::SeqCst);
error!(worker_index, "delete worker task panicked: {}", e);
error_list
.lock()
.unwrap()
.push_back(anyhow::anyhow!("delete worker panicked: {}", e));
}
}
});
}
drop(sender);
next_stage_receiver
}
fn terminate(&self, deleted_objects: Receiver<S3Object>) -> JoinHandle<()> {
let terminator = Terminator::new(deleted_objects);
tokio::spawn(async move {
terminator.terminate().await;
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::TracingConfig;
use crate::filters::tests::{create_mock_storage, create_test_config};
use crate::storage::StorageTrait;
use crate::test_utils::init_dummy_tracing_subscriber;
use crate::types::DeletionStatistics;
use crate::types::token::create_pipeline_cancellation_token;
use async_trait::async_trait;
use aws_sdk_s3::primitives::DateTime;
use aws_sdk_s3::types::Object;
use proptest::prelude::*;
use std::sync::atomic::AtomicBool;
#[tokio::test]
async fn spsc_stage_creates_connected_channels() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token: cancellation_token.clone(),
stats_receiver: async_channel::unbounded().1,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
let (input_sender, input_receiver) = async_channel::bounded::<S3Object>(10);
let (stage, output_receiver) =
pipeline.create_spsc_stage(Some(input_receiver), pipeline.has_warning.clone());
assert!(stage.receiver.is_some());
assert!(stage.sender.is_some());
let obj = S3Object::NotVersioning(Object::builder().key("test-key").build());
input_sender.send(obj).await.unwrap();
input_sender.close();
let received = stage.receiver.as_ref().unwrap().recv().await.unwrap();
assert_eq!(received.key(), "test-key");
stage.sender.as_ref().unwrap().send(received).await.unwrap();
let output = output_receiver.recv().await.unwrap();
assert_eq!(output.key(), "test-key");
}
#[tokio::test]
async fn mpmc_stage_shares_channels() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token: cancellation_token.clone(),
stats_receiver: async_channel::unbounded().1,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
let (shared_sender, shared_receiver) = async_channel::bounded::<S3Object>(10);
let (output_sender, _output_receiver) = async_channel::bounded::<S3Object>(10);
let stage1 = pipeline.create_mpmc_stage(
output_sender.clone(),
shared_receiver.clone(),
pipeline.has_warning.clone(),
);
let stage2 = pipeline.create_mpmc_stage(
output_sender,
shared_receiver,
pipeline.has_warning.clone(),
);
assert!(stage1.receiver.is_some());
assert!(stage1.sender.is_some());
assert!(stage2.receiver.is_some());
assert!(stage2.sender.is_some());
let obj1 = S3Object::NotVersioning(Object::builder().key("key1").build());
let obj2 = S3Object::NotVersioning(Object::builder().key("key2").build());
shared_sender.send(obj1).await.unwrap();
shared_sender.send(obj2).await.unwrap();
let r1 = stage1.receiver.as_ref().unwrap().recv().await.unwrap();
let r2 = stage2.receiver.as_ref().unwrap().recv().await.unwrap();
let mut keys: Vec<String> = vec![r1.key().to_string(), r2.key().to_string()];
keys.sort();
assert_eq!(keys, vec!["key1", "key2"]);
}
#[tokio::test]
async fn has_error_initially_false() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
assert!(!pipeline.has_error());
assert!(!pipeline.has_warning());
assert!(pipeline.get_errors_and_consume().is_none());
}
#[tokio::test]
async fn record_error_sets_flag_and_stores_error() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.record_error(anyhow::anyhow!("test error"));
assert!(pipeline.has_error());
let errors = pipeline.get_errors_and_consume().unwrap();
assert_eq!(errors.len(), 1);
assert_eq!(errors[0].to_string(), "test error");
assert!(pipeline.has_error());
}
#[tokio::test]
async fn deletion_stats_initially_zero() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
let stats = pipeline.get_deletion_stats();
assert_eq!(stats.stats_deleted_objects, 0);
assert_eq!(stats.stats_failed_objects, 0);
assert_eq!(stats.stats_deleted_bytes, 0);
}
#[tokio::test]
async fn cancellation_propagates_to_stages() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token: cancellation_token.clone(),
stats_receiver: async_channel::unbounded().1,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
let (stage, _receiver) = pipeline.create_spsc_stage(None, pipeline.has_warning.clone());
assert!(!stage.cancellation_token.is_cancelled());
cancellation_token.cancel();
assert!(stage.cancellation_token.is_cancelled());
}
#[derive(Clone)]
struct ListingMockStorage {
objects: Vec<S3Object>,
stats_sender: async_channel::Sender<DeletionStatistics>,
has_warning: Arc<AtomicBool>,
}
#[async_trait]
impl StorageTrait for ListingMockStorage {
fn is_express_onezone_storage(&self) -> bool {
false
}
async fn list_objects(
&self,
sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
for obj in &self.objects {
sender.send(obj.clone()).await.unwrap();
}
Ok(())
}
async fn list_object_versions(
&self,
_sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
Ok(())
}
async fn head_object(
&self,
_relative_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::head_object::HeadObjectOutput> {
Ok(aws_sdk_s3::operation::head_object::HeadObjectOutput::builder().build())
}
async fn get_object_tagging(
&self,
_relative_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput> {
Ok(
aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput::builder()
.set_tag_set(Some(vec![]))
.build()
.unwrap(),
)
}
async fn delete_object(
&self,
_relative_key: &str,
_version_id: Option<String>,
_if_match: Option<String>,
) -> Result<aws_sdk_s3::operation::delete_object::DeleteObjectOutput> {
Ok(aws_sdk_s3::operation::delete_object::DeleteObjectOutput::builder().build())
}
async fn delete_objects(
&self,
objects: Vec<aws_sdk_s3::types::ObjectIdentifier>,
) -> Result<aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput> {
use aws_sdk_s3::types::DeletedObject;
let deleted: Vec<DeletedObject> = objects
.into_iter()
.map(|oi| {
DeletedObject::builder()
.key(oi.key())
.set_version_id(oi.version_id().map(|s| s.to_string()))
.build()
})
.collect();
Ok(
aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput::builder()
.set_deleted(Some(deleted))
.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) -> async_channel::Sender<DeletionStatistics> {
self.stats_sender.clone()
}
async fn send_stats(&self, stats: DeletionStatistics) {
let _ = self.stats_sender.send(stats).await;
}
fn set_warning(&self) {
self.has_warning
.store(true, std::sync::atomic::Ordering::SeqCst);
}
}
#[tokio::test]
async fn pipeline_runs_with_mock_storage_dry_run() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("file1.txt")
.size(100)
.last_modified(DateTime::from_secs(0))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("file2.txt")
.size(200)
.last_modified(DateTime::from_secs(0))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.dry_run = true;
config.force = true;
config.batch_size = 1;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
}
#[tokio::test]
async fn pipeline_runs_with_mock_storage_and_deletes() {
init_dummy_tracing_subscriber();
let objects = vec![S3Object::NotVersioning(
Object::builder()
.key("delete-me.txt")
.size(50)
.last_modified(DateTime::from_secs(0))
.build(),
)];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true; config.batch_size = 1;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(stats.stats_deleted_objects, 1);
assert_eq!(stats.stats_deleted_bytes, 50);
}
#[tokio::test]
async fn pipeline_empty_listing_completes() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(stats.stats_deleted_objects, 0);
}
#[tokio::test]
async fn pipeline_cancellation_stops_processing() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("key1")
.size(10)
.last_modified(DateTime::from_secs(0))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("key2")
.size(20)
.last_modified(DateTime::from_secs(0))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("key3")
.size(30)
.last_modified(DateTime::from_secs(0))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
let cancellation_token = create_pipeline_cancellation_token();
cancellation_token.cancel();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
let report = pipeline.deletion_stats_report.clone();
assert_eq!(
report.stats_deleted_objects.load(Ordering::Relaxed),
0,
"Pre-cancelled pipeline must not delete objects"
);
assert_eq!(
report.stats_deleted_bytes.load(Ordering::Relaxed),
0,
"Pre-cancelled pipeline must not report deleted bytes"
);
assert!(
!pipeline.has_error(),
"Clean cancellation should not set error flag"
);
}
#[tokio::test]
async fn pipeline_with_filters() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("include-me.log")
.size(100)
.last_modified(DateTime::from_secs(0))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("exclude-me.txt")
.size(200)
.last_modified(DateTime::from_secs(0))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
config.filter_config.include_regex = Some(fancy_regex::Regex::new(r"\.log$").unwrap());
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(stats.stats_deleted_objects, 1);
assert_eq!(stats.stats_deleted_bytes, 100);
}
#[tokio::test]
async fn pipeline_multiple_workers() {
init_dummy_tracing_subscriber();
let mut objects = Vec::new();
for i in 0..20 {
objects.push(S3Object::NotVersioning(
Object::builder()
.key(format!("file{i}.txt"))
.size(10)
.last_modified(DateTime::from_secs(0))
.build(),
));
}
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
config.worker_size = 4;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(stats.stats_deleted_objects, 20);
assert_eq!(stats.stats_deleted_bytes, 200);
}
#[tokio::test]
async fn pipeline_batch_deletion() {
init_dummy_tracing_subscriber();
let mut objects = Vec::new();
for i in 0..5 {
objects.push(S3Object::NotVersioning(
Object::builder()
.key(format!("batch{i}.txt"))
.size(25)
.last_modified(DateTime::from_secs(0))
.build(),
));
}
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1000;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(stats.stats_deleted_objects, 5);
assert_eq!(stats.stats_deleted_bytes, 125);
}
#[tokio::test]
#[should_panic(expected = "called more than once")]
async fn pipeline_panics_on_double_run() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
pipeline.run().await; }
#[tokio::test]
async fn pipeline_safety_check_cancelled() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
let cancellation_token = create_pipeline_cancellation_token();
cancellation_token.cancel();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
}
#[tokio::test]
async fn pipeline_delete_all_versions_on_unversioned_bucket_clears_flag() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.delete_all_versions = true;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
assert!(!pipeline.config.delete_all_versions);
}
#[tokio::test]
async fn pipeline_keep_latest_only_without_delete_all_versions_errors() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.filter_config.keep_latest_only = true;
config.delete_all_versions = false;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
let result = pipeline.check_prerequisites().await;
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("--keep-latest-only requires --delete-all-versions")
);
}
#[tokio::test]
async fn pipeline_if_match_with_delete_all_versions_errors() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.if_match = true;
config.delete_all_versions = true;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
let result = pipeline.check_prerequisites().await;
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("--if-match cannot be used with --delete-all-versions")
);
}
#[tokio::test]
async fn pipeline_delete_marker_only_without_delete_all_versions_errors() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.filter_config.delete_marker_only = true;
config.delete_all_versions = false;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
let result = pipeline.check_prerequisites().await;
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("--filter-delete-marker-only requires --delete-all-versions")
);
}
#[tokio::test]
async fn pipeline_delete_marker_only_on_non_versioned_bucket_errors() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.filter_config.delete_marker_only = true;
config.delete_all_versions = true;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
let result = pipeline.check_prerequisites().await;
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("--filter-delete-marker-only cannot be used on a non-versioned bucket")
);
}
#[tokio::test]
async fn pipeline_max_delete_cancels_when_threshold_exceeded() {
init_dummy_tracing_subscriber();
let mut objects = Vec::new();
for i in 0..10 {
objects.push(S3Object::NotVersioning(
Object::builder()
.key(format!("file{i}.txt"))
.size(10)
.last_modified(DateTime::from_secs(0))
.build(),
));
}
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1; config.worker_size = 1; config.max_delete = Some(3);
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning: has_warning.clone(),
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
assert!(pipeline.has_warning());
let stats = pipeline.get_deletion_stats();
assert!(
stats.stats_deleted_objects <= 3,
"Expected at most 3 deletions with max_delete=3, got {}",
stats.stats_deleted_objects
);
}
struct RecordingEventCallback {
events: Arc<Mutex<Vec<EventData>>>,
}
#[async_trait]
impl crate::types::event_callback::EventCallback for RecordingEventCallback {
async fn on_event(&mut self, event_data: EventData) {
self.events.lock().unwrap().push(event_data);
}
}
#[tokio::test]
async fn pipeline_event_callback_fires_during_execution() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("event1.txt")
.size(50)
.last_modified(DateTime::from_secs(0))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("event2.txt")
.size(75)
.last_modified(DateTime::from_secs(0))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let recorded_events: Arc<Mutex<Vec<EventData>>> = Arc::new(Mutex::new(Vec::new()));
let callback = RecordingEventCallback {
events: recorded_events.clone(),
};
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
config.event_manager.register_callback(
crate::types::event_callback::EventType::ALL_EVENTS,
callback,
false,
);
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let events = recorded_events.lock().unwrap();
let delete_complete_events: Vec<&EventData> = events
.iter()
.filter(|e| {
e.event_type
.contains(crate::types::event_callback::EventType::DELETE_COMPLETE)
})
.collect();
assert_eq!(
delete_complete_events.len(),
2,
"Expected 2 DELETE_COMPLETE events, got {}",
delete_complete_events.len()
);
let keys: Vec<&str> = delete_complete_events
.iter()
.filter_map(|e| e.key.as_deref())
.collect();
assert!(keys.contains(&"event1.txt"), "Missing event1.txt in events");
assert!(keys.contains(&"event2.txt"), "Missing event2.txt in events");
let end_events: Vec<&EventData> = events
.iter()
.filter(|e| {
e.event_type
.contains(crate::types::event_callback::EventType::PIPELINE_END)
})
.collect();
assert!(
!end_events.is_empty(),
"Expected at least one PIPELINE_END event"
);
}
#[tokio::test]
async fn pipeline_rust_filter_callback_integration() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("keep-me.log")
.size(100)
.last_modified(DateTime::from_secs(0))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("remove-me.txt")
.size(200)
.last_modified(DateTime::from_secs(0))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("also-keep.log")
.size(300)
.last_modified(DateTime::from_secs(0))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
struct LogOnlyFilter;
#[async_trait]
impl crate::types::filter_callback::FilterCallback for LogOnlyFilter {
async fn filter(&mut self, object: &S3Object) -> Result<bool> {
Ok(object.key().ends_with(".log"))
}
}
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
config.filter_manager.register_callback(LogOnlyFilter);
config.test_user_defined_callback = true;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(
stats.stats_deleted_objects, 2,
"Expected 2 deletions (only .log files), got {}",
stats.stats_deleted_objects
);
assert_eq!(
stats.stats_deleted_bytes, 400,
"Expected 400 bytes (100 + 300), got {}",
stats.stats_deleted_bytes
);
}
#[derive(Clone)]
struct PartialFailureMockStorage {
objects: Vec<S3Object>,
stats_sender: async_channel::Sender<DeletionStatistics>,
has_warning: Arc<AtomicBool>,
fail_in_batch: Arc<Vec<String>>,
}
#[async_trait]
impl StorageTrait for PartialFailureMockStorage {
fn is_express_onezone_storage(&self) -> bool {
false
}
async fn list_objects(
&self,
sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
for obj in &self.objects {
sender.send(obj.clone()).await.unwrap();
}
Ok(())
}
async fn list_object_versions(
&self,
_sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
Ok(())
}
async fn head_object(
&self,
_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::head_object::HeadObjectOutput> {
Ok(aws_sdk_s3::operation::head_object::HeadObjectOutput::builder().build())
}
async fn get_object_tagging(
&self,
_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput> {
Ok(
aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput::builder()
.set_tag_set(Some(vec![]))
.build()
.unwrap(),
)
}
async fn delete_object(
&self,
_key: &str,
_version_id: Option<String>,
_if_match: Option<String>,
) -> Result<aws_sdk_s3::operation::delete_object::DeleteObjectOutput> {
Ok(aws_sdk_s3::operation::delete_object::DeleteObjectOutput::builder().build())
}
async fn delete_objects(
&self,
objects: Vec<aws_sdk_s3::types::ObjectIdentifier>,
) -> Result<aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput> {
use aws_sdk_s3::types::{DeletedObject, Error as AwsS3Error};
let mut builder = aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput::builder();
for oi in &objects {
let key = oi.key().to_string();
if self.fail_in_batch.contains(&key) {
builder = builder.errors(
AwsS3Error::builder()
.key(&key)
.code("InternalError")
.message("Transient internal error")
.build(),
);
} else {
builder = builder.deleted(
DeletedObject::builder()
.key(oi.key())
.set_version_id(oi.version_id().map(|s| s.to_string()))
.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) -> async_channel::Sender<DeletionStatistics> {
self.stats_sender.clone()
}
async fn send_stats(&self, stats: DeletionStatistics) {
let _ = self.stats_sender.send(stats).await;
}
fn set_warning(&self) {
self.has_warning
.store(true, std::sync::atomic::Ordering::SeqCst);
}
}
#[tokio::test]
async fn pipeline_partial_batch_failure_retries_via_single_delete() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("ok1.txt")
.size(10)
.last_modified(DateTime::from_secs(0))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("fail1.txt")
.size(20)
.last_modified(DateTime::from_secs(0))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("ok2.txt")
.size(30)
.last_modified(DateTime::from_secs(0))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("fail2.txt")
.size(40)
.last_modified(DateTime::from_secs(0))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let fail_in_batch = Arc::new(vec!["fail1.txt".to_string(), "fail2.txt".to_string()]);
let storage: Storage = Box::new(PartialFailureMockStorage {
objects: objects.clone(),
stats_sender,
has_warning: has_warning.clone(),
fail_in_batch,
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1000; config.worker_size = 1;
config.force_retry_config.force_retry_count = 1;
config.force_retry_config.force_retry_interval_milliseconds = 0;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(
stats.stats_deleted_objects, 4,
"All 4 objects should be deleted (2 batch + 2 fallback), got {}",
stats.stats_deleted_objects
);
assert_eq!(
stats.stats_deleted_bytes, 100,
"Expected 100 bytes (10+20+30+40), got {}",
stats.stats_deleted_bytes
);
}
#[derive(Clone)]
struct FailingListerStorage {
stats_sender: async_channel::Sender<DeletionStatistics>,
has_warning: Arc<AtomicBool>,
}
#[async_trait]
impl StorageTrait for FailingListerStorage {
fn is_express_onezone_storage(&self) -> bool {
false
}
async fn list_objects(
&self,
_sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
Err(anyhow::anyhow!("S3 ListObjects failed: AccessDenied"))
}
async fn list_object_versions(
&self,
_sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
Err(anyhow::anyhow!("S3 ListObjectVersions failed"))
}
async fn head_object(
&self,
_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::head_object::HeadObjectOutput> {
Ok(aws_sdk_s3::operation::head_object::HeadObjectOutput::builder().build())
}
async fn get_object_tagging(
&self,
_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput> {
Ok(
aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput::builder()
.set_tag_set(Some(vec![]))
.build()
.unwrap(),
)
}
async fn delete_object(
&self,
_key: &str,
_version_id: Option<String>,
_if_match: Option<String>,
) -> Result<aws_sdk_s3::operation::delete_object::DeleteObjectOutput> {
Ok(aws_sdk_s3::operation::delete_object::DeleteObjectOutput::builder().build())
}
async fn delete_objects(
&self,
_objects: Vec<aws_sdk_s3::types::ObjectIdentifier>,
) -> Result<aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput> {
Ok(aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput::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) -> async_channel::Sender<DeletionStatistics> {
self.stats_sender.clone()
}
async fn send_stats(&self, stats: DeletionStatistics) {
let _ = self.stats_sender.send(stats).await;
}
fn set_warning(&self) {
self.has_warning
.store(true, std::sync::atomic::Ordering::SeqCst);
}
}
#[tokio::test]
async fn pipeline_listing_failure_produces_error() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(FailingListerStorage {
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(
pipeline.has_error(),
"Pipeline must report error when listing fails"
);
let errors = pipeline.get_errors_and_consume().unwrap();
assert!(
!errors.is_empty(),
"Error list must contain the listing failure"
);
let error_msg = errors[0].to_string();
assert!(
error_msg.contains("AccessDenied") || error_msg.contains("ListObjects"),
"Error message should describe the listing failure, got: {}",
error_msg
);
}
#[tokio::test]
async fn pipeline_event_callback_fires_on_error() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(FailingListerStorage {
stats_sender,
has_warning: has_warning.clone(),
});
let recorded_events: Arc<Mutex<Vec<EventData>>> = Arc::new(Mutex::new(Vec::new()));
let callback = RecordingEventCallback {
events: recorded_events.clone(),
};
let mut config = create_test_config();
config.force = true;
config.event_manager.register_callback(
crate::types::event_callback::EventType::ALL_EVENTS,
callback,
false,
);
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(pipeline.has_error());
let events = recorded_events.lock().unwrap();
let error_events: Vec<&EventData> = events
.iter()
.filter(|e| {
e.event_type
.contains(crate::types::event_callback::EventType::PIPELINE_ERROR)
})
.collect();
assert!(
!error_events.is_empty(),
"Expected PIPELINE_ERROR event on listing failure"
);
let end_events: Vec<&EventData> = events
.iter()
.filter(|e| {
e.event_type
.contains(crate::types::event_callback::EventType::PIPELINE_END)
})
.collect();
assert!(
!end_events.is_empty(),
"Expected PIPELINE_END event even on error path"
);
}
#[derive(Clone)]
struct NonRetryableFailureMockStorage {
objects: Vec<S3Object>,
stats_sender: async_channel::Sender<DeletionStatistics>,
has_warning: Arc<AtomicBool>,
fail_in_batch: Arc<Vec<String>>,
}
#[async_trait]
impl StorageTrait for NonRetryableFailureMockStorage {
fn is_express_onezone_storage(&self) -> bool {
false
}
async fn list_objects(
&self,
sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
for obj in &self.objects {
sender.send(obj.clone()).await.unwrap();
}
Ok(())
}
async fn list_object_versions(
&self,
_sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
Ok(())
}
async fn head_object(
&self,
_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::head_object::HeadObjectOutput> {
Ok(aws_sdk_s3::operation::head_object::HeadObjectOutput::builder().build())
}
async fn get_object_tagging(
&self,
_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput> {
Ok(
aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput::builder()
.set_tag_set(Some(vec![]))
.build()
.unwrap(),
)
}
async fn delete_object(
&self,
_key: &str,
_version_id: Option<String>,
_if_match: Option<String>,
) -> Result<aws_sdk_s3::operation::delete_object::DeleteObjectOutput> {
Ok(aws_sdk_s3::operation::delete_object::DeleteObjectOutput::builder().build())
}
async fn delete_objects(
&self,
objects: Vec<aws_sdk_s3::types::ObjectIdentifier>,
) -> Result<aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput> {
use aws_sdk_s3::types::{DeletedObject, Error as AwsS3Error};
let mut builder = aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput::builder();
for oi in &objects {
let key = oi.key().to_string();
if self.fail_in_batch.contains(&key) {
builder = builder.errors(
AwsS3Error::builder()
.key(&key)
.code("AccessDenied")
.message("Access denied")
.build(),
);
} else {
builder = builder.deleted(
DeletedObject::builder()
.key(oi.key())
.set_version_id(oi.version_id().map(|s| s.to_string()))
.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) -> async_channel::Sender<DeletionStatistics> {
self.stats_sender.clone()
}
async fn send_stats(&self, stats: DeletionStatistics) {
let _ = self.stats_sender.send(stats).await;
}
fn set_warning(&self) {
self.has_warning
.store(true, std::sync::atomic::Ordering::SeqCst);
}
}
#[tokio::test]
async fn pipeline_warn_as_error_true_promotes_warnings_to_errors() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("ok.txt")
.size(10)
.last_modified(DateTime::from_secs(1000))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("fail.txt")
.size(20)
.last_modified(DateTime::from_secs(1000))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let fail_in_batch = Arc::new(vec!["fail.txt".to_string()]);
let storage: Storage = Box::new(NonRetryableFailureMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
fail_in_batch,
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1000;
config.worker_size = 1;
config.warn_as_error = true;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(
pipeline.has_error(),
"Pipeline must report error when warn_as_error=true and warnings exist"
);
assert!(
pipeline.has_warning(),
"Warning flag must be set when partial failures occur"
);
let errors = pipeline.get_errors_and_consume().unwrap();
let error_messages: Vec<String> = errors.iter().map(|e| e.to_string()).collect();
assert!(
error_messages
.iter()
.any(|m| m.contains("warnings promoted to errors")),
"Error message should mention promotion, got: {:?}",
error_messages
);
}
#[tokio::test]
async fn pipeline_warn_as_error_false_does_not_promote_warnings() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("ok.txt")
.size(10)
.last_modified(DateTime::from_secs(1000))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("fail.txt")
.size(20)
.last_modified(DateTime::from_secs(1000))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let fail_in_batch = Arc::new(vec!["fail.txt".to_string()]);
let storage: Storage = Box::new(NonRetryableFailureMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
fail_in_batch,
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1000;
config.worker_size = 1;
config.warn_as_error = false;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(
!pipeline.has_error(),
"Pipeline should not report error when warn_as_error=false"
);
assert!(
pipeline.has_warning(),
"Warning flag should be set when partial failures occur"
);
let stats = pipeline.get_deletion_stats();
assert_eq!(stats.stats_deleted_objects, 1); assert_eq!(stats.stats_failed_objects, 1); }
#[tokio::test]
async fn pipeline_max_delete_stops_at_threshold() {
init_dummy_tracing_subscriber();
let objects: Vec<S3Object> = (0..10)
.map(|i| {
S3Object::NotVersioning(
Object::builder()
.key(format!("obj/{i}"))
.size(50)
.last_modified(DateTime::from_secs(0))
.build(),
)
})
.collect();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
config.max_delete = Some(3);
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
let stats = pipeline.get_deletion_stats();
assert!(
stats.stats_deleted_objects <= 4,
"Expected at most 4 deleted objects (max_delete=3 + 1 triggering), got {}",
stats.stats_deleted_objects
);
assert!(
pipeline.has_warning(),
"Warning flag should be set when max-delete threshold is exceeded"
);
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(10))]
#[test]
fn prop_listing_failure_always_sets_error(_seed in 0u32..50) {
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 has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(FailingListerStorage {
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
prop_assert!(
pipeline.has_error(),
"Pipeline must set error flag when listing fails"
);
Ok(())
})?;
}
}
#[tokio::test]
async fn pipeline_has_panic_initially_false() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
assert!(!pipeline.has_panic());
}
#[tokio::test]
async fn pipeline_has_panic_reflects_atomic_flag() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let has_panic = Arc::new(AtomicBool::new(false));
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: has_panic.clone(),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
assert!(!pipeline.has_panic());
has_panic.store(true, Ordering::SeqCst);
assert!(pipeline.has_panic());
}
#[tokio::test]
async fn pipeline_get_stats_receiver_returns_working_channel() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender.clone(), has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
let receiver = pipeline.get_stats_receiver();
let stats = DeletionStatistics::DeleteComplete {
key: "test-key".to_string(),
};
stats_sender.send(stats).await.unwrap();
let received = receiver.recv().await.unwrap();
match received {
DeletionStatistics::DeleteComplete { key } => assert_eq!(key, "test-key"),
other => panic!("Expected DeleteComplete, got {:?}", other),
}
}
#[tokio::test]
async fn pipeline_get_error_messages_none_when_no_errors() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
assert!(pipeline.get_error_messages().is_none());
}
#[tokio::test]
async fn pipeline_get_error_messages_returns_messages() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.record_error(anyhow::anyhow!("first error"));
pipeline.record_error(anyhow::anyhow!("second error"));
let messages = pipeline.get_error_messages().unwrap();
assert_eq!(messages.len(), 2);
assert_eq!(messages[0], "first error");
assert_eq!(messages[1], "second error");
}
#[tokio::test]
async fn pipeline_prerequisites_failure_records_error_and_stops() {
init_dummy_tracing_subscriber();
let objects = vec![S3Object::NotVersioning(
Object::builder()
.key("should-not-be-deleted.txt")
.size(100)
.last_modified(DateTime::from_secs(0))
.build(),
)];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = false;
config.dry_run = false;
config.tracing_config = Some(TracingConfig {
tracing_level: log::Level::Info,
json_tracing: true,
aws_sdk_tracing: false,
span_events_tracing: false,
disable_color_tracing: false,
});
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false, deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(
pipeline.has_error(),
"Pipeline must set error when prerequisites fail"
);
let errors = pipeline.get_errors_and_consume().unwrap();
assert!(
errors.iter().any(|e| e
.to_string()
.contains("Cannot run destructive operation without --force")),
"Error should mention --force requirement, got: {:?}",
errors.iter().map(|e| e.to_string()).collect::<Vec<_>>()
);
let stats = pipeline.get_deletion_stats();
assert_eq!(stats.stats_deleted_objects, 0);
}
#[tokio::test]
async fn pipeline_mtime_before_filter_applied() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("old.txt")
.size(10)
.last_modified(DateTime::from_secs(1000))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("new.txt")
.size(20)
.last_modified(DateTime::from_secs(2000))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
config.worker_size = 1;
config.filter_config.before_time = Some(chrono::DateTime::from_timestamp(1500, 0).unwrap());
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(
stats.stats_deleted_objects, 1,
"Only one object should pass mtime_before filter, got {}",
stats.stats_deleted_objects
);
assert_eq!(stats.stats_deleted_bytes, 10);
}
#[tokio::test]
async fn pipeline_mtime_after_filter_applied() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("old.txt")
.size(10)
.last_modified(DateTime::from_secs(1000))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("new.txt")
.size(20)
.last_modified(DateTime::from_secs(2000))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
config.worker_size = 1;
config.filter_config.after_time = Some(chrono::DateTime::from_timestamp(1500, 0).unwrap());
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(
stats.stats_deleted_objects, 1,
"Only one object should pass mtime_after filter, got {}",
stats.stats_deleted_objects
);
assert_eq!(stats.stats_deleted_bytes, 20);
}
#[tokio::test]
async fn pipeline_smaller_size_filter_applied() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("small.txt")
.size(50)
.last_modified(DateTime::from_secs(1000))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("large.txt")
.size(500)
.last_modified(DateTime::from_secs(1000))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
config.worker_size = 1;
config.filter_config.smaller_size = Some(100);
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(
stats.stats_deleted_objects, 1,
"Only small.txt should pass smaller_size filter, got {}",
stats.stats_deleted_objects
);
assert_eq!(stats.stats_deleted_bytes, 50);
}
#[tokio::test]
async fn pipeline_larger_size_filter_applied() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("small.txt")
.size(50)
.last_modified(DateTime::from_secs(1000))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("large.txt")
.size(500)
.last_modified(DateTime::from_secs(1000))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
config.worker_size = 1;
config.filter_config.larger_size = Some(100);
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(
stats.stats_deleted_objects, 1,
"Only large.txt should pass larger_size filter, got {}",
stats.stats_deleted_objects
);
assert_eq!(stats.stats_deleted_bytes, 500);
}
#[tokio::test]
async fn pipeline_exclude_regex_filter_applied() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("keep.txt")
.size(10)
.last_modified(DateTime::from_secs(1000))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("remove.log")
.size(20)
.last_modified(DateTime::from_secs(1000))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
config.worker_size = 1;
config.filter_config.exclude_regex = Some(fancy_regex::Regex::new(r".*\.log$").unwrap());
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(
stats.stats_deleted_objects, 1,
"Only keep.txt should remain after exclude_regex, got {}",
stats.stats_deleted_objects
);
assert_eq!(stats.stats_deleted_bytes, 10);
}
#[tokio::test]
async fn pipeline_multiple_filters_combined() {
init_dummy_tracing_subscriber();
let objects = vec![
S3Object::NotVersioning(
Object::builder()
.key("small-old.txt")
.size(50)
.last_modified(DateTime::from_secs(1000))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("small-new.txt")
.size(50)
.last_modified(DateTime::from_secs(3000))
.build(),
),
S3Object::NotVersioning(
Object::builder()
.key("large-old.txt")
.size(500)
.last_modified(DateTime::from_secs(1000))
.build(),
),
];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.batch_size = 1;
config.worker_size = 1;
config.filter_config.smaller_size = Some(100);
config.filter_config.before_time = Some(chrono::DateTime::from_timestamp(2000, 0).unwrap());
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(!pipeline.has_error());
let stats = pipeline.get_deletion_stats();
assert_eq!(
stats.stats_deleted_objects, 1,
"Only small-old.txt should pass both filters, got {}",
stats.stats_deleted_objects
);
assert_eq!(stats.stats_deleted_bytes, 50);
}
#[derive(Clone)]
struct PanickingListerMockStorage {
stats_sender: async_channel::Sender<DeletionStatistics>,
has_warning: Arc<AtomicBool>,
}
#[async_trait]
impl StorageTrait for PanickingListerMockStorage {
fn is_express_onezone_storage(&self) -> bool {
false
}
async fn list_objects(
&self,
_sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
panic!("simulated lister panic");
}
async fn list_object_versions(
&self,
_sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
panic!("simulated lister panic");
}
async fn head_object(
&self,
_relative_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::head_object::HeadObjectOutput> {
unimplemented!()
}
async fn get_object_tagging(
&self,
_relative_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput> {
unimplemented!()
}
async fn delete_object(
&self,
_relative_key: &str,
_version_id: Option<String>,
_if_match: Option<String>,
) -> Result<aws_sdk_s3::operation::delete_object::DeleteObjectOutput> {
unimplemented!()
}
async fn delete_objects(
&self,
_objects: Vec<aws_sdk_s3::types::ObjectIdentifier>,
) -> Result<aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput> {
unimplemented!()
}
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) -> async_channel::Sender<DeletionStatistics> {
self.stats_sender.clone()
}
async fn send_stats(&self, stats: DeletionStatistics) {
let _ = self.stats_sender.send(stats).await;
}
fn set_warning(&self) {
self.has_warning
.store(true, std::sync::atomic::Ordering::SeqCst);
}
}
#[tokio::test]
#[should_panic(expected = "not implemented")]
async fn panicking_lister_mock_head_object_panics_unimplemented() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingListerMockStorage {
stats_sender,
has_warning,
};
let _ = mock.head_object("key", None).await;
}
#[tokio::test]
#[should_panic(expected = "not implemented")]
async fn panicking_lister_mock_get_object_tagging_panics_unimplemented() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingListerMockStorage {
stats_sender,
has_warning,
};
let _ = mock.get_object_tagging("key", None).await;
}
#[tokio::test]
#[should_panic(expected = "not implemented")]
async fn panicking_lister_mock_delete_object_panics_unimplemented() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingListerMockStorage {
stats_sender,
has_warning,
};
let _ = mock.delete_object("key", None, None).await;
}
#[tokio::test]
#[should_panic(expected = "not implemented")]
async fn panicking_lister_mock_delete_objects_panics_unimplemented() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingListerMockStorage {
stats_sender,
has_warning,
};
let _ = mock.delete_objects(vec![]).await;
}
#[tokio::test]
async fn list_target_panic_sets_has_error_and_has_panic() {
init_dummy_tracing_subscriber();
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Box<dyn StorageTrait + Send + Sync> = Box::new(PanickingListerMockStorage {
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.worker_size = 1;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token,
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: true,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(
pipeline.has_error(),
"Pipeline should have error after lister panic"
);
assert!(pipeline.has_panic(), "Pipeline should have panic flag set");
let errors = pipeline.errors.lock().unwrap();
assert_eq!(errors.len(), 1);
let msg = errors[0].to_string();
assert!(
msg.contains("object lister task panicked"),
"Error should mention lister panic, got: {}",
msg
);
}
#[test]
fn panicking_lister_mock_is_express_onezone_returns_false() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingListerMockStorage {
stats_sender,
has_warning,
};
assert!(!mock.is_express_onezone_storage());
}
#[tokio::test]
async fn panicking_lister_mock_is_versioning_enabled_returns_false() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingListerMockStorage {
stats_sender,
has_warning,
};
assert!(!mock.is_versioning_enabled().await.unwrap());
}
#[test]
fn panicking_lister_mock_get_client_returns_none() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingListerMockStorage {
stats_sender,
has_warning,
};
assert!(mock.get_client().is_none());
}
#[tokio::test]
async fn panicking_lister_mock_get_stats_sender_works() {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingListerMockStorage {
stats_sender,
has_warning,
};
let sender = mock.get_stats_sender();
sender
.send(DeletionStatistics::DeleteBytes(55))
.await
.unwrap();
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(received, DeletionStatistics::DeleteBytes(55)));
}
#[tokio::test]
async fn panicking_lister_mock_send_stats_delivers_stat() {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingListerMockStorage {
stats_sender,
has_warning,
};
mock.send_stats(DeletionStatistics::DeleteComplete {
key: "k".to_string(),
})
.await;
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(
received,
DeletionStatistics::DeleteComplete { .. }
));
}
#[test]
fn panicking_lister_mock_set_warning_sets_flag() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingListerMockStorage {
stats_sender,
has_warning: has_warning.clone(),
};
assert!(!has_warning.load(std::sync::atomic::Ordering::SeqCst));
mock.set_warning();
assert!(has_warning.load(std::sync::atomic::Ordering::SeqCst));
}
struct ErrorFilter;
#[async_trait]
impl ObjectFilter for ErrorFilter {
async fn filter(&self) -> Result<()> {
Err(anyhow::anyhow!("simulated filter error"))
}
}
struct PanickingFilter;
#[async_trait]
impl ObjectFilter for PanickingFilter {
async fn filter(&self) -> Result<()> {
panic!("simulated filter panic");
}
}
#[tokio::test]
async fn spawn_filter_error_sets_has_error_and_records_error() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token: cancellation_token.clone(),
stats_receiver: async_channel::unbounded().1,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.spawn_filter(Box::new(ErrorFilter));
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
assert!(
pipeline.has_error(),
"has_error should be set after filter error"
);
assert!(
!pipeline.has_panic(),
"has_panic should NOT be set for a normal error"
);
assert!(
cancellation_token.is_cancelled(),
"cancellation_token should be cancelled"
);
let errors = pipeline.errors.lock().unwrap();
assert_eq!(errors.len(), 1);
let msg = errors[0].to_string();
assert!(
msg.contains("simulated filter error"),
"Error message should contain the filter error, got: {}",
msg
);
}
#[tokio::test]
async fn spawn_filter_panic_sets_has_error_and_has_panic() {
init_dummy_tracing_subscriber();
let (stats_sender, _stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage = create_mock_storage(stats_sender, has_warning.clone());
let config = create_test_config();
let cancellation_token = create_pipeline_cancellation_token();
let pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token: cancellation_token.clone(),
stats_receiver: async_channel::unbounded().1,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: false,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.spawn_filter(Box::new(PanickingFilter));
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
assert!(
pipeline.has_error(),
"has_error should be set after filter panic"
);
assert!(
pipeline.has_panic(),
"has_panic should be set after filter panic"
);
assert!(
cancellation_token.is_cancelled(),
"cancellation_token should be cancelled"
);
let errors = pipeline.errors.lock().unwrap();
assert_eq!(errors.len(), 1);
let msg = errors[0].to_string();
assert!(
msg.contains("filter task panicked"),
"Error message should mention filter panic, got: {}",
msg
);
}
#[tokio::test]
async fn user_defined_filter_panic_sets_has_error_and_has_panic() {
init_dummy_tracing_subscriber();
let objects = vec![S3Object::NotVersioning(
Object::builder()
.key("trigger-panic.txt")
.size(100)
.last_modified(DateTime::from_secs(0))
.build(),
)];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ListingMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
struct PanickingFilterCallback;
#[async_trait]
impl crate::types::filter_callback::FilterCallback for PanickingFilterCallback {
async fn filter(&mut self, _object: &S3Object) -> Result<bool> {
panic!("simulated user defined filter panic");
}
}
let mut config = create_test_config();
config.force = true;
config.worker_size = 1;
config
.filter_manager
.register_callback(PanickingFilterCallback);
config.test_user_defined_callback = true;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token: cancellation_token.clone(),
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: true,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(
pipeline.has_error(),
"has_error should be set after user defined filter panic"
);
assert!(
pipeline.has_panic(),
"has_panic should be set after user defined filter panic"
);
let errors = pipeline.errors.lock().unwrap();
let panic_error = errors.iter().find(|e| {
let msg = e.to_string();
msg.contains("user defined filter panicked")
});
assert!(
panic_error.is_some(),
"Should contain 'user defined filter panicked' error, got: {:?}",
errors.iter().map(|e| e.to_string()).collect::<Vec<_>>()
);
}
#[derive(Clone)]
struct ErrorDeleteMockStorage {
objects: Vec<S3Object>,
stats_sender: async_channel::Sender<DeletionStatistics>,
has_warning: Arc<AtomicBool>,
}
#[async_trait]
impl StorageTrait for ErrorDeleteMockStorage {
fn is_express_onezone_storage(&self) -> bool {
false
}
async fn list_objects(
&self,
sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
for obj in &self.objects {
sender.send(obj.clone()).await.unwrap();
}
Ok(())
}
async fn list_object_versions(
&self,
_sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
Ok(())
}
async fn head_object(
&self,
_relative_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::head_object::HeadObjectOutput> {
Ok(aws_sdk_s3::operation::head_object::HeadObjectOutput::builder().build())
}
async fn get_object_tagging(
&self,
_relative_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput> {
Ok(
aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput::builder()
.set_tag_set(Some(vec![]))
.build()
.unwrap(),
)
}
async fn delete_object(
&self,
_relative_key: &str,
_version_id: Option<String>,
_if_match: Option<String>,
) -> Result<aws_sdk_s3::operation::delete_object::DeleteObjectOutput> {
Err(anyhow::anyhow!("simulated delete error"))
}
async fn delete_objects(
&self,
_objects: Vec<aws_sdk_s3::types::ObjectIdentifier>,
) -> Result<aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput> {
Err(anyhow::anyhow!("simulated batch delete error"))
}
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) -> async_channel::Sender<DeletionStatistics> {
self.stats_sender.clone()
}
async fn send_stats(&self, stats: DeletionStatistics) {
let _ = self.stats_sender.send(stats).await;
}
fn set_warning(&self) {
self.has_warning
.store(true, std::sync::atomic::Ordering::SeqCst);
}
}
#[derive(Clone)]
struct PanickingDeleteMockStorage {
objects: Vec<S3Object>,
stats_sender: async_channel::Sender<DeletionStatistics>,
has_warning: Arc<AtomicBool>,
}
#[async_trait]
impl StorageTrait for PanickingDeleteMockStorage {
fn is_express_onezone_storage(&self) -> bool {
false
}
async fn list_objects(
&self,
sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
for obj in &self.objects {
sender.send(obj.clone()).await.unwrap();
}
Ok(())
}
async fn list_object_versions(
&self,
_sender: &async_channel::Sender<S3Object>,
_max_keys: i32,
) -> Result<()> {
Ok(())
}
async fn head_object(
&self,
_relative_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::head_object::HeadObjectOutput> {
Ok(aws_sdk_s3::operation::head_object::HeadObjectOutput::builder().build())
}
async fn get_object_tagging(
&self,
_relative_key: &str,
_version_id: Option<String>,
) -> Result<aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput> {
Ok(
aws_sdk_s3::operation::get_object_tagging::GetObjectTaggingOutput::builder()
.set_tag_set(Some(vec![]))
.build()
.unwrap(),
)
}
async fn delete_object(
&self,
_relative_key: &str,
_version_id: Option<String>,
_if_match: Option<String>,
) -> Result<aws_sdk_s3::operation::delete_object::DeleteObjectOutput> {
panic!("simulated delete worker panic");
}
async fn delete_objects(
&self,
_objects: Vec<aws_sdk_s3::types::ObjectIdentifier>,
) -> Result<aws_sdk_s3::operation::delete_objects::DeleteObjectsOutput> {
panic!("simulated batch delete worker panic");
}
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) -> async_channel::Sender<DeletionStatistics> {
self.stats_sender.clone()
}
async fn send_stats(&self, stats: DeletionStatistics) {
let _ = self.stats_sender.send(stats).await;
}
fn set_warning(&self) {
self.has_warning
.store(true, std::sync::atomic::Ordering::SeqCst);
}
}
#[tokio::test]
async fn delete_worker_error_sets_has_error() {
init_dummy_tracing_subscriber();
let objects = vec![S3Object::NotVersioning(
Object::builder()
.key("error-obj.txt")
.size(100)
.last_modified(DateTime::from_secs(0))
.build(),
)];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(ErrorDeleteMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.worker_size = 1;
config.batch_size = 1000;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token: cancellation_token.clone(),
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: true,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(
pipeline.has_error(),
"has_error should be set after delete worker error"
);
assert!(
!pipeline.has_panic(),
"has_panic should NOT be set for a normal error"
);
let errors = pipeline.errors.lock().unwrap();
assert!(!errors.is_empty(), "Should have recorded an error");
}
#[tokio::test]
async fn delete_worker_panic_sets_has_error_and_has_panic() {
init_dummy_tracing_subscriber();
let objects = vec![S3Object::NotVersioning(
Object::builder()
.key("panic-obj.txt")
.size(100)
.last_modified(DateTime::from_secs(0))
.build(),
)];
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let storage: Storage = Box::new(PanickingDeleteMockStorage {
objects,
stats_sender,
has_warning: has_warning.clone(),
});
let mut config = create_test_config();
config.force = true;
config.worker_size = 1;
config.batch_size = 1000;
let cancellation_token = create_pipeline_cancellation_token();
let mut pipeline = DeletionPipeline {
config,
target: storage,
cancellation_token: cancellation_token.clone(),
stats_receiver,
has_error: Arc::new(AtomicBool::new(false)),
has_panic: Arc::new(AtomicBool::new(false)),
has_warning,
errors: Arc::new(Mutex::new(VecDeque::new())),
ready: true,
prerequisites_checked: true,
deletion_stats_report: Arc::new(DeletionStatsReport::new()),
};
pipeline.run().await;
assert!(
pipeline.has_error(),
"has_error should be set after delete worker panic"
);
assert!(
pipeline.has_panic(),
"has_panic should be set after delete worker panic"
);
let errors = pipeline.errors.lock().unwrap();
let panic_error = errors.iter().find(|e| {
let msg = e.to_string();
msg.contains("delete worker panicked")
});
assert!(
panic_error.is_some(),
"Should contain 'delete worker panicked' error, got: {:?}",
errors.iter().map(|e| e.to_string()).collect::<Vec<_>>()
);
}
fn make_error_delete_mock() -> (
ErrorDeleteMockStorage,
async_channel::Receiver<DeletionStatistics>,
) {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = ErrorDeleteMockStorage {
objects: vec![],
stats_sender,
has_warning,
};
(mock, stats_receiver)
}
#[test]
fn error_delete_mock_is_express_onezone_returns_false() {
let (mock, _) = make_error_delete_mock();
assert!(!mock.is_express_onezone_storage());
}
#[tokio::test]
async fn error_delete_mock_list_objects_returns_ok() {
let (mock, _) = make_error_delete_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_objects(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn error_delete_mock_list_object_versions_returns_ok() {
let (mock, _) = make_error_delete_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_object_versions(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn error_delete_mock_head_object_returns_ok() {
let (mock, _) = make_error_delete_mock();
assert!(mock.head_object("key", None).await.is_ok());
}
#[tokio::test]
async fn error_delete_mock_get_object_tagging_returns_ok() {
let (mock, _) = make_error_delete_mock();
assert!(mock.get_object_tagging("key", None).await.is_ok());
}
#[tokio::test]
async fn error_delete_mock_delete_object_returns_err() {
let (mock, _) = make_error_delete_mock();
let result = mock.delete_object("key", None, None).await;
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("simulated delete error")
);
}
#[tokio::test]
async fn error_delete_mock_delete_objects_returns_err() {
let (mock, _) = make_error_delete_mock();
let result = mock.delete_objects(vec![]).await;
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("simulated batch delete error")
);
}
#[tokio::test]
async fn error_delete_mock_is_versioning_enabled_returns_false() {
let (mock, _) = make_error_delete_mock();
assert!(!mock.is_versioning_enabled().await.unwrap());
}
#[test]
fn error_delete_mock_get_client_returns_none() {
let (mock, _) = make_error_delete_mock();
assert!(mock.get_client().is_none());
}
#[tokio::test]
async fn error_delete_mock_get_stats_sender_works() {
let (mock, stats_receiver) = make_error_delete_mock();
let sender = mock.get_stats_sender();
sender
.send(DeletionStatistics::DeleteBytes(10))
.await
.unwrap();
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(received, DeletionStatistics::DeleteBytes(10)));
}
#[tokio::test]
async fn error_delete_mock_send_stats_delivers_stat() {
let (mock, stats_receiver) = make_error_delete_mock();
mock.send_stats(DeletionStatistics::DeleteComplete {
key: "k".to_string(),
})
.await;
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(
received,
DeletionStatistics::DeleteComplete { .. }
));
}
#[test]
fn error_delete_mock_set_warning_sets_flag() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = ErrorDeleteMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
};
assert!(!has_warning.load(std::sync::atomic::Ordering::SeqCst));
mock.set_warning();
assert!(has_warning.load(std::sync::atomic::Ordering::SeqCst));
}
fn make_panicking_delete_mock() -> (
PanickingDeleteMockStorage,
async_channel::Receiver<DeletionStatistics>,
) {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingDeleteMockStorage {
objects: vec![],
stats_sender,
has_warning,
};
(mock, stats_receiver)
}
#[test]
fn panicking_delete_mock_is_express_onezone_returns_false() {
let (mock, _) = make_panicking_delete_mock();
assert!(!mock.is_express_onezone_storage());
}
#[tokio::test]
async fn panicking_delete_mock_list_objects_returns_ok() {
let (mock, _) = make_panicking_delete_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_objects(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn panicking_delete_mock_list_object_versions_returns_ok() {
let (mock, _) = make_panicking_delete_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_object_versions(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn panicking_delete_mock_head_object_returns_ok() {
let (mock, _) = make_panicking_delete_mock();
assert!(mock.head_object("key", None).await.is_ok());
}
#[tokio::test]
async fn panicking_delete_mock_get_object_tagging_returns_ok() {
let (mock, _) = make_panicking_delete_mock();
assert!(mock.get_object_tagging("key", None).await.is_ok());
}
#[tokio::test]
async fn panicking_delete_mock_is_versioning_enabled_returns_false() {
let (mock, _) = make_panicking_delete_mock();
assert!(!mock.is_versioning_enabled().await.unwrap());
}
#[test]
fn panicking_delete_mock_get_client_returns_none() {
let (mock, _) = make_panicking_delete_mock();
assert!(mock.get_client().is_none());
}
#[tokio::test]
async fn panicking_delete_mock_get_stats_sender_works() {
let (mock, stats_receiver) = make_panicking_delete_mock();
let sender = mock.get_stats_sender();
sender
.send(DeletionStatistics::DeleteBytes(20))
.await
.unwrap();
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(received, DeletionStatistics::DeleteBytes(20)));
}
#[tokio::test]
async fn panicking_delete_mock_send_stats_delivers_stat() {
let (mock, stats_receiver) = make_panicking_delete_mock();
mock.send_stats(DeletionStatistics::DeleteComplete {
key: "k".to_string(),
})
.await;
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(
received,
DeletionStatistics::DeleteComplete { .. }
));
}
#[test]
fn panicking_delete_mock_set_warning_sets_flag() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PanickingDeleteMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
};
assert!(!has_warning.load(std::sync::atomic::Ordering::SeqCst));
mock.set_warning();
assert!(has_warning.load(std::sync::atomic::Ordering::SeqCst));
}
fn make_listing_mock() -> (
ListingMockStorage,
async_channel::Receiver<DeletionStatistics>,
) {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = ListingMockStorage {
objects: vec![],
stats_sender,
has_warning,
};
(mock, stats_receiver)
}
#[test]
fn listing_mock_is_express_onezone_returns_false() {
let (mock, _) = make_listing_mock();
assert!(!mock.is_express_onezone_storage());
}
#[tokio::test]
async fn listing_mock_list_objects_sends_objects() {
let obj = S3Object::NotVersioning(
Object::builder()
.key("a.txt")
.size(10)
.last_modified(DateTime::from_secs(0))
.build(),
);
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = ListingMockStorage {
objects: vec![obj],
stats_sender,
has_warning,
};
let (sender, receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_objects(&sender, 1000).await.is_ok());
let received = receiver.try_recv().unwrap();
assert_eq!(received.key(), "a.txt");
}
#[tokio::test]
async fn listing_mock_list_object_versions_returns_ok() {
let (mock, _) = make_listing_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_object_versions(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn listing_mock_head_object_returns_ok() {
let (mock, _) = make_listing_mock();
assert!(mock.head_object("key", None).await.is_ok());
}
#[tokio::test]
async fn listing_mock_get_object_tagging_returns_ok() {
let (mock, _) = make_listing_mock();
let result = mock.get_object_tagging("key", None).await;
assert!(result.is_ok());
assert!(result.unwrap().tag_set().is_empty());
}
#[tokio::test]
async fn listing_mock_delete_object_returns_ok() {
let (mock, _) = make_listing_mock();
assert!(mock.delete_object("key", None, None).await.is_ok());
}
#[tokio::test]
async fn listing_mock_delete_objects_returns_deleted_keys() {
let (mock, _) = make_listing_mock();
let idents = vec![
aws_sdk_s3::types::ObjectIdentifier::builder()
.key("x.txt")
.build()
.unwrap(),
aws_sdk_s3::types::ObjectIdentifier::builder()
.key("y.txt")
.build()
.unwrap(),
];
let result = mock.delete_objects(idents).await;
assert!(result.is_ok());
let output = result.unwrap();
let deleted: Vec<&str> = output.deleted().iter().filter_map(|d| d.key()).collect();
assert_eq!(deleted, vec!["x.txt", "y.txt"]);
}
#[tokio::test]
async fn listing_mock_is_versioning_enabled_returns_false() {
let (mock, _) = make_listing_mock();
assert!(!mock.is_versioning_enabled().await.unwrap());
}
#[test]
fn listing_mock_get_client_returns_none() {
let (mock, _) = make_listing_mock();
assert!(mock.get_client().is_none());
}
#[tokio::test]
async fn listing_mock_get_stats_sender_works() {
let (mock, stats_receiver) = make_listing_mock();
let sender = mock.get_stats_sender();
sender
.send(DeletionStatistics::DeleteBytes(55))
.await
.unwrap();
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(received, DeletionStatistics::DeleteBytes(55)));
}
#[tokio::test]
async fn listing_mock_send_stats_delivers_stat() {
let (mock, stats_receiver) = make_listing_mock();
mock.send_stats(DeletionStatistics::DeleteComplete {
key: "k".to_string(),
})
.await;
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(
received,
DeletionStatistics::DeleteComplete { .. }
));
}
#[test]
fn listing_mock_set_warning_sets_flag() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = ListingMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
};
assert!(!has_warning.load(std::sync::atomic::Ordering::SeqCst));
mock.set_warning();
assert!(has_warning.load(std::sync::atomic::Ordering::SeqCst));
}
fn make_partial_failure_mock() -> (
PartialFailureMockStorage,
async_channel::Receiver<DeletionStatistics>,
) {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PartialFailureMockStorage {
objects: vec![],
stats_sender,
has_warning,
fail_in_batch: Arc::new(vec!["fail.txt".to_string()]),
};
(mock, stats_receiver)
}
#[test]
fn partial_failure_mock_is_express_onezone_returns_false() {
let (mock, _) = make_partial_failure_mock();
assert!(!mock.is_express_onezone_storage());
}
#[tokio::test]
async fn partial_failure_mock_list_objects_returns_ok() {
let (mock, _) = make_partial_failure_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_objects(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn partial_failure_mock_list_object_versions_returns_ok() {
let (mock, _) = make_partial_failure_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_object_versions(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn partial_failure_mock_head_object_returns_ok() {
let (mock, _) = make_partial_failure_mock();
assert!(mock.head_object("key", None).await.is_ok());
}
#[tokio::test]
async fn partial_failure_mock_get_object_tagging_returns_ok() {
let (mock, _) = make_partial_failure_mock();
let result = mock.get_object_tagging("key", None).await;
assert!(result.is_ok());
assert!(result.unwrap().tag_set().is_empty());
}
#[tokio::test]
async fn partial_failure_mock_delete_object_returns_ok() {
let (mock, _) = make_partial_failure_mock();
assert!(mock.delete_object("key", None, None).await.is_ok());
}
#[tokio::test]
async fn partial_failure_mock_delete_objects_partial_failure() {
let (mock, _) = make_partial_failure_mock();
let idents = vec![
aws_sdk_s3::types::ObjectIdentifier::builder()
.key("ok.txt")
.build()
.unwrap(),
aws_sdk_s3::types::ObjectIdentifier::builder()
.key("fail.txt")
.build()
.unwrap(),
];
let result = mock.delete_objects(idents).await.unwrap();
assert_eq!(result.deleted().len(), 1);
assert_eq!(result.deleted()[0].key(), Some("ok.txt"));
assert_eq!(result.errors().len(), 1);
assert_eq!(result.errors()[0].key(), Some("fail.txt"));
assert_eq!(result.errors()[0].code(), Some("InternalError"));
}
#[tokio::test]
async fn partial_failure_mock_is_versioning_enabled_returns_false() {
let (mock, _) = make_partial_failure_mock();
assert!(!mock.is_versioning_enabled().await.unwrap());
}
#[test]
fn partial_failure_mock_get_client_returns_none() {
let (mock, _) = make_partial_failure_mock();
assert!(mock.get_client().is_none());
}
#[tokio::test]
async fn partial_failure_mock_get_stats_sender_works() {
let (mock, stats_receiver) = make_partial_failure_mock();
let sender = mock.get_stats_sender();
sender
.send(DeletionStatistics::DeleteBytes(33))
.await
.unwrap();
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(received, DeletionStatistics::DeleteBytes(33)));
}
#[tokio::test]
async fn partial_failure_mock_send_stats_delivers_stat() {
let (mock, stats_receiver) = make_partial_failure_mock();
mock.send_stats(DeletionStatistics::DeleteComplete {
key: "k".to_string(),
})
.await;
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(
received,
DeletionStatistics::DeleteComplete { .. }
));
}
#[test]
fn partial_failure_mock_set_warning_sets_flag() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = PartialFailureMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
fail_in_batch: Arc::new(vec![]),
};
assert!(!has_warning.load(std::sync::atomic::Ordering::SeqCst));
mock.set_warning();
assert!(has_warning.load(std::sync::atomic::Ordering::SeqCst));
}
fn make_failing_lister_mock() -> (
FailingListerStorage,
async_channel::Receiver<DeletionStatistics>,
) {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = FailingListerStorage {
stats_sender,
has_warning,
};
(mock, stats_receiver)
}
#[test]
fn failing_lister_mock_is_express_onezone_returns_false() {
let (mock, _) = make_failing_lister_mock();
assert!(!mock.is_express_onezone_storage());
}
#[tokio::test]
async fn failing_lister_mock_list_objects_returns_err() {
let (mock, _) = make_failing_lister_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
let result = mock.list_objects(&sender, 1000).await;
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("AccessDenied"));
}
#[tokio::test]
async fn failing_lister_mock_list_object_versions_returns_err() {
let (mock, _) = make_failing_lister_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
let result = mock.list_object_versions(&sender, 1000).await;
assert!(result.is_err());
}
#[tokio::test]
async fn failing_lister_mock_head_object_returns_ok() {
let (mock, _) = make_failing_lister_mock();
assert!(mock.head_object("key", None).await.is_ok());
}
#[tokio::test]
async fn failing_lister_mock_get_object_tagging_returns_ok() {
let (mock, _) = make_failing_lister_mock();
let result = mock.get_object_tagging("key", None).await;
assert!(result.is_ok());
assert!(result.unwrap().tag_set().is_empty());
}
#[tokio::test]
async fn failing_lister_mock_delete_object_returns_ok() {
let (mock, _) = make_failing_lister_mock();
assert!(mock.delete_object("key", None, None).await.is_ok());
}
#[tokio::test]
async fn failing_lister_mock_delete_objects_returns_ok() {
let (mock, _) = make_failing_lister_mock();
assert!(mock.delete_objects(vec![]).await.is_ok());
}
#[tokio::test]
async fn failing_lister_mock_is_versioning_enabled_returns_false() {
let (mock, _) = make_failing_lister_mock();
assert!(!mock.is_versioning_enabled().await.unwrap());
}
#[test]
fn failing_lister_mock_get_client_returns_none() {
let (mock, _) = make_failing_lister_mock();
assert!(mock.get_client().is_none());
}
#[tokio::test]
async fn failing_lister_mock_get_stats_sender_works() {
let (mock, stats_receiver) = make_failing_lister_mock();
let sender = mock.get_stats_sender();
sender
.send(DeletionStatistics::DeleteBytes(11))
.await
.unwrap();
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(received, DeletionStatistics::DeleteBytes(11)));
}
#[tokio::test]
async fn failing_lister_mock_send_stats_delivers_stat() {
let (mock, stats_receiver) = make_failing_lister_mock();
mock.send_stats(DeletionStatistics::DeleteComplete {
key: "k".to_string(),
})
.await;
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(
received,
DeletionStatistics::DeleteComplete { .. }
));
}
#[test]
fn failing_lister_mock_set_warning_sets_flag() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = FailingListerStorage {
stats_sender,
has_warning: has_warning.clone(),
};
assert!(!has_warning.load(std::sync::atomic::Ordering::SeqCst));
mock.set_warning();
assert!(has_warning.load(std::sync::atomic::Ordering::SeqCst));
}
fn make_non_retryable_failure_mock() -> (
NonRetryableFailureMockStorage,
async_channel::Receiver<DeletionStatistics>,
) {
let (stats_sender, stats_receiver) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = NonRetryableFailureMockStorage {
objects: vec![],
stats_sender,
has_warning,
fail_in_batch: Arc::new(vec!["denied.txt".to_string()]),
};
(mock, stats_receiver)
}
#[test]
fn non_retryable_failure_mock_is_express_onezone_returns_false() {
let (mock, _) = make_non_retryable_failure_mock();
assert!(!mock.is_express_onezone_storage());
}
#[tokio::test]
async fn non_retryable_failure_mock_list_objects_returns_ok() {
let (mock, _) = make_non_retryable_failure_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_objects(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn non_retryable_failure_mock_list_object_versions_returns_ok() {
let (mock, _) = make_non_retryable_failure_mock();
let (sender, _receiver) = async_channel::bounded::<S3Object>(10);
assert!(mock.list_object_versions(&sender, 1000).await.is_ok());
}
#[tokio::test]
async fn non_retryable_failure_mock_head_object_returns_ok() {
let (mock, _) = make_non_retryable_failure_mock();
assert!(mock.head_object("key", None).await.is_ok());
}
#[tokio::test]
async fn non_retryable_failure_mock_get_object_tagging_returns_ok() {
let (mock, _) = make_non_retryable_failure_mock();
let result = mock.get_object_tagging("key", None).await;
assert!(result.is_ok());
assert!(result.unwrap().tag_set().is_empty());
}
#[tokio::test]
async fn non_retryable_failure_mock_delete_object_returns_ok() {
let (mock, _) = make_non_retryable_failure_mock();
assert!(mock.delete_object("key", None, None).await.is_ok());
}
#[tokio::test]
async fn non_retryable_failure_mock_delete_objects_access_denied() {
let (mock, _) = make_non_retryable_failure_mock();
let idents = vec![
aws_sdk_s3::types::ObjectIdentifier::builder()
.key("ok.txt")
.build()
.unwrap(),
aws_sdk_s3::types::ObjectIdentifier::builder()
.key("denied.txt")
.build()
.unwrap(),
];
let result = mock.delete_objects(idents).await.unwrap();
assert_eq!(result.deleted().len(), 1);
assert_eq!(result.deleted()[0].key(), Some("ok.txt"));
assert_eq!(result.errors().len(), 1);
assert_eq!(result.errors()[0].key(), Some("denied.txt"));
assert_eq!(result.errors()[0].code(), Some("AccessDenied"));
}
#[tokio::test]
async fn non_retryable_failure_mock_is_versioning_enabled_returns_false() {
let (mock, _) = make_non_retryable_failure_mock();
assert!(!mock.is_versioning_enabled().await.unwrap());
}
#[test]
fn non_retryable_failure_mock_get_client_returns_none() {
let (mock, _) = make_non_retryable_failure_mock();
assert!(mock.get_client().is_none());
}
#[tokio::test]
async fn non_retryable_failure_mock_get_stats_sender_works() {
let (mock, stats_receiver) = make_non_retryable_failure_mock();
let sender = mock.get_stats_sender();
sender
.send(DeletionStatistics::DeleteBytes(77))
.await
.unwrap();
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(received, DeletionStatistics::DeleteBytes(77)));
}
#[tokio::test]
async fn non_retryable_failure_mock_send_stats_delivers_stat() {
let (mock, stats_receiver) = make_non_retryable_failure_mock();
mock.send_stats(DeletionStatistics::DeleteComplete {
key: "k".to_string(),
})
.await;
let received = stats_receiver.recv().await.unwrap();
assert!(matches!(
received,
DeletionStatistics::DeleteComplete { .. }
));
}
#[test]
fn non_retryable_failure_mock_set_warning_sets_flag() {
let (stats_sender, _) = async_channel::unbounded();
let has_warning = Arc::new(AtomicBool::new(false));
let mock = NonRetryableFailureMockStorage {
objects: vec![],
stats_sender,
has_warning: has_warning.clone(),
fail_in_batch: Arc::new(vec![]),
};
assert!(!has_warning.load(std::sync::atomic::Ordering::SeqCst));
mock.set_warning();
assert!(has_warning.load(std::sync::atomic::Ordering::SeqCst));
}
}