diff --git a/Cargo.lock b/Cargo.lock index ec017c28e..dc355d7cb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8129,6 +8129,7 @@ dependencies = [ "serde_json", "starshard", "thiserror 2.0.18", + "time", "tokio", "tracing", "tracing-subscriber", diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 18745dbbc..22c08b544 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -817,12 +817,17 @@ pub async fn expire_transitioned_object( //defer auditLogLifecycle(ctx, *oi, ILMExpiry, tags, traceFn) - let mut event_name = EventName::ObjectRemovedDelete; - if oi.delete_marker { - event_name = EventName::ObjectRemovedDeleteMarkerCreated; - } + let event_name = if oi.delete_marker { + EventName::LifecycleExpirationDelete + } else if dobj.delete_marker { + EventName::LifecycleExpirationDeleteMarkerCreated + } else { + EventName::LifecycleExpirationDelete + }; let obj_info = ObjectInfo { + bucket: oi.bucket.clone(), name: oi.name.clone(), + size: oi.size, version_id: oi.version_id, delete_marker: oi.delete_marker, ..Default::default() @@ -1230,15 +1235,12 @@ pub async fn apply_expiry_on_non_transitioned_objects( //let tags = LcAuditEvent::new(lc_event.clone(), src.clone()).tags(); //tags["version-id"] = dobj.version_id; - let mut event_name = EventName::ObjectRemovedDelete; - if oi.delete_marker { - event_name = EventName::ObjectRemovedDeleteMarkerCreated; - } - match lc_event.action { - IlmAction::DeleteAllVersionsAction => event_name = EventName::ObjectRemovedDeleteAllVersions, - IlmAction::DelMarkerDeleteAllVersionsAction => event_name = EventName::LifecycleDelMarkerExpirationDelete, - _ => (), - } + let event_name = match lc_event.action { + IlmAction::DeleteAllVersionsAction | IlmAction::DelMarkerDeleteAllVersionsAction => EventName::LifecycleExpirationDelete, + _ if oi.delete_marker => EventName::LifecycleExpirationDelete, + _ if dobj.delete_marker => EventName::LifecycleExpirationDeleteMarkerCreated, + _ => EventName::LifecycleExpirationDelete, + }; send_event(EventArgs { event_name: event_name.to_string(), bucket_name: dobj.bucket.clone(), diff --git a/crates/ecstore/src/event/name.rs b/crates/ecstore/src/event/name.rs index 71075b036..43da8bc75 100644 --- a/crates/ecstore/src/event/name.rs +++ b/crates/ecstore/src/event/name.rs @@ -1,4 +1,3 @@ -#![allow(unused_variables)] // Copyright 2024 RustFS Team // // Licensed under the Apache License, Version 2.0 (the "License"); @@ -13,274 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -//! Defines the EventName enum which represents the various S3 event types that can trigger notifications. -//! This enum includes both specific event types (e.g., ObjectCreated:Put) and aggregate types (e.g., ObjectCreated:*). Each variant has methods to expand into its constituent event types and to compute a bitmask for efficient filtering. -//! The EventName enum is used in the event notification system to determine which events should trigger notifications based on the configured rules. -//! -//! @Deprecated: This module is currently not fully implemented and serves as a placeholder for future development of the event notification system. The EventName enum and its associated methods are defined, but the actual logic for handling events and sending notifications is not yet implemented. +//! Compatibility re-export for the legacy `rustfs_ecstore::event::name::EventName` path. +//! The canonical event definition now lives in `rustfs_s3_common::EventName`. -#[derive(Default, Clone)] -pub enum EventName { - ObjectAccessedGet, - ObjectAccessedGetRetention, - ObjectAccessedGetLegalHold, - ObjectAccessedHead, - ObjectAccessedAttributes, - ObjectCreatedCompleteMultipartUpload, - ObjectCreatedCopy, - ObjectCreatedPost, - ObjectCreatedPut, - ObjectCreatedPutRetention, - ObjectCreatedPutLegalHold, - ObjectCreatedPutTagging, - ObjectCreatedDeleteTagging, - ObjectRemovedDelete, - ObjectRemovedDeleteMarkerCreated, - ObjectRemovedDeleteAllVersions, - ObjectRemovedNoOP, - BucketCreated, - BucketRemoved, - ObjectReplicationFailed, - ObjectReplicationComplete, - ObjectReplicationMissedThreshold, - ObjectReplicationReplicatedAfterThreshold, - ObjectReplicationNotTracked, - ObjectRestorePost, - ObjectRestoreCompleted, - ObjectTransitionFailed, - ObjectTransitionComplete, - ObjectManyVersions, - ObjectLargeVersions, - PrefixManyFolders, - ILMDelMarkerExpirationDelete, - ObjectSingleTypesEnd, - ObjectAccessedAll, - ObjectCreatedAll, - ObjectRemovedAll, - ObjectReplicationAll, - ObjectRestoreAll, - ObjectTransitionAll, - ObjectScannerAll, - #[default] - Everything, -} - -impl EventName { - fn expand(&self) -> Vec { - match self.clone() { - EventName::Everything => vec![ - EventName::BucketCreated, - EventName::BucketRemoved, - EventName::ObjectAccessedAll, - EventName::ObjectCreatedAll, - EventName::ObjectRemovedAll, - EventName::ObjectManyVersions, - EventName::ObjectLargeVersions, - EventName::PrefixManyFolders, - EventName::ILMDelMarkerExpirationDelete, - EventName::ObjectReplicationAll, - EventName::ObjectRestoreAll, - EventName::ObjectTransitionAll, - ], - EventName::ObjectAccessedAll => vec![ - EventName::ObjectAccessedGet, - EventName::ObjectAccessedGetRetention, - EventName::ObjectAccessedGetLegalHold, - EventName::ObjectAccessedHead, - EventName::ObjectAccessedAttributes, - ], - EventName::ObjectCreatedAll => vec![ - EventName::ObjectCreatedCompleteMultipartUpload, - EventName::ObjectCreatedCopy, - EventName::ObjectCreatedPost, - EventName::ObjectCreatedPut, - EventName::ObjectCreatedPutRetention, - EventName::ObjectCreatedPutLegalHold, - EventName::ObjectCreatedPutTagging, - EventName::ObjectCreatedDeleteTagging, - ], - EventName::ObjectRemovedAll => vec![ - EventName::ObjectRemovedDelete, - EventName::ObjectRemovedDeleteMarkerCreated, - EventName::ObjectRemovedNoOP, - EventName::ObjectRemovedDeleteAllVersions, - ], - EventName::ObjectReplicationAll => vec![ - EventName::ObjectReplicationFailed, - EventName::ObjectReplicationComplete, - EventName::ObjectReplicationNotTracked, - EventName::ObjectReplicationMissedThreshold, - EventName::ObjectReplicationReplicatedAfterThreshold, - ], - EventName::ObjectRestoreAll => vec![EventName::ObjectRestorePost, EventName::ObjectRestoreCompleted], - EventName::ObjectTransitionAll => vec![EventName::ObjectTransitionFailed, EventName::ObjectTransitionComplete], - EventName::ObjectSingleTypesEnd | EventName::ObjectScannerAll => vec![self.clone()], - _ => vec![self.clone()], - } - } - - fn mask(&self) -> u64 { - match self { - EventName::Everything => u64::MAX, - EventName::BucketCreated => 1_u64 << 0, - EventName::BucketRemoved => 1_u64 << 1, - EventName::ObjectAccessedGet => 1_u64 << 2, - EventName::ObjectAccessedGetRetention => 1_u64 << 3, - EventName::ObjectAccessedGetLegalHold => 1_u64 << 4, - EventName::ObjectAccessedHead => 1_u64 << 5, - EventName::ObjectAccessedAttributes => 1_u64 << 6, - EventName::ObjectCreatedCompleteMultipartUpload => 1_u64 << 7, - EventName::ObjectCreatedCopy => 1_u64 << 8, - EventName::ObjectCreatedPost => 1_u64 << 9, - EventName::ObjectCreatedPut => 1_u64 << 10, - EventName::ObjectCreatedPutRetention => 1_u64 << 11, - EventName::ObjectCreatedPutLegalHold => 1_u64 << 12, - EventName::ObjectCreatedPutTagging => 1_u64 << 13, - EventName::ObjectCreatedDeleteTagging => 1_u64 << 14, - EventName::ObjectRemovedDelete => 1_u64 << 15, - EventName::ObjectRemovedDeleteMarkerCreated => 1_u64 << 16, - EventName::ObjectRemovedDeleteAllVersions => 1_u64 << 17, - EventName::ObjectRemovedNoOP => 1_u64 << 18, - EventName::ObjectManyVersions => 1_u64 << 19, - EventName::ObjectLargeVersions => 1_u64 << 20, - EventName::PrefixManyFolders => 1_u64 << 21, - EventName::ILMDelMarkerExpirationDelete => 1_u64 << 22, - EventName::ObjectReplicationFailed => 1_u64 << 23, - EventName::ObjectReplicationComplete => 1_u64 << 24, - EventName::ObjectReplicationMissedThreshold => 1_u64 << 25, - EventName::ObjectReplicationReplicatedAfterThreshold => 1_u64 << 26, - EventName::ObjectReplicationNotTracked => 1_u64 << 27, - EventName::ObjectRestorePost => 1_u64 << 28, - EventName::ObjectRestoreCompleted => 1_u64 << 29, - EventName::ObjectRestoreAll => 1_u64 << 30, - EventName::ObjectTransitionFailed => 1_u64 << 31, - EventName::ObjectTransitionComplete => 1_u64 << 32, - EventName::ObjectAccessedAll => { - EventName::ObjectAccessedGet.mask() - | EventName::ObjectAccessedGetRetention.mask() - | EventName::ObjectAccessedGetLegalHold.mask() - | EventName::ObjectAccessedHead.mask() - | EventName::ObjectAccessedAttributes.mask() - } - EventName::ObjectCreatedAll => { - EventName::ObjectCreatedCompleteMultipartUpload.mask() - | EventName::ObjectCreatedCopy.mask() - | EventName::ObjectCreatedPost.mask() - | EventName::ObjectCreatedPut.mask() - | EventName::ObjectCreatedPutRetention.mask() - | EventName::ObjectCreatedPutLegalHold.mask() - | EventName::ObjectCreatedPutTagging.mask() - | EventName::ObjectCreatedDeleteTagging.mask() - } - EventName::ObjectRemovedAll => { - EventName::ObjectRemovedDelete.mask() - | EventName::ObjectRemovedDeleteMarkerCreated.mask() - | EventName::ObjectRemovedNoOP.mask() - | EventName::ObjectRemovedDeleteAllVersions.mask() - } - EventName::ObjectReplicationAll => { - EventName::ObjectReplicationFailed.mask() - | EventName::ObjectReplicationComplete.mask() - | EventName::ObjectReplicationMissedThreshold.mask() - | EventName::ObjectReplicationReplicatedAfterThreshold.mask() - | EventName::ObjectReplicationNotTracked.mask() - } - EventName::ObjectTransitionAll => { - EventName::ObjectTransitionFailed.mask() | EventName::ObjectTransitionComplete.mask() - } - EventName::ObjectSingleTypesEnd | EventName::ObjectScannerAll => 0, - } - } -} - -impl AsRef for EventName { - fn as_ref(&self) -> &str { - match self { - EventName::BucketCreated => "s3:BucketCreated:*", - EventName::BucketRemoved => "s3:BucketRemoved:*", - EventName::ObjectAccessedAll => "s3:ObjectAccessed:*", - EventName::ObjectAccessedGet => "s3:ObjectAccessed:Get", - EventName::ObjectAccessedGetRetention => "s3:ObjectAccessed:GetRetention", - EventName::ObjectAccessedGetLegalHold => "s3:ObjectAccessed:GetLegalHold", - EventName::ObjectAccessedHead => "s3:ObjectAccessed:Head", - EventName::ObjectAccessedAttributes => "s3:ObjectAccessed:Attributes", - EventName::ObjectCreatedAll => "s3:ObjectCreated:*", - EventName::ObjectCreatedCompleteMultipartUpload => "s3:ObjectCreated:CompleteMultipartUpload", - EventName::ObjectCreatedCopy => "s3:ObjectCreated:Copy", - EventName::ObjectCreatedPost => "s3:ObjectCreated:Post", - EventName::ObjectCreatedPut => "s3:ObjectCreated:Put", - EventName::ObjectCreatedPutTagging => "s3:ObjectCreated:PutTagging", - EventName::ObjectCreatedDeleteTagging => "s3:ObjectCreated:DeleteTagging", - EventName::ObjectCreatedPutRetention => "s3:ObjectCreated:PutRetention", - EventName::ObjectCreatedPutLegalHold => "s3:ObjectCreated:PutLegalHold", - EventName::ObjectRemovedAll => "s3:ObjectRemoved:*", - EventName::ObjectRemovedDelete => "s3:ObjectRemoved:Delete", - EventName::ObjectRemovedDeleteMarkerCreated => "s3:ObjectRemoved:DeleteMarkerCreated", - EventName::ObjectRemovedNoOP => "s3:ObjectRemoved:NoOP", - EventName::ObjectRemovedDeleteAllVersions => "s3:ObjectRemoved:DeleteAllVersions", - EventName::ILMDelMarkerExpirationDelete => "s3:LifecycleDelMarkerExpiration:Delete", - EventName::ObjectReplicationAll => "s3:Replication:*", - EventName::ObjectReplicationFailed => "s3:Replication:OperationFailedReplication", - EventName::ObjectReplicationComplete => "s3:Replication:OperationCompletedReplication", - EventName::ObjectReplicationNotTracked => "s3:Replication:OperationNotTracked", - EventName::ObjectReplicationMissedThreshold => "s3:Replication:OperationMissedThreshold", - EventName::ObjectReplicationReplicatedAfterThreshold => "s3:Replication:OperationReplicatedAfterThreshold", - EventName::ObjectRestoreAll => "s3:ObjectRestore:*", - EventName::ObjectRestorePost => "s3:ObjectRestore:Post", - EventName::ObjectRestoreCompleted => "s3:ObjectRestore:Completed", - EventName::ObjectTransitionAll => "s3:ObjectTransition:*", - EventName::ObjectTransitionFailed => "s3:ObjectTransition:Failed", - EventName::ObjectTransitionComplete => "s3:ObjectTransition:Complete", - EventName::ObjectManyVersions => "s3:Scanner:ManyVersions", - EventName::ObjectLargeVersions => "s3:Scanner:LargeVersions", - EventName::PrefixManyFolders => "s3:Scanner:BigPrefix", - _ => "", - } - } -} - -impl From<&str> for EventName { - fn from(s: &str) -> Self { - match s { - "s3:BucketCreated:*" => EventName::BucketCreated, - "s3:BucketRemoved:*" => EventName::BucketRemoved, - "s3:ObjectAccessed:*" => EventName::ObjectAccessedAll, - "s3:ObjectAccessed:Get" => EventName::ObjectAccessedGet, - "s3:ObjectAccessed:GetRetention" => EventName::ObjectAccessedGetRetention, - "s3:ObjectAccessed:GetLegalHold" => EventName::ObjectAccessedGetLegalHold, - "s3:ObjectAccessed:Head" => EventName::ObjectAccessedHead, - "s3:ObjectAccessed:Attributes" => EventName::ObjectAccessedAttributes, - "s3:ObjectCreated:*" => EventName::ObjectCreatedAll, - "s3:ObjectCreated:CompleteMultipartUpload" => EventName::ObjectCreatedCompleteMultipartUpload, - "s3:ObjectCreated:Copy" => EventName::ObjectCreatedCopy, - "s3:ObjectCreated:Post" => EventName::ObjectCreatedPost, - "s3:ObjectCreated:Put" => EventName::ObjectCreatedPut, - "s3:ObjectCreated:PutRetention" => EventName::ObjectCreatedPutRetention, - "s3:ObjectCreated:PutLegalHold" => EventName::ObjectCreatedPutLegalHold, - "s3:ObjectCreated:PutTagging" => EventName::ObjectCreatedPutTagging, - "s3:ObjectCreated:DeleteTagging" => EventName::ObjectCreatedDeleteTagging, - "s3:ObjectRemoved:*" => EventName::ObjectRemovedAll, - "s3:ObjectRemoved:Delete" => EventName::ObjectRemovedDelete, - "s3:ObjectRemoved:DeleteMarkerCreated" => EventName::ObjectRemovedDeleteMarkerCreated, - "s3:ObjectRemoved:NoOP" => EventName::ObjectRemovedNoOP, - "s3:ObjectRemoved:DeleteAllVersions" => EventName::ObjectRemovedDeleteAllVersions, - "s3:LifecycleDelMarkerExpiration:Delete" => EventName::ILMDelMarkerExpirationDelete, - "s3:Replication:*" => EventName::ObjectReplicationAll, - "s3:Replication:OperationFailedReplication" => EventName::ObjectReplicationFailed, - "s3:Replication:OperationCompletedReplication" => EventName::ObjectReplicationComplete, - "s3:Replication:OperationMissedThreshold" => EventName::ObjectReplicationMissedThreshold, - "s3:Replication:OperationReplicatedAfterThreshold" => EventName::ObjectReplicationReplicatedAfterThreshold, - "s3:Replication:OperationNotTracked" => EventName::ObjectReplicationNotTracked, - "s3:ObjectRestore:*" => EventName::ObjectRestoreAll, - "s3:ObjectRestore:Post" => EventName::ObjectRestorePost, - "s3:ObjectRestore:Completed" => EventName::ObjectRestoreCompleted, - "s3:ObjectTransition:Failed" => EventName::ObjectTransitionFailed, - "s3:ObjectTransition:Complete" => EventName::ObjectTransitionComplete, - "s3:ObjectTransition:*" => EventName::ObjectTransitionAll, - "s3:Scanner:ManyVersions" => EventName::ObjectManyVersions, - "s3:Scanner:LargeVersions" => EventName::ObjectLargeVersions, - "s3:Scanner:BigPrefix" => EventName::PrefixManyFolders, - _ => EventName::Everything, - } - } -} +pub use rustfs_s3_common::EventName; diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index dbfad0789..5dda520dc 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -1882,12 +1882,19 @@ impl ObjectOperations for SetDisks { fi.transitioned_objname = dest_obj; fi.transition_tier = opts.transition.tier.clone(); fi.transition_version_id = if rv.is_empty() { None } else { Some(Uuid::parse_str(&rv)?) }; - let mut event_name = EventName::ObjectTransitionComplete.as_str(); + let event_name = EventName::LifecycleTransition.as_str(); + let mut should_notify_transition = true; let disks = self.get_disks(0, 0).await?; if let Err(err) = self.delete_object_version(bucket, object, &fi, false).await { - event_name = EventName::ObjectTransitionFailed.as_str(); + should_notify_transition = false; + warn!( + bucket = bucket, + object = object, + error = ?err, + "transition completed on remote tier but source cleanup failed; skipping external lifecycle transition notification" + ); } for disk in disks.iter() { @@ -1898,15 +1905,17 @@ impl ObjectOperations for SetDisks { break; } - let obj_info = ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended); - send_event(EventArgs { - event_name: event_name.to_string(), - bucket_name: bucket.to_string(), - object: obj_info, - user_agent: "Internal: [ILM-Transition]".to_string(), - host: GLOBAL_LocalNodeName.to_string(), - ..Default::default() - }); + if should_notify_transition { + let obj_info = ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended); + send_event(EventArgs { + event_name: event_name.to_string(), + bucket_name: bucket.to_string(), + object: obj_info, + user_agent: "Internal: [ILM-Transition]".to_string(), + host: GLOBAL_LocalNodeName.to_string(), + ..Default::default() + }); + } //let tags = opts.lifecycle_audit_event.tags(); //auditLogLifecycle(ctx, objInfo, ILMTransition, tags, traceFn) Ok(()) @@ -1961,10 +1970,19 @@ impl ObjectOperations for SetDisks { false, )?; let mut p_reader = PutObjReader::new(hash_reader); - return if let Err(err) = self_.clone().put_object(bucket, object, &mut p_reader, &ropts).await { - set_restore_header_fn(&mut oi, Some(to_object_err(err, vec![bucket, object]))).await - } else { - Ok(()) + return match self_.clone().put_object(bucket, object, &mut p_reader, &ropts).await { + Ok(restored_info) => { + send_event(EventArgs { + event_name: EventName::ObjectRestoreCompleted.as_str().to_string(), + bucket_name: bucket.to_string(), + object: restored_info, + user_agent: "Internal: [Restore-Completed]".to_string(), + host: GLOBAL_LocalNodeName.to_string(), + ..Default::default() + }); + Ok(()) + } + Err(err) => set_restore_header_fn(&mut oi, Some(to_object_err(err, vec![bucket, object]))).await, }; } @@ -2055,7 +2073,7 @@ impl ObjectOperations for SetDisks { checksum_crc64nvme: None, }); } - if let Err(err) = self_ + let restored_info = match self_ .clone() .complete_multipart_upload( bucket, @@ -2069,8 +2087,17 @@ impl ObjectOperations for SetDisks { ) .await { - return set_restore_header_fn(&mut oi, Some(err)).await; - } + Ok(info) => info, + Err(err) => return set_restore_header_fn(&mut oi, Some(err)).await, + }; + send_event(EventArgs { + event_name: EventName::ObjectRestoreCompleted.as_str().to_string(), + bucket_name: bucket.to_string(), + object: restored_info, + user_agent: "Internal: [Restore-Completed]".to_string(), + host: GLOBAL_LocalNodeName.to_string(), + ..Default::default() + }); Ok(()) } diff --git a/crates/notify/Cargo.toml b/crates/notify/Cargo.toml index df929bf60..ce13f6a81 100644 --- a/crates/notify/Cargo.toml +++ b/crates/notify/Cargo.toml @@ -61,6 +61,7 @@ tracing-subscriber = { workspace = true, features = ["env-filter"] } axum = { workspace = true } rustfs-utils = { workspace = true, features = ["path", "sys"] } serde_json = { workspace = true } +time = { workspace = true } [lints] workspace = true diff --git a/crates/notify/src/event.rs b/crates/notify/src/event.rs index 567856bdf..0cf4662ee 100644 --- a/crates/notify/src/event.rs +++ b/crates/notify/src/event.rs @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -use chrono::{DateTime, Utc}; +use chrono::{DateTime, SecondsFormat, Utc}; use hashbrown::HashMap; use rustfs_s3_common::EventName; use serde::{Deserialize, Serialize}; @@ -90,6 +90,21 @@ pub struct Source { pub user_agent: String, } +/// Additional data included for restore-completed events. +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct GlacierEventData { + pub restore_event_data: RestoreEventData, +} + +/// Restore-specific event attributes for `s3:ObjectRestore:Completed`. +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct RestoreEventData { + pub lifecycle_restoration_expiry_time: String, + pub lifecycle_restore_storage_class: String, +} + /// Represents a storage event #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] @@ -112,11 +127,33 @@ pub struct Event { pub response_elements: HashMap, /// Metadata about the event pub s3: Metadata, + /// Additional restore event data when present. + #[serde(skip_serializing_if = "Option::is_none")] + pub glacier_event_data: Option, /// Information about the source of the event pub source: Source, } impl Event { + fn event_version_for(event_name: EventName) -> &'static str { + match event_name { + EventName::ObjectReplicationFailed + | EventName::ObjectReplicationComplete + | EventName::ObjectReplicationMissedThreshold + | EventName::ObjectReplicationReplicatedAfterThreshold + | EventName::ObjectReplicationNotTracked => "2.2", + EventName::ObjectRestoreCompleted + | EventName::ObjectAclPut + | EventName::ObjectTaggingPut + | EventName::ObjectTaggingDelete + | EventName::LifecycleExpirationDelete + | EventName::LifecycleExpirationDeleteMarkerCreated + | EventName::LifecycleTransition + | EventName::IntelligentTiering => "2.3", + _ => "2.1", + } + } + /// Creates a test event for a given bucket and object pub fn new_test_event(bucket: &str, key: &str, event_name: EventName) -> Self { let mut user_metadata = HashMap::new(); @@ -139,7 +176,7 @@ impl Event { user_metadata.insert("x-request-time".to_string(), Utc::now().to_rfc3339()); Event { - event_version: "2.1".to_string(), + event_version: Self::event_version_for(event_name).to_string(), event_source: "rustfs:s3".to_string(), aws_region: "us-east-1".to_string(), event_time: Utc::now(), @@ -169,6 +206,7 @@ impl Event { sequencer: "0055AED6DCD90281E5".to_string(), }, }, + glacier_event_data: None, source: Source { host: "127.0.0.1".to_string(), port: "9000".to_string(), @@ -237,8 +275,25 @@ impl Event { s3_metadata.object.user_metadata = Some(user_metadata); } + let glacier_event_data = if args.event_name == EventName::ObjectRestoreCompleted { + args.object.restore_expires.and_then(|expiry| { + let expiry_time = DateTime::::from_timestamp(expiry.unix_timestamp(), expiry.nanosecond())?; + let storage_class = args.object.storage_class.clone().or_else(|| { + (!args.object.transitioned_object.tier.is_empty()).then_some(args.object.transitioned_object.tier.clone()) + })?; + Some(GlacierEventData { + restore_event_data: RestoreEventData { + lifecycle_restoration_expiry_time: expiry_time.to_rfc3339_opts(SecondsFormat::Millis, true), + lifecycle_restore_storage_class: storage_class, + }, + }) + }) + } else { + None + }; + Self { - event_version: "2.1".to_string(), + event_version: Self::event_version_for(args.event_name).to_string(), event_source: "rustfs:s3".to_string(), aws_region: args.req_params.get("region").cloned().unwrap_or_default(), event_time: event_time.and_utc(), @@ -247,6 +302,7 @@ impl Event { request_parameters: args.req_params, response_elements: resp_elements, s3: s3_metadata, + glacier_event_data, source: Source { host: args.host, port: if args.port == 0 { @@ -410,3 +466,61 @@ impl EventArgsBuilder { } } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn new_test_event_uses_aws_compatible_event_versions() { + let acl_event = Event::new_test_event("bucket", "key", EventName::ObjectAclPut); + assert_eq!(acl_event.event_version, "2.3"); + + let tagging_event = Event::new_test_event("bucket", "key", EventName::ObjectTaggingPut); + assert_eq!(tagging_event.event_version, "2.3"); + + let lifecycle_event = Event::new_test_event("bucket", "key", EventName::LifecycleExpirationDelete); + assert_eq!(lifecycle_event.event_version, "2.3"); + + let put_event = Event::new_test_event("bucket", "key", EventName::ObjectCreatedPut); + assert_eq!(put_event.event_version, "2.1"); + } + + #[test] + fn event_new_uses_aws_compatible_event_versions() { + let args = EventArgsBuilder::new( + EventName::LifecycleTransition, + "bucket", + rustfs_ecstore::store_api::ObjectInfo { + bucket: "bucket".to_string(), + name: "key".to_string(), + ..Default::default() + }, + ) + .build(); + let event = Event::new(args); + assert_eq!(event.event_version, "2.3"); + } + + #[test] + fn object_restore_completed_includes_glacier_event_data() { + let args = EventArgsBuilder::new( + EventName::ObjectRestoreCompleted, + "bucket", + rustfs_ecstore::store_api::ObjectInfo { + bucket: "bucket".to_string(), + name: "key".to_string(), + restore_expires: Some(time::OffsetDateTime::from_unix_timestamp(1_700_000_000).unwrap()), + storage_class: Some("GLACIER".to_string()), + ..Default::default() + }, + ) + .build(); + let event = Event::new(args); + + assert_eq!(event.event_version, "2.3"); + let glacier = event.glacier_event_data.expect("glacier event data should be present"); + assert_eq!(glacier.restore_event_data.lifecycle_restoration_expiry_time, "2023-11-14T22:13:20.000Z"); + assert_eq!(glacier.restore_event_data.lifecycle_restore_storage_class, "GLACIER"); + } +} diff --git a/crates/notify/src/rules/config_test.rs b/crates/notify/src/rules/config_test.rs index 801b60e2b..a2d2dbef1 100644 --- a/crates/notify/src/rules/config_test.rs +++ b/crates/notify/src/rules/config_test.rs @@ -417,16 +417,12 @@ mod integration_tests { let rules_map = config.get_rules_map(); - // ObjectCreated:* should be expanded to all ObjectCreated events + // AWS ObjectCreated:* should only include object creation operations. let event_types = [ EventName::ObjectCreatedPut, EventName::ObjectCreatedPost, EventName::ObjectCreatedCopy, EventName::ObjectCreatedCompleteMultipartUpload, - EventName::ObjectCreatedPutRetention, - EventName::ObjectCreatedPutLegalHold, - EventName::ObjectCreatedPutTagging, - EventName::ObjectCreatedDeleteTagging, ]; for event_type in event_types { @@ -436,5 +432,8 @@ mod integration_tests { let targets = rules_map.match_rules(event_type, "data/file.csv"); assert!(!targets.is_empty(), "Event {:?} should match", event_type); } + + assert!(!rules_map.has_subscriber(&EventName::ObjectTaggingPut)); + assert!(!rules_map.has_subscriber(&EventName::ObjectTaggingDelete)); } } diff --git a/crates/s3-common/src/event_name.rs b/crates/s3-common/src/event_name.rs index 99c436e7d..6b22a656e 100644 --- a/crates/s3-common/src/event_name.rs +++ b/crates/s3-common/src/event_name.rs @@ -30,7 +30,7 @@ impl std::error::Error for ParseEventNameError {} /// Based on AWS S3 event type and includes RustFS extension. #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)] pub enum EventName { - // Single event type (values are 1-32 for compatible mask logic) + // Single event type (values are sequential for compatible mask logic) ObjectAccessedGet = 1, ObjectAccessedGetRetention = 2, ObjectAccessedGetLegalHold = 3, @@ -42,8 +42,8 @@ pub enum EventName { ObjectCreatedPut = 9, ObjectCreatedPutRetention = 10, ObjectCreatedPutLegalHold = 11, - ObjectCreatedPutTagging = 12, - ObjectCreatedDeleteTagging = 13, + ObjectTaggingPut = 12, + ObjectTaggingDelete = 13, ObjectRemovedDelete = 14, ObjectRemovedDeleteMarkerCreated = 15, ObjectRemovedDeleteAllVersions = 16, @@ -63,6 +63,11 @@ pub enum EventName { ScannerLargeVersions = 30, // ObjectLargeVersions corresponding to Go ScannerBigPrefix = 31, // PrefixManyFolders corresponding to Go LifecycleDelMarkerExpirationDelete = 32, // ILMDelMarkerExpirationDelete corresponding to Go + ObjectAclPut = 33, + LifecycleExpirationDelete = 34, + LifecycleExpirationDeleteMarkerCreated = 35, + LifecycleTransition = 36, + IntelligentTiering = 37, // Compound "All" event type (no sequential value for mask) ObjectAccessedAll, @@ -70,6 +75,8 @@ pub enum EventName { ObjectRemovedAll, ObjectReplicationAll, ObjectRestoreAll, + ObjectTaggingAll, + LifecycleExpirationAll, ObjectTransitionAll, ObjectScannerAll, // New, from Go #[default] @@ -94,8 +101,8 @@ const SINGLE_EVENT_NAMES_IN_ORDER: [EventName; 32] = [ EventName::ObjectCreatedPut, EventName::ObjectCreatedPutRetention, EventName::ObjectCreatedPutLegalHold, - EventName::ObjectCreatedPutTagging, - EventName::ObjectCreatedDeleteTagging, + EventName::ObjectTaggingPut, + EventName::ObjectTaggingDelete, EventName::ObjectRemovedDelete, EventName::ObjectRemovedDeleteMarkerCreated, EventName::ObjectRemovedDeleteAllVersions, @@ -117,7 +124,15 @@ const SINGLE_EVENT_NAMES_IN_ORDER: [EventName; 32] = [ EventName::LifecycleDelMarkerExpirationDelete, ]; -const LAST_SINGLE_TYPE_VALUE: u32 = EventName::LifecycleDelMarkerExpirationDelete as u32; +const SINGLE_AWS_AND_EXTENSION_EVENTS_AFTER_COMPAT: [EventName; 5] = [ + EventName::ObjectAclPut, + EventName::LifecycleExpirationDelete, + EventName::LifecycleExpirationDeleteMarkerCreated, + EventName::LifecycleTransition, + EventName::IntelligentTiering, +]; + +const LAST_SINGLE_TYPE_VALUE: u32 = EventName::IntelligentTiering as u32; impl EventName { /// The parsed string is EventName. @@ -138,14 +153,21 @@ impl EventName { "s3:ObjectCreated:Put" => Ok(EventName::ObjectCreatedPut), "s3:ObjectCreated:PutRetention" => Ok(EventName::ObjectCreatedPutRetention), "s3:ObjectCreated:PutLegalHold" => Ok(EventName::ObjectCreatedPutLegalHold), - "s3:ObjectCreated:PutTagging" => Ok(EventName::ObjectCreatedPutTagging), - "s3:ObjectCreated:DeleteTagging" => Ok(EventName::ObjectCreatedDeleteTagging), + "s3:ObjectCreated:PutTagging" => Ok(EventName::ObjectTaggingPut), + "s3:ObjectCreated:DeleteTagging" => Ok(EventName::ObjectTaggingDelete), + "s3:ObjectTagging:*" => Ok(EventName::ObjectTaggingAll), + "s3:ObjectTagging:Put" => Ok(EventName::ObjectTaggingPut), + "s3:ObjectTagging:Delete" => Ok(EventName::ObjectTaggingDelete), + "s3:ObjectAcl:Put" => Ok(EventName::ObjectAclPut), "s3:ObjectRemoved:*" => Ok(EventName::ObjectRemovedAll), "s3:ObjectRemoved:Delete" => Ok(EventName::ObjectRemovedDelete), "s3:ObjectRemoved:DeleteMarkerCreated" => Ok(EventName::ObjectRemovedDeleteMarkerCreated), "s3:ObjectRemoved:NoOP" => Ok(EventName::ObjectRemovedNoOP), "s3:ObjectRemoved:DeleteAllVersions" => Ok(EventName::ObjectRemovedDeleteAllVersions), - "s3:LifecycleDelMarkerExpiration:Delete" => Ok(EventName::LifecycleDelMarkerExpirationDelete), + "s3:LifecycleDelMarkerExpiration:Delete" => Ok(EventName::LifecycleExpirationDeleteMarkerCreated), + "s3:LifecycleExpiration:*" => Ok(EventName::LifecycleExpirationAll), + "s3:LifecycleExpiration:Delete" => Ok(EventName::LifecycleExpirationDelete), + "s3:LifecycleExpiration:DeleteMarkerCreated" => Ok(EventName::LifecycleExpirationDeleteMarkerCreated), "s3:Replication:*" => Ok(EventName::ObjectReplicationAll), "s3:Replication:OperationFailedReplication" => Ok(EventName::ObjectReplicationFailed), "s3:Replication:OperationCompletedReplication" => Ok(EventName::ObjectReplicationComplete), @@ -156,8 +178,10 @@ impl EventName { "s3:ObjectRestore:Post" => Ok(EventName::ObjectRestorePost), "s3:ObjectRestore:Completed" => Ok(EventName::ObjectRestoreCompleted), "s3:ObjectTransition:Failed" => Ok(EventName::ObjectTransitionFailed), - "s3:ObjectTransition:Complete" => Ok(EventName::ObjectTransitionComplete), + "s3:ObjectTransition:Complete" => Ok(EventName::LifecycleTransition), "s3:ObjectTransition:*" => Ok(EventName::ObjectTransitionAll), + "s3:LifecycleTransition" => Ok(EventName::LifecycleTransition), + "s3:IntelligentTiering" => Ok(EventName::IntelligentTiering), "s3:Scanner:ManyVersions" => Ok(EventName::ScannerManyVersions), "s3:Scanner:LargeVersions" => Ok(EventName::ScannerLargeVersions), "s3:Scanner:BigPrefix" => Ok(EventName::ScannerBigPrefix), @@ -182,16 +206,21 @@ impl EventName { EventName::ObjectCreatedCopy => "s3:ObjectCreated:Copy", EventName::ObjectCreatedPost => "s3:ObjectCreated:Post", EventName::ObjectCreatedPut => "s3:ObjectCreated:Put", - EventName::ObjectCreatedPutTagging => "s3:ObjectCreated:PutTagging", - EventName::ObjectCreatedDeleteTagging => "s3:ObjectCreated:DeleteTagging", EventName::ObjectCreatedPutRetention => "s3:ObjectCreated:PutRetention", EventName::ObjectCreatedPutLegalHold => "s3:ObjectCreated:PutLegalHold", + EventName::ObjectTaggingAll => "s3:ObjectTagging:*", + EventName::ObjectTaggingPut => "s3:ObjectTagging:Put", + EventName::ObjectTaggingDelete => "s3:ObjectTagging:Delete", + EventName::ObjectAclPut => "s3:ObjectAcl:Put", EventName::ObjectRemovedAll => "s3:ObjectRemoved:*", EventName::ObjectRemovedDelete => "s3:ObjectRemoved:Delete", EventName::ObjectRemovedDeleteMarkerCreated => "s3:ObjectRemoved:DeleteMarkerCreated", EventName::ObjectRemovedNoOP => "s3:ObjectRemoved:NoOP", EventName::ObjectRemovedDeleteAllVersions => "s3:ObjectRemoved:DeleteAllVersions", EventName::LifecycleDelMarkerExpirationDelete => "s3:LifecycleDelMarkerExpiration:Delete", + EventName::LifecycleExpirationAll => "s3:LifecycleExpiration:*", + EventName::LifecycleExpirationDelete => "s3:LifecycleExpiration:Delete", + EventName::LifecycleExpirationDeleteMarkerCreated => "s3:LifecycleExpiration:DeleteMarkerCreated", EventName::ObjectReplicationAll => "s3:Replication:*", EventName::ObjectReplicationFailed => "s3:Replication:OperationFailedReplication", EventName::ObjectReplicationComplete => "s3:Replication:OperationCompletedReplication", @@ -204,6 +233,8 @@ impl EventName { EventName::ObjectTransitionAll => "s3:ObjectTransition:*", EventName::ObjectTransitionFailed => "s3:ObjectTransition:Failed", EventName::ObjectTransitionComplete => "s3:ObjectTransition:Complete", + EventName::LifecycleTransition => "s3:LifecycleTransition", + EventName::IntelligentTiering => "s3:IntelligentTiering", EventName::ScannerManyVersions => "s3:Scanner:ManyVersions", EventName::ScannerLargeVersions => "s3:Scanner:LargeVersions", EventName::ScannerBigPrefix => "s3:Scanner:BigPrefix", @@ -231,17 +262,9 @@ impl EventName { EventName::ObjectCreatedCopy, EventName::ObjectCreatedPost, EventName::ObjectCreatedPut, - EventName::ObjectCreatedPutRetention, - EventName::ObjectCreatedPutLegalHold, - EventName::ObjectCreatedPutTagging, - EventName::ObjectCreatedDeleteTagging, - ], - EventName::ObjectRemovedAll => vec![ - EventName::ObjectRemovedDelete, - EventName::ObjectRemovedDeleteMarkerCreated, - EventName::ObjectRemovedNoOP, - EventName::ObjectRemovedDeleteAllVersions, ], + EventName::ObjectTaggingAll => vec![EventName::ObjectTaggingPut, EventName::ObjectTaggingDelete], + EventName::ObjectRemovedAll => vec![EventName::ObjectRemovedDelete, EventName::ObjectRemovedDeleteMarkerCreated], EventName::ObjectReplicationAll => vec![ EventName::ObjectReplicationFailed, EventName::ObjectReplicationComplete, @@ -250,7 +273,15 @@ impl EventName { EventName::ObjectReplicationReplicatedAfterThreshold, ], EventName::ObjectRestoreAll => vec![EventName::ObjectRestorePost, EventName::ObjectRestoreCompleted], - EventName::ObjectTransitionAll => vec![EventName::ObjectTransitionFailed, EventName::ObjectTransitionComplete], + EventName::LifecycleExpirationAll => vec![ + EventName::LifecycleExpirationDelete, + EventName::LifecycleExpirationDeleteMarkerCreated, + ], + EventName::ObjectTransitionAll => vec![ + EventName::ObjectTransitionFailed, + EventName::ObjectTransitionComplete, + EventName::LifecycleTransition, + ], EventName::ObjectScannerAll => vec![ // New EventName::ScannerManyVersions, @@ -259,7 +290,9 @@ impl EventName { ], EventName::Everything => { // New - SINGLE_EVENT_NAMES_IN_ORDER.to_vec() + let mut all = SINGLE_EVENT_NAMES_IN_ORDER.to_vec(); + all.extend(SINGLE_AWS_AND_EXTENSION_EVENTS_AFTER_COMPAT); + all } // A single type returns to itself directly _ => vec![*self], @@ -299,8 +332,9 @@ impl EventName { EventName::ObjectCreatedPut => Some(S3Operation::PutObject), EventName::ObjectCreatedPutRetention => Some(S3Operation::PutObjectRetention), EventName::ObjectCreatedPutLegalHold => Some(S3Operation::PutObjectLegalHold), - EventName::ObjectCreatedPutTagging => Some(S3Operation::PutObjectTagging), - EventName::ObjectCreatedDeleteTagging => Some(S3Operation::DeleteObjectTagging), + EventName::ObjectTaggingPut => Some(S3Operation::PutObjectTagging), + EventName::ObjectTaggingDelete => Some(S3Operation::DeleteObjectTagging), + EventName::ObjectAclPut => Some(S3Operation::PutObjectAcl), EventName::ObjectRemovedDelete => Some(S3Operation::DeleteObject), EventName::ObjectRemovedDeleteMarkerCreated => Some(S3Operation::DeleteObject), EventName::ObjectRemovedDeleteAllVersions => Some(S3Operation::DeleteObject), @@ -497,16 +531,17 @@ impl S3Operation { Self::DeleteBucket => Some(EventName::BucketRemoved), Self::DeleteObject => Some(EventName::ObjectRemovedDelete), Self::DeleteObjects => Some(EventName::ObjectRemovedDeleteObjects), - Self::DeleteObjectTagging => Some(EventName::ObjectCreatedDeleteTagging), + Self::DeleteObjectTagging => Some(EventName::ObjectTaggingDelete), Self::GetObject => Some(EventName::ObjectAccessedGet), Self::GetObjectAttributes => Some(EventName::ObjectAccessedAttributes), Self::GetObjectLegalHold => Some(EventName::ObjectAccessedGetLegalHold), Self::GetObjectRetention => Some(EventName::ObjectAccessedGetRetention), Self::HeadObject => Some(EventName::ObjectAccessedHead), Self::PutObject => Some(EventName::ObjectCreatedPut), + Self::PutObjectAcl => Some(EventName::ObjectAclPut), Self::PutObjectLegalHold => Some(EventName::ObjectCreatedPutLegalHold), Self::PutObjectRetention => Some(EventName::ObjectCreatedPutRetention), - Self::PutObjectTagging => Some(EventName::ObjectCreatedPutTagging), + Self::PutObjectTagging => Some(EventName::ObjectTaggingPut), Self::RestoreObject => Some(EventName::ObjectRestorePost), Self::SelectObjectContent => Some(EventName::ObjectAccessedGet), Self::AbortMultipartUpload => Some(EventName::ObjectRemovedAbortMultipartUpload), @@ -541,6 +576,10 @@ mod tests { event: EventName::ObjectCreatedPut, serialized_str: "\"s3:ObjectCreated:Put\"", }, + TestCase { + event: EventName::ObjectTaggingPut, + serialized_str: "\"s3:ObjectTagging:Put\"", + }, ]; for case in &test_cases { @@ -574,6 +613,9 @@ mod tests { #[test] fn test_s3_operation_to_event_name() { assert_eq!(S3Operation::PutObject.to_event_name(), Some(EventName::ObjectCreatedPut)); + assert_eq!(S3Operation::PutObjectAcl.to_event_name(), Some(EventName::ObjectAclPut)); + assert_eq!(S3Operation::PutObjectTagging.to_event_name(), Some(EventName::ObjectTaggingPut)); + assert_eq!(S3Operation::DeleteObjectTagging.to_event_name(), Some(EventName::ObjectTaggingDelete)); assert_eq!(S3Operation::GetObject.to_event_name(), Some(EventName::ObjectAccessedGet)); assert_eq!(S3Operation::ListBuckets.to_event_name(), None); assert_eq!(S3Operation::RestoreObject.to_event_name(), Some(EventName::ObjectRestorePost)); @@ -587,6 +629,9 @@ mod tests { #[test] fn test_event_name_to_s3_operation() { assert_eq!(EventName::ObjectCreatedPut.to_s3_operation(), Some(S3Operation::PutObject)); + assert_eq!(EventName::ObjectAclPut.to_s3_operation(), Some(S3Operation::PutObjectAcl)); + assert_eq!(EventName::ObjectTaggingPut.to_s3_operation(), Some(S3Operation::PutObjectTagging)); + assert_eq!(EventName::ObjectTaggingDelete.to_s3_operation(), Some(S3Operation::DeleteObjectTagging)); assert_eq!(EventName::ObjectAccessedGet.to_s3_operation(), Some(S3Operation::GetObject)); assert_eq!(EventName::BucketCreated.to_s3_operation(), Some(S3Operation::CreateBucket)); assert_eq!(EventName::Everything.to_s3_operation(), None); @@ -597,4 +642,32 @@ mod tests { Some(S3Operation::AbortMultipartUpload) ); } + + #[test] + fn test_event_name_aliases_parse_to_aws_compatible_variants() { + assert_eq!(EventName::parse("s3:ObjectCreated:PutTagging").unwrap(), EventName::ObjectTaggingPut); + assert_eq!( + EventName::parse("s3:ObjectCreated:DeleteTagging").unwrap(), + EventName::ObjectTaggingDelete + ); + assert_eq!(EventName::parse("s3:ObjectTransition:Complete").unwrap(), EventName::LifecycleTransition); + assert_eq!( + EventName::parse("s3:LifecycleDelMarkerExpiration:Delete").unwrap(), + EventName::LifecycleExpirationDeleteMarkerCreated + ); + } + + #[test] + fn test_object_created_all_expansion_matches_aws_scope() { + let expanded = EventName::ObjectCreatedAll.expand(); + assert_eq!( + expanded, + vec![ + EventName::ObjectCreatedCompleteMultipartUpload, + EventName::ObjectCreatedCopy, + EventName::ObjectCreatedPost, + EventName::ObjectCreatedPut, + ] + ); + } } diff --git a/rustfs/src/app/multipart_usecase.rs b/rustfs/src/app/multipart_usecase.rs index 6079c5892..f9c9ba5a3 100644 --- a/rustfs/src/app/multipart_usecase.rs +++ b/rustfs/src/app/multipart_usecase.rs @@ -517,7 +517,9 @@ impl DefaultMultipartUsecase { let _ = context.object_store(); } - let helper = OperationHelper::new(&req, EventName::ObjectCreatedPut, S3Operation::CreateMultipartUpload); + let helper = + OperationHelper::new(&req, EventName::ObjectCreatedCreateMultipartUpload, S3Operation::CreateMultipartUpload) + .suppress_event(); let CreateMultipartUploadInput { bucket, key, diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index ccd379d54..813af36e6 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -2010,13 +2010,14 @@ impl DefaultObjectUsecase { let _ = context.object_store(); } + let mut helper = OperationHelper::new(&req, EventName::ObjectAclPut, S3Operation::PutObjectAcl); let PutObjectAclInput { bucket, key, access_control_policy, version_id, .. - } = req.input; + } = req.input.clone(); let Some(store) = new_object_layer_fn() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); @@ -2025,7 +2026,7 @@ impl DefaultObjectUsecase { let opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), None, &req.headers) .await .map_err(ApiError::from)?; - store.get_object_info(&bucket, &key, &opts).await.map_err(ApiError::from)?; + let object_info = store.get_object_info(&bucket, &key, &opts).await.map_err(ApiError::from)?; if access_control_policy.is_some() { return Err(s3_error!( @@ -2034,7 +2035,14 @@ impl DefaultObjectUsecase { )); } - Ok(S3Response::new(PutObjectAclOutput::default())) + let event_version_id = version_id + .or_else(|| object_info.version_id.map(|version_id| version_id.to_string())) + .unwrap_or_default(); + helper = helper.object(object_info).version_id(event_version_id); + + let result = Ok(S3Response::new(PutObjectAclOutput::default())); + let _ = helper.complete(&result); + result } pub async fn execute_put_object_legal_hold( @@ -2045,7 +2053,8 @@ impl DefaultObjectUsecase { let _ = context.object_store(); } - let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedPutLegalHold, S3Operation::PutObjectLegalHold); + let mut helper = + OperationHelper::new(&req, EventName::ObjectCreatedPutLegalHold, S3Operation::PutObjectLegalHold).suppress_event(); let PutObjectLegalHoldInput { bucket, key, @@ -2176,7 +2185,8 @@ impl DefaultObjectUsecase { let _ = context.object_store(); } - let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedPutRetention, S3Operation::PutObjectRetention); + let mut helper = + OperationHelper::new(&req, EventName::ObjectCreatedPutRetention, S3Operation::PutObjectRetention).suppress_event(); let PutObjectRetentionInput { bucket, key, @@ -2269,7 +2279,7 @@ impl DefaultObjectUsecase { } let start_time = std::time::Instant::now(); - let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedPutTagging, S3Operation::PutObjectTagging); + let mut helper = OperationHelper::new(&req, EventName::ObjectTaggingPut, S3Operation::PutObjectTagging); let PutObjectTaggingInput { bucket, key: object, @@ -2329,20 +2339,50 @@ impl DefaultObjectUsecase { ApiError::from(e) })?; + let event_object_info = match store.get_object_info(&bucket, &object, &opts).await { + Ok(info) => Some(info), + Err(err) => { + warn!( + bucket = %bucket, + object = %object, + version_id = ?req.input.version_id, + error = %err, + "failed to load object info for put-object-tagging notification; falling back to request context" + ); + None + } + }; + let manager = get_concurrency_manager(); let version_id = req.input.version_id.clone(); let cache_key = ConcurrencyManager::make_cache_key(&bucket, &object, version_id.clone().as_deref()); + let cache_bucket = bucket.clone(); + let cache_object = object.clone(); tokio::spawn(async move { manager - .invalidate_cache_versioned(&bucket, &object, version_id.as_deref()) + .invalidate_cache_versioned(&cache_bucket, &cache_object, version_id.as_deref()) .await; debug!("Cache invalidated for tagged object: {}", cache_key); }); counter!("rustfs.put_object_tagging.success").increment(1); - let version_id_resp = req.input.version_id.clone().unwrap_or_default(); - helper = helper.version_id(version_id_resp); + let event_version_id = req + .input + .version_id + .as_deref() + .filter(|version_id| !version_id.is_empty()) + .map(str::to_string) + .or_else(|| { + event_object_info + .as_ref() + .and_then(|info| info.version_id.map(|version_id| version_id.to_string())) + }) + .unwrap_or_default(); + if let Some(event_object_info) = event_object_info { + helper = helper.object(event_object_info); + } + helper = helper.version_id(event_version_id); let result = Ok(S3Response::new(PutObjectTaggingOutput { version_id: req.input.version_id.clone(), @@ -2555,7 +2595,7 @@ impl DefaultObjectUsecase { let concurrent_requests = bootstrap.concurrent_requests; let mut request_guard = bootstrap.request_guard; - let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedGet, S3Operation::GetObject); + let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedGet, S3Operation::GetObject).suppress_event(); // mc get 3 let request_context = Self::prepare_get_object_request_context(&req).await?; @@ -2709,7 +2749,8 @@ impl DefaultObjectUsecase { let _ = context.object_store(); } - let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedAttributes, S3Operation::GetObjectAttributes); + let mut helper = + OperationHelper::new(&req, EventName::ObjectAccessedAttributes, S3Operation::GetObjectAttributes).suppress_event(); let GetObjectAttributesInput { bucket, key, @@ -2941,7 +2982,8 @@ impl DefaultObjectUsecase { let _ = context.object_store(); } - let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedGetLegalHold, S3Operation::GetObjectLegalHold); + let mut helper = + OperationHelper::new(&req, EventName::ObjectAccessedGetLegalHold, S3Operation::GetObjectLegalHold).suppress_event(); let GetObjectLegalHoldInput { bucket, key, version_id, .. } = req.input.clone(); @@ -3032,7 +3074,8 @@ impl DefaultObjectUsecase { let _ = context.object_store(); } - let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedGetRetention, S3Operation::GetObjectRetention); + let mut helper = + OperationHelper::new(&req, EventName::ObjectAccessedGetRetention, S3Operation::GetObjectRetention).suppress_event(); let GetObjectRetentionInput { bucket, key, version_id, .. } = req.input.clone(); @@ -3949,7 +3992,7 @@ impl DefaultObjectUsecase { } let start_time = std::time::Instant::now(); - let mut helper = OperationHelper::new(&req, EventName::ObjectCreatedDeleteTagging, S3Operation::DeleteObjectTagging); + let mut helper = OperationHelper::new(&req, EventName::ObjectTaggingDelete, S3Operation::DeleteObjectTagging); let DeleteObjectTaggingInput { bucket, key: object, @@ -3973,22 +4016,50 @@ impl DefaultObjectUsecase { ApiError::from(e) })?; + let event_object_info = match store.get_object_info(&bucket, &object, &opts).await { + Ok(info) => Some(info), + Err(err) => { + warn!( + bucket = %bucket, + object = %object, + version_id = ?version_id, + error = %err, + "failed to load object info for delete-object-tagging notification; falling back to request context" + ); + None + } + }; + let manager = get_concurrency_manager(); let version_id_clone = version_id.clone(); + let cache_bucket = bucket.clone(); + let cache_object = object.clone(); tokio::spawn(async move { manager - .invalidate_cache_versioned(&bucket, &object, version_id_clone.as_deref()) + .invalidate_cache_versioned(&cache_bucket, &cache_object, version_id_clone.as_deref()) .await; debug!( "Cache invalidated for deleted tagged object: bucket={}, object={}, version_id={:?}", - bucket, object, version_id_clone + cache_bucket, cache_object, version_id_clone ); }); counter!("rustfs.delete_object_tagging.success").increment(1); - let version_id_resp = version_id.clone().unwrap_or_default(); - helper = helper.version_id(version_id_resp); + let event_version_id = version_id + .as_deref() + .filter(|value| !value.is_empty()) + .map(str::to_string) + .or_else(|| { + event_object_info + .as_ref() + .and_then(|info| info.version_id.map(|version_id| version_id.to_string())) + }) + .unwrap_or_default(); + if let Some(event_object_info) = event_object_info { + helper = helper.object(event_object_info); + } + helper = helper.version_id(event_version_id); let result = Ok(S3Response::new(DeleteObjectTaggingOutput { version_id })); let _ = helper.complete(&result); @@ -4003,7 +4074,7 @@ impl DefaultObjectUsecase { let _ = context.object_store(); } - let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedHead, S3Operation::HeadObject); + let mut helper = OperationHelper::new(&req, EventName::ObjectAccessedHead, S3Operation::HeadObject).suppress_event(); // mc get 2 let HeadObjectInput { bucket, @@ -4329,6 +4400,7 @@ impl DefaultObjectUsecase { let _ = context.object_store(); } + let mut helper = OperationHelper::new(&req, EventName::ObjectRestorePost, S3Operation::RestoreObject); let RestoreObjectInput { bucket, key: object, @@ -4392,6 +4464,7 @@ impl DefaultObjectUsecase { let mut header = HeaderMap::new(); + let event_object_info = obj_info.clone(); let obj_info_ = obj_info.clone(); if rreq.type_.as_ref().is_none_or(|t| t.as_str() != "SELECT") { obj_info.metadata_only = true; @@ -4445,7 +4518,13 @@ impl DefaultObjectUsecase { if already_restored { let output = restore::build_restore_object_output(Some(RequestCharged::from_static(RequestCharged::REQUESTER)), None); - return Ok(S3Response::new(output)); + helper = helper + .object(event_object_info.clone()) + .version_id(version_id_str.clone()) + .suppress_event(); + let result = Ok(S3Response::new(output)); + let _ = helper.complete(&result); + return result; } } @@ -4496,8 +4575,10 @@ impl DefaultObjectUsecase { }); let output = restore::build_restore_object_output(Some(RequestCharged::from_static(RequestCharged::REQUESTER)), None); - - Ok(S3Response::with_headers(output, header)) + helper = helper.object(event_object_info).version_id(version_id_str); + let result = Ok(S3Response::with_headers(output, header)); + let _ = helper.complete(&result); + result } #[instrument(level = "debug", skip(self, req))] diff --git a/rustfs/src/storage/helper.rs b/rustfs/src/storage/helper.rs index fa5230da9..b4466f538 100644 --- a/rustfs/src/storage/helper.rs +++ b/rustfs/src/storage/helper.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use crate::storage::access::ReqInfo; use http::StatusCode; use rustfs_audit::{ entity::{ApiDetails, ApiDetailsBuilder, AuditEntryBuilder}, @@ -59,8 +60,17 @@ impl OperationHelper { // Parse path -> bucket/object let path = req.uri.path().trim_start_matches('/'); let mut segs = path.splitn(2, '/'); - let bucket = segs.next().unwrap_or("").to_string(); - let object_key = segs.next().unwrap_or("").to_string(); + let path_bucket = segs.next().unwrap_or("").to_string(); + let path_object_key = segs.next().unwrap_or("").to_string(); + let req_info = req.extensions.get::(); + let bucket = req_info + .and_then(|info| info.bucket.clone()) + .filter(|value| !value.is_empty()) + .unwrap_or(path_bucket); + let object_key = req_info + .and_then(|info| info.object.clone()) + .filter(|value| !value.is_empty()) + .unwrap_or(path_object_key); // Infer remote address let remote_host = req @@ -98,13 +108,34 @@ impl OperationHelper { audit_builder = audit_builder.request_id(id_str); } + let event_object = ObjectInfo { + bucket: bucket.clone(), + name: object_key.clone(), + ..Default::default() + }; + + let mut req_params = extract_params_header(&req.headers); + if let Some(principal_id) = req_info + .and_then(|info| info.cred.as_ref()) + .map(|cred| cred.access_key.clone()) + .filter(|value| !value.is_empty()) + { + req_params.entry("principalId".to_string()).or_insert(principal_id); + } + // initialize event builder // object is a placeholder that must be set later using the `object()` method. - let event_builder = EventArgsBuilder::new(event, bucket, ObjectInfo::default()) + let mut event_builder = EventArgsBuilder::new(event, bucket, event_object) .host(get_request_host(&req.headers)) .port(get_request_port(&req.headers)) .user_agent(get_request_user_agent(&req.headers)) - .req_params(extract_params_header(&req.headers)); + .req_params(req_params); + if let Some(version_id) = req_info + .and_then(|info| info.version_id.clone()) + .filter(|value| !value.is_empty()) + { + event_builder = event_builder.version_id(version_id); + } Self { audit_builder: Some(audit_builder), @@ -222,3 +253,56 @@ impl Drop for OperationHelper { } } } + +#[cfg(test)] +mod tests { + use super::*; + use http::{Extensions, HeaderMap, HeaderValue, Method, Uri}; + use rustfs_credentials::Credentials; + use s3s::dto::DeleteObjectTaggingInput; + + fn build_request(input: T, method: Method, uri: Uri) -> S3Request { + S3Request { + input, + method, + uri, + headers: HeaderMap::new(), + extensions: Extensions::new(), + credentials: None, + region: None, + service: None, + trailing_headers: None, + } + } + + #[test] + fn operation_helper_uses_req_info_for_notification_context() { + let input = DeleteObjectTaggingInput::builder() + .bucket("input-bucket".to_string()) + .key("input-object".to_string()) + .build() + .unwrap(); + let mut req = build_request(input, Method::DELETE, Uri::from_static("/from-uri/ignored")); + req.headers.insert("host", HeaderValue::from_static("example.com")); + req.headers.insert("user-agent", HeaderValue::from_static("rustfs-test")); + req.extensions.insert(ReqInfo { + cred: Some(Credentials { + access_key: "notifyTag".to_string(), + ..Default::default() + }), + bucket: Some("issue-2292-bucket".to_string()), + object: Some("prefix/issue-2292.txt".to_string()), + version_id: Some("version-123".to_string()), + ..Default::default() + }); + + let helper = OperationHelper::new(&req, EventName::ObjectTaggingPut, S3Operation::PutObjectTagging); + let event_args = helper.event_builder.clone().expect("event builder should exist").build(); + + assert_eq!(event_args.bucket_name, "issue-2292-bucket"); + assert_eq!(event_args.object.bucket, "issue-2292-bucket"); + assert_eq!(event_args.object.name, "prefix/issue-2292.txt"); + assert_eq!(event_args.version_id, "version-123"); + assert_eq!(event_args.req_params.get("principalId").map(String::as_str), Some("notifyTag")); + } +}