diff --git a/crates/protocols/src/swift/expiration_worker.rs b/crates/protocols/src/swift/expiration_worker.rs index 164b694c0..7aeb88a86 100644 --- a/crates/protocols/src/swift/expiration_worker.rs +++ b/crates/protocols/src/swift/expiration_worker.rs @@ -54,7 +54,7 @@ use super::storage_api::object::ObjectOperations as _; use super::{SwiftError, SwiftObjectOptions, SwiftResult, resolve_swift_object_store_handle}; use std::cmp::Reverse; use std::collections::{BinaryHeap, HashMap}; -use std::sync::Arc; +use std::sync::{Arc, OnceLock}; use std::time::{Duration, SystemTime, UNIX_EPOCH}; use tokio::sync::RwLock; use tokio::time::interval; @@ -69,19 +69,51 @@ const EVENT_SWIFT_EXPIRATION_DELETE_STATE: &str = "swift_expiration_delete_state const EVENT_SWIFT_EXPIRATION_SCAN_STATE: &str = "swift_expiration_scan_state"; const SWIFT_DELETE_AT_METADATA: &str = "x-delete-at"; +static GLOBAL_EXPIRATION_WORKER: OnceLock> = OnceLock::new(); + +fn global_expiration_worker() -> Arc { + Arc::clone(GLOBAL_EXPIRATION_WORKER.get_or_init(|| Arc::new(ExpirationWorker::new(ExpirationWorkerConfig::default())))) +} + +pub async fn track_object_expiration(account: &str, container: &str, object: &str, expires_at: u64) { + let worker = global_expiration_worker(); + worker.track_object(account, container, object, expires_at).await; + worker.ensure_started().await; +} + +pub async fn untrack_object_expiration(account: &str, container: &str, object: &str) { + if let Some(worker) = GLOBAL_EXPIRATION_WORKER.get() { + worker.untrack_object(account, container, object).await; + } +} + #[async_trait::async_trait] trait ExpirationObjectBackend: Send + Sync { + async fn expiring_objects(&self) -> SwiftResult>; + async fn object_metadata(&self, account: &str, container: &str, object: &str) -> SwiftResult>>; async fn delete_object(&self, account: &str, container: &str, object: &str) -> SwiftResult<()>; } +#[derive(Debug, Clone, Eq, PartialEq)] +struct ExpirationCandidate { + account: String, + container: String, + object: String, + expires_at: u64, +} + #[derive(Debug, Default)] struct SwiftStorageExpirationBackend; #[async_trait::async_trait] impl ExpirationObjectBackend for SwiftStorageExpirationBackend { + async fn expiring_objects(&self) -> SwiftResult> { + Ok(Vec::new()) + } + async fn object_metadata( &self, account: &str, @@ -314,6 +346,14 @@ impl ExpirationWorker { }); } + pub async fn ensure_started(&self) { + if *self.running.read().await { + return; + } + + self.start().await; + } + /// Stop the background worker pub async fn stop(&self) { let mut running = self.running.write().await; @@ -623,6 +663,14 @@ impl ExpirationWorker { /// This is used for initial population or recovery after restart. /// In production, objects should be tracked incrementally via track_object(). pub async fn scan_all_objects(&self) -> SwiftResult<()> { + self.scan_all_objects_with_backend(&SwiftStorageExpirationBackend).await?; + Ok(()) + } + + async fn scan_all_objects_with_backend(&self, backend: &B) -> SwiftResult + where + B: ExpirationObjectBackend + ?Sized, + { info!( event = EVENT_SWIFT_EXPIRATION_SCAN_STATE, component = LOG_COMPONENT_PROTOCOLS, @@ -633,20 +681,33 @@ impl ExpirationWorker { "swift expiration scan state changed" ); - // TODO: This would integrate with the storage layer to list all objects - // For each object with X-Delete-At metadata, call track_object() + let mut tracked_count = 0; + for candidate in backend.expiring_objects().await? { + let path = format!("{}/{}/{}", candidate.account, candidate.container, candidate.object); + if !self.should_handle_object(&path) { + continue; + } - // Placeholder implementation - warn!( + self.track_object(&candidate.account, &candidate.container, &candidate.object, candidate.expires_at) + .await; + tracked_count += 1; + } + + let queue_size = self.priority_queue.read().await.len(); + self.metrics.write().await.queue_size = queue_size; + + info!( event = EVENT_SWIFT_EXPIRATION_SCAN_STATE, component = LOG_COMPONENT_PROTOCOLS, subsystem = LOG_SUBSYSTEM_SWIFT_EXPIRATION, - result = "unimplemented", + state = "completed", worker_id = self.config.worker_id, + tracked_count, + queue_size, "swift expiration scan state changed" ); - Ok(()) + Ok(tracked_count) } } @@ -661,6 +722,7 @@ mod tests { #[derive(Default)] #[allow(clippy::type_complexity)] struct MockExpirationObjectBackend { + expiring_objects: Mutex>, metadata_results: Mutex>, delete_results: Mutex>>, deleted_objects: Mutex>, @@ -674,6 +736,13 @@ mod tests { } } + fn with_expiring_objects(objects: Vec) -> Self { + Self { + expiring_objects: Mutex::new(objects), + ..Default::default() + } + } + fn with_delete_result(self, result: SwiftResult<()>) -> Self { self.delete_results .lock() @@ -692,6 +761,14 @@ mod tests { #[async_trait::async_trait] impl ExpirationObjectBackend for MockExpirationObjectBackend { + async fn expiring_objects(&self) -> SwiftResult> { + let objects = match self.expiring_objects.lock() { + Ok(objects) => objects, + Err(err) => err.into_inner(), + }; + Ok(objects.clone()) + } + async fn object_metadata( &self, _account: &str, @@ -887,6 +964,37 @@ mod tests { assert_eq!(metrics.queue_size, 2); } + #[tokio::test] + async fn test_scan_all_objects_tracks_backend_candidates() { + let worker = ExpirationWorker::new(ExpirationWorkerConfig::default()); + let backend = MockExpirationObjectBackend::with_expiring_objects(vec![ + ExpirationCandidate { + account: "AUTH_test".to_string(), + container: "container".to_string(), + object: "object-a".to_string(), + expires_at: 1000, + }, + ExpirationCandidate { + account: "AUTH_test".to_string(), + container: "container".to_string(), + object: "object-b".to_string(), + expires_at: 2000, + }, + ]); + + let tracked = match worker.scan_all_objects_with_backend(&backend).await { + Ok(tracked) => tracked, + Err(err) => { + assert!(false, "scan candidates should be tracked: {err}"); + 0 + } + }; + + assert_eq!(tracked, 2); + assert_eq!(worker.priority_queue.read().await.len(), 2); + assert_eq!(worker.get_metrics().await.queue_size, 2); + } + #[tokio::test] async fn test_delete_expired_object_deletes_when_expired_metadata_matches() { let expires_at = now_secs().saturating_sub(1); diff --git a/crates/protocols/src/swift/object.rs b/crates/protocols/src/swift/object.rs index f3f592259..66fa370c3 100644 --- a/crates/protocols/src/swift/object.rs +++ b/crates/protocols/src/swift/object.rs @@ -51,6 +51,7 @@ use super::account::validate_account_access; use super::container::ContainerMapper; +use super::expiration_worker::{track_object_expiration, untrack_object_expiration}; use super::storage_api::object::{BucketOperations, BucketOptions, HTTPRangeSpec, ObjectIO as _, ObjectOperations as _}; use super::{SwiftError, SwiftResult, resolve_swift_object_store_handle}; use axum::http::HeaderMap; @@ -65,6 +66,7 @@ pub use super::{SwiftGetObjectReader, SwiftObjectInfo, SwiftObjectOptions, Swift const LOG_COMPONENT_PROTOCOLS: &str = "protocols"; const LOG_SUBSYSTEM_SWIFT_OBJECT: &str = "swift_object"; const EVENT_SWIFT_OBJECT_STORAGE_STATE: &str = "swift_object_storage_state"; +const SWIFT_DELETE_AT_METADATA: &str = "x-delete-at"; /// Maximum number of metadata headers allowed per object (Swift standard) const MAX_METADATA_COUNT: usize = 90; @@ -264,6 +266,20 @@ fn validate_metadata(metadata: &HashMap) -> SwiftResult<()> { Ok(()) } +fn metadata_delete_at(metadata: &HashMap) -> Option { + metadata + .get(SWIFT_DELETE_AT_METADATA) + .and_then(|value| value.parse::().ok()) +} + +async fn update_object_expiration_tracking(account: &str, container: &str, object: &str, delete_at: Option) { + if let Some(delete_at) = delete_at { + track_object_expiration(account, container, object, delete_at).await; + } else { + untrack_object_expiration(account, container, object).await; + } +} + /// Sanitize storage layer errors for client responses /// /// Logs detailed error server-side while returning generic message to client. @@ -344,9 +360,10 @@ where } // 7. Extract and validate expiration headers (X-Delete-At / X-Delete-After) - if let Some(delete_at) = super::expiration::extract_expiration(headers)? { + let delete_at = super::expiration::extract_expiration(headers)?; + if let Some(delete_at) = delete_at { super::expiration::validate_expiration(delete_at)?; - user_metadata.insert("x-delete-at".to_string(), delete_at.to_string()); + user_metadata.insert(SWIFT_DELETE_AT_METADATA.to_string(), delete_at.to_string()); } // 8. Extract symlink target if creating a symlink @@ -418,6 +435,8 @@ where .await .map_err(|e| sanitize_storage_error("Object upload", e))?; + update_object_expiration_tracking(account, container, object, delete_at).await; + // 17. Return ETag (MD5 hash in hex format) Ok(obj_info.etag.unwrap_or_default()) } @@ -461,6 +480,7 @@ where // Validate metadata limits validate_metadata(metadata)?; + let delete_at = metadata_delete_at(metadata); // Get storage layer let Some(store) = resolve_swift_object_store_handle() else { @@ -501,6 +521,8 @@ where .await .map_err(|e| sanitize_storage_error("Object upload", e))?; + update_object_expiration_tracking(account, container, object, delete_at).await; + // Return ETag Ok(obj_info.etag.unwrap_or_default()) } @@ -675,7 +697,10 @@ pub async fn delete_object(account: &str, container: &str, object: &str, credent // 7. Delete object from storage // Swift DELETE is idempotent - returns success even if object doesn't exist match store.delete_object(&bucket, &s3_key, opts).await { - Ok(_) => Ok(()), + Ok(_) => { + untrack_object_expiration(account, container, object).await; + Ok(()) + } Err(e) => { let err_str = e.to_string(); // Only fail if the container (bucket) doesn't exist @@ -683,6 +708,7 @@ pub async fn delete_object(account: &str, container: &str, object: &str, credent Err(SwiftError::NotFound(format!("Container '{}' not found", container))) } else if err_str.contains("Object not found") || err_str.contains("does not exist") { // Object already gone - this is success for idempotent DELETE + untrack_object_expiration(account, container, object).await; Ok(()) } else { Err(sanitize_storage_error("Object deletion", e)) @@ -769,6 +795,12 @@ pub async fn update_object_metadata( new_metadata.insert("content-type".to_string(), ct_str.to_string()); } + let delete_at = super::expiration::extract_expiration(headers)?; + if let Some(delete_at) = delete_at { + super::expiration::validate_expiration(delete_at)?; + new_metadata.insert(SWIFT_DELETE_AT_METADATA.to_string(), delete_at.to_string()); + } + // 10. Validate metadata limits validate_metadata(&new_metadata)?; @@ -787,6 +819,8 @@ pub async fn update_object_metadata( .await .map_err(|e| sanitize_storage_error("Metadata update", e))?; + update_object_expiration_tracking(account, container, object, delete_at).await; + Ok(()) } @@ -1201,6 +1235,22 @@ mod tests { assert!(ObjectKeyMapper::validate_object_name("my..file.txt").is_ok()); } + #[test] + fn test_metadata_delete_at_parses_valid_timestamp() { + let mut metadata = HashMap::new(); + metadata.insert(SWIFT_DELETE_AT_METADATA.to_string(), "1740000000".to_string()); + + assert_eq!(metadata_delete_at(&metadata), Some(1740000000)); + } + + #[test] + fn test_metadata_delete_at_ignores_invalid_timestamp() { + let mut metadata = HashMap::new(); + metadata.insert(SWIFT_DELETE_AT_METADATA.to_string(), "invalid".to_string()); + + assert_eq!(metadata_delete_at(&metadata), None); + } + #[test] fn test_swift_to_s3_key() { assert_eq!(ObjectKeyMapper::swift_to_s3_key("file.txt").unwrap(), "file.txt");