#![cfg(e2e_test)]
mod common;
use common::{CollectingEventCallback, TestHelper};
use s3rm_rs::{EventType, FilterCallback, S3Object};
use std::io::Write as IoWrite;
use std::sync::{Arc, Mutex};
struct PrefixFilterCallback {
prefix: String,
}
#[async_trait::async_trait]
impl FilterCallback for PrefixFilterCallback {
async fn filter(&mut self, object: &S3Object) -> anyhow::Result<bool> {
Ok(object.key().starts_with(&self.prefix))
}
}
struct SizeFilterCallback {
min_size: i64,
}
#[async_trait::async_trait]
impl FilterCallback for SizeFilterCallback {
async fn filter(&mut self, object: &S3Object) -> anyhow::Result<bool> {
Ok(object.size() >= self.min_size)
}
}
#[tokio::test]
async fn e2e_rust_filter_callback() {
e2e_timeout!(async {
let helper = TestHelper::new().await;
let bucket = helper.generate_bucket_name();
helper.create_bucket(&bucket).await;
let guard = helper.bucket_guard(&bucket);
for i in 0..10 {
helper
.put_object(&bucket, &format!("delete-file{i}.dat"), vec![b'd'; 100])
.await;
}
for i in 0..10 {
helper
.put_object(&bucket, &format!("keep-file{i}.dat"), vec![b'k'; 100])
.await;
}
let mut config = TestHelper::build_config(vec![&format!("s3://{bucket}/"), "--force"]);
config
.filter_manager
.register_callback(PrefixFilterCallback {
prefix: "delete-".to_string(),
});
let result = TestHelper::run_pipeline(config).await;
assert!(!result.has_error, "Pipeline should complete without errors");
assert_eq!(
result.stats.stats_deleted_objects, 10,
"Should delete exactly 10 delete- objects"
);
let remaining = helper.list_objects(&bucket, "keep-").await;
assert_eq!(remaining.len(), 10, "All keep- objects should remain");
let deleted = helper.list_objects(&bucket, "delete-").await;
assert_eq!(deleted.len(), 0, "All delete- objects should be removed");
guard.cleanup().await;
});
}
#[tokio::test]
async fn e2e_rust_event_callback() {
e2e_timeout!(async {
let helper = TestHelper::new().await;
let bucket = helper.generate_bucket_name();
helper.create_bucket(&bucket).await;
let guard = helper.bucket_guard(&bucket);
for i in 0..10 {
helper
.put_object(&bucket, &format!("event/file{i}.dat"), vec![b'e'; 256])
.await;
}
let collected_events = Arc::new(Mutex::new(Vec::new()));
let callback = CollectingEventCallback {
events: Arc::clone(&collected_events),
};
let mut config =
TestHelper::build_config(vec![&format!("s3://{bucket}/event/"), "--force"]);
config
.event_manager
.register_callback(EventType::ALL_EVENTS, callback, false);
let result = TestHelper::run_pipeline(config).await;
assert!(!result.has_error, "Pipeline should complete without errors");
assert_eq!(result.stats.stats_deleted_objects, 10);
{
let events = collected_events.lock().unwrap();
let starts: Vec<_> = events
.iter()
.filter(|e| e.event_type == EventType::PIPELINE_START)
.collect();
assert_eq!(starts.len(), 1, "Should receive exactly 1 PIPELINE_START");
let completes: Vec<_> = events
.iter()
.filter(|e| e.event_type == EventType::DELETE_COMPLETE)
.collect();
assert_eq!(
completes.len(),
10,
"Should receive 10 DELETE_COMPLETE events"
);
let ends: Vec<_> = events
.iter()
.filter(|e| e.event_type == EventType::PIPELINE_END)
.collect();
assert_eq!(ends.len(), 1, "Should receive exactly 1 PIPELINE_END");
}
guard.cleanup().await;
});
}
#[tokio::test]
async fn e2e_lua_filter_callback() {
e2e_timeout!(async {
let helper = TestHelper::new().await;
let bucket = helper.generate_bucket_name();
helper.create_bucket(&bucket).await;
let guard = helper.bucket_guard(&bucket);
for i in 0..10 {
helper
.put_object(&bucket, &format!("lua/file{i}.tmp"), vec![b't'; 100])
.await;
}
for i in 0..10 {
helper
.put_object(&bucket, &format!("lua/file{i}.dat"), vec![b'd'; 100])
.await;
}
let lua_script = tempfile::NamedTempFile::new().unwrap();
writeln!(
lua_script.as_file(),
r#"
function filter(object)
local key = object["key"]
if string.match(key, "%.tmp$") then
return true
end
return false
end
"#
)
.unwrap();
let lua_path = lua_script.path().to_str().unwrap();
let config = TestHelper::build_config(vec![
&format!("s3://{bucket}/lua/"),
"--filter-callback-lua-script",
lua_path,
"--force",
]);
let result = TestHelper::run_pipeline(config).await;
assert!(!result.has_error, "Pipeline should complete without errors");
assert_eq!(
result.stats.stats_deleted_objects, 10,
"Should delete exactly 10 .tmp objects"
);
let remaining_dat = helper.list_objects(&bucket, "lua/").await;
let dat_count = remaining_dat.iter().filter(|k| k.ends_with(".dat")).count();
assert_eq!(dat_count, 10, "All .dat objects should remain");
guard.cleanup().await;
});
}
#[tokio::test]
async fn e2e_lua_event_callback() {
e2e_timeout!(async {
let helper = TestHelper::new().await;
let bucket = helper.generate_bucket_name();
helper.create_bucket(&bucket).await;
let guard = helper.bucket_guard(&bucket);
for i in 0..10 {
helper
.put_object(&bucket, &format!("lua-event/file{i}.dat"), vec![b'e'; 100])
.await;
}
let output_file = tempfile::NamedTempFile::new().unwrap();
let output_path = output_file.path().to_str().unwrap().replace('\\', "/");
let lua_script = tempfile::NamedTempFile::new().unwrap();
writeln!(
lua_script.as_file(),
r#"
local count = 0
function on_event(event_data)
count = count + 1
local f = io.open("{output_path}", "w")
if f then
f:write("events=" .. tostring(count))
f:close()
end
end
"#,
output_path = output_path
)
.unwrap();
let lua_path = lua_script.path().to_str().unwrap();
let config = TestHelper::build_config(vec![
&format!("s3://{bucket}/lua-event/"),
"--event-callback-lua-script",
lua_path,
"--allow-lua-os-library",
"--force",
]);
let result = TestHelper::run_pipeline(config).await;
assert!(!result.has_error, "Pipeline should complete without errors");
assert_eq!(result.stats.stats_deleted_objects, 10);
let content = std::fs::read_to_string(&output_path)
.expect("Lua event output file should exist and be readable");
assert!(
!content.is_empty(),
"Lua event output file should be non-empty (event script should have written to it)"
);
assert!(
content.starts_with("events="),
"Lua event output should contain 'events=<count>', got: {content}"
);
let count_str = content.trim_start_matches("events=");
let count: u64 = count_str
.parse()
.unwrap_or_else(|_| panic!("Failed to parse event count from: {content}"));
assert!(
count > 0,
"Lua event script should have received at least 1 event, got {count}"
);
guard.cleanup().await;
});
}
#[tokio::test]
async fn e2e_lua_sandbox_blocks_os_access() {
e2e_timeout!(async {
let helper = TestHelper::new().await;
let bucket = helper.generate_bucket_name();
helper.create_bucket(&bucket).await;
let guard = helper.bucket_guard(&bucket);
for i in 0..5 {
helper
.put_object(&bucket, &format!("sandbox/file{i}.dat"), vec![b's'; 100])
.await;
}
let lua_script = tempfile::NamedTempFile::new().unwrap();
writeln!(
lua_script.as_file(),
r#"
function filter(object)
os.execute("echo sandbox_escape")
return true
end
"#
)
.unwrap();
let lua_path = lua_script.path().to_str().unwrap();
let config = TestHelper::build_config(vec![
&format!("s3://{bucket}/sandbox/"),
"--filter-callback-lua-script",
lua_path,
"--force",
]);
let result = TestHelper::run_pipeline(config).await;
assert!(
result.has_error,
"Pipeline should report error for Lua sandbox violation"
);
guard.cleanup().await;
});
}
#[tokio::test]
async fn e2e_lua_vm_memory_limit() {
e2e_timeout!(async {
let helper = TestHelper::new().await;
let bucket = helper.generate_bucket_name();
helper.create_bucket(&bucket).await;
let guard = helper.bucket_guard(&bucket);
for i in 0..5 {
helper
.put_object(&bucket, &format!("memlimit/file{i}.dat"), vec![b'm'; 100])
.await;
}
let lua_script = tempfile::NamedTempFile::new().unwrap();
writeln!(
lua_script.as_file(),
r#"
function filter(object)
local big = {{}}
for i = 1, 1000000 do
big[i] = string.rep("x", 100)
end
return true
end
"#
)
.unwrap();
let lua_path = lua_script.path().to_str().unwrap();
let config = TestHelper::build_config(vec![
&format!("s3://{bucket}/memlimit/"),
"--filter-callback-lua-script",
lua_path,
"--lua-vm-memory-limit",
"1MB",
"--force",
]);
let result = TestHelper::run_pipeline(config).await;
assert!(
result.has_error,
"Pipeline should report error for Lua memory limit exceeded"
);
guard.cleanup().await;
});
}
#[tokio::test]
async fn e2e_rust_filter_and_event_callbacks_combined() {
e2e_timeout!(async {
let helper = TestHelper::new().await;
let bucket = helper.generate_bucket_name();
helper.create_bucket(&bucket).await;
let guard = helper.bucket_guard(&bucket);
for i in 0..10 {
helper
.put_object(
&bucket,
&format!("combined/large{i}.dat"),
vec![b'L'; 5 * 1024],
)
.await;
}
for i in 0..10 {
helper
.put_object(&bucket, &format!("combined/small{i}.dat"), vec![b's'; 100])
.await;
}
let collected_events = Arc::new(Mutex::new(Vec::new()));
let event_callback = CollectingEventCallback {
events: Arc::clone(&collected_events),
};
let filter_callback = SizeFilterCallback { min_size: 1024 };
let mut config =
TestHelper::build_config(vec![&format!("s3://{bucket}/combined/"), "--force"]);
config.filter_manager.register_callback(filter_callback);
config
.event_manager
.register_callback(EventType::ALL_EVENTS, event_callback, false);
let result = TestHelper::run_pipeline(config).await;
assert!(!result.has_error, "Pipeline should complete without errors");
assert_eq!(
result.stats.stats_deleted_objects, 10,
"Should delete exactly 10 large objects"
);
let remaining = helper.list_objects(&bucket, "combined/small").await;
assert_eq!(remaining.len(), 10, "All small objects should remain");
{
let events = collected_events.lock().unwrap();
let delete_completes: Vec<_> = events
.iter()
.filter(|e| e.event_type == EventType::DELETE_COMPLETE)
.collect();
assert_eq!(
delete_completes.len(),
10,
"Should receive 10 DELETE_COMPLETE events for the large objects"
);
let delete_filtered: Vec<_> = events
.iter()
.filter(|e| e.event_type == EventType::DELETE_FILTERED)
.collect();
assert_eq!(
delete_filtered.len(),
10,
"Should receive 10 DELETE_FILTERED events for the small objects"
);
}
guard.cleanup().await;
});
}