diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_audit.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_audit.rs index d3420fc1e..cfeb644a8 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_audit.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_audit.rs @@ -12,7 +12,15 @@ // See the License for the specific language governing permissions and // limitations under the License. +use super::runtime_boundary as runtime_sources; use crate::bucket::lifecycle::lifecycle; +use crate::object_api::ObjectInfo; +use crate::services::event_notification::{EventArgs, send_event}; +use rustfs_common::metrics::IlmAction; +use rustfs_s3_types::EventName; + +const LIFECYCLE_EXPIRY_USER_AGENT: &str = "Internal: [ILM-Expiry]"; +const LIFECYCLE_TRANSITION_USER_AGENT: &str = "Internal: [ILM-Transition]"; #[derive(Debug, Clone, Default)] pub enum LcEventSrc { @@ -41,3 +49,84 @@ impl LcAuditEvent { Self { event, source } } } + +pub(crate) fn emit_transition_failed_event(object: ObjectInfo) { + emit_lifecycle_event(EventName::ObjectTransitionFailed, object, LIFECYCLE_TRANSITION_USER_AGENT); +} + +pub(crate) fn emit_transition_complete_event(object: ObjectInfo) { + emit_lifecycle_event(EventName::ObjectTransitionComplete, object, LIFECYCLE_TRANSITION_USER_AGENT); +} + +pub(crate) fn emit_transitioned_expiration_event(source: &ObjectInfo, deleted: &ObjectInfo) { + let event_name = transitioned_expiration_event_name(source.delete_marker, deleted.delete_marker); + let object = ObjectInfo { + bucket: source.bucket.clone(), + name: source.name.clone(), + size: source.size, + version_id: source.version_id, + delete_marker: source.delete_marker, + ..Default::default() + }; + emit_lifecycle_event(event_name, object, LIFECYCLE_EXPIRY_USER_AGENT); +} + +pub(crate) fn emit_non_transitioned_expiration_event(action: IlmAction, source: &ObjectInfo, deleted: ObjectInfo) { + let event_name = non_transitioned_expiration_event_name(action, source.delete_marker, deleted.delete_marker); + emit_lifecycle_event(event_name, deleted, LIFECYCLE_EXPIRY_USER_AGENT); +} + +fn emit_lifecycle_event(event_name: EventName, object: ObjectInfo, user_agent: &str) { + send_event(EventArgs { + event_name: event_name.to_string(), + bucket_name: object.bucket.clone(), + object, + user_agent: user_agent.to_string(), + host: runtime_sources::default_local_node_name(), + ..Default::default() + }); +} + +fn transitioned_expiration_event_name(source_delete_marker: bool, deleted_delete_marker: bool) -> EventName { + if source_delete_marker { + EventName::LifecycleExpirationDelete + } else if deleted_delete_marker { + EventName::LifecycleExpirationDeleteMarkerCreated + } else { + EventName::LifecycleExpirationDelete + } +} + +fn non_transitioned_expiration_event_name( + action: IlmAction, + source_delete_marker: bool, + deleted_delete_marker: bool, +) -> EventName { + match action { + IlmAction::DeleteAllVersionsAction | IlmAction::DelMarkerDeleteAllVersionsAction => EventName::LifecycleExpirationDelete, + _ if source_delete_marker => EventName::LifecycleExpirationDelete, + _ if deleted_delete_marker => EventName::LifecycleExpirationDeleteMarkerCreated, + _ => EventName::LifecycleExpirationDelete, + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn transitioned_expiration_event_marks_delete_marker_creation() { + assert_eq!( + transitioned_expiration_event_name(false, true), + EventName::LifecycleExpirationDeleteMarkerCreated + ); + } + + #[test] + fn non_transitioned_delete_all_uses_delete_event() { + assert_eq!( + non_transitioned_expiration_event_name(IlmAction::DeleteAllVersionsAction, false, true), + EventName::LifecycleExpirationDelete + ); + } +} diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 784110b67..21543460f 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -13,7 +13,10 @@ // limitations under the License. use super::runtime_boundary as runtime_sources; -use crate::bucket::lifecycle::bucket_lifecycle_audit::{LcAuditEvent, LcEventSrc}; +use crate::bucket::lifecycle::bucket_lifecycle_audit::{ + LcAuditEvent, LcEventSrc, emit_non_transitioned_expiration_event, emit_transition_complete_event, + emit_transition_failed_event, emit_transitioned_expiration_event, +}; use crate::bucket::lifecycle::evaluator::Evaluator; use crate::bucket::lifecycle::lifecycle::{ self, Lifecycle, ObjectOpts, TransitionOptions, abort_incomplete_multipart_upload_due, @@ -34,7 +37,6 @@ use crate::error::Error; use crate::error::StorageError; use crate::error::{error_resp_to_object_err, is_err_object_not_found, is_err_version_not_found, is_network_or_host_down}; use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions}; -use crate::services::event_notification::{EventArgs, send_event}; use crate::services::tier::warm_backend::WarmBackendGetOpts; use crate::set_disk::{MAX_PARTS_COUNT, RUSTFS_MULTIPART_BUCKET_KEY, RUSTFS_MULTIPART_OBJECT_KEY, SetDisks}; use crate::storage_api_contracts::{ @@ -63,7 +65,6 @@ use rustfs_filemeta::{ FileInfo, FileInfoOpts, NULL_VERSION_ID, REPLICATE_INCOMING_DELETE, ReplicateDecision, ReplicationState, RestoreStatusOps, VersionPurgeStatusType, get_file_info, is_restored_object_on_disk, }; -use rustfs_s3_types::EventName; use rustfs_utils::{get_env_i64, get_env_usize, path::encode_dir_object, string::strings_has_prefix_fold}; use s3s::dto::{ BucketLifecycleConfiguration, DefaultRetention, ExpirationStatus, ReplicationConfiguration, RestoreRequest, @@ -1230,15 +1231,7 @@ impl TransitionState { "Lifecycle tier operation failed" ); } - // Send s3:ObjectTransition:Failed event - send_event(EventArgs { - event_name: EventName::ObjectTransitionFailed.to_string(), - bucket_name: obj_info_for_event.bucket.clone(), - object: obj_info_for_event, - user_agent: "Internal: [ILM-Transition]".to_string(), - host: runtime_sources::default_local_node_name(), - ..Default::default() - }); + emit_transition_failed_event(obj_info_for_event); } else { global_metrics().record_scanner_transition_completed(1); let mut ts = TierStats { @@ -1251,15 +1244,7 @@ impl TransitionState { } transition_state.add_lastday_stats(&task.event.storage_class, ts); - // Send s3:ObjectTransition:Complete event - send_event(EventArgs { - event_name: EventName::ObjectTransitionComplete.to_string(), - bucket_name: obj_info_for_event.bucket.clone(), - object: obj_info_for_event, - user_agent: "Internal: [ILM-Transition]".to_string(), - host: runtime_sources::default_local_node_name(), - ..Default::default() - }); + emit_transition_complete_event(obj_info_for_event); } TransitionState::add_counter(&transition_state.active_tasks, -1); transition_state.record_scanner_transition_state(); @@ -2199,39 +2184,7 @@ pub async fn expire_transitioned_object( //audit_log_lifecycle(oi, ILMExpiry, tags); - 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() - }; - send_event(EventArgs { - event_name: event_name.to_string(), - bucket_name: obj_info.bucket.clone(), - object: obj_info, - user_agent: "Internal: [ILM-Expiry]".to_string(), - host: runtime_sources::default_local_node_name(), - ..Default::default() - }); - /*let system = match notification_system() { - Some(sys) => sys, - None => { - let config = Config::new(); - initialize(config).await?; - notification_system().expect("Failed to initialize notification system") - } - }; - let event = Arc::new(Event::new_test_event("my-bucket", "document.pdf", EventName::ObjectCreatedPut)); - system.send_event(event).await;*/ + emit_transitioned_expiration_event(oi, &dobj); Ok(dobj) } @@ -2681,20 +2634,7 @@ 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 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(), - object: dobj, - user_agent: "Internal: [ILM-Expiry]".to_string(), - host: runtime_sources::default_local_node_name(), - ..Default::default() - }); + emit_non_transitioned_expiration_event(lc_event.action, oi, dobj); if lc_event.action != IlmAction::NoneAction { let mut num_versions = 1_u64; diff --git a/scripts/check_architecture_migration_rules.sh b/scripts/check_architecture_migration_rules.sh index 20749f83f..1928a0370 100755 --- a/scripts/check_architecture_migration_rules.sh +++ b/scripts/check_architecture_migration_rules.sh @@ -122,6 +122,7 @@ STORE_API_DELETE_DTO_REEXPORTS_FILE="${TMP_DIR}/store_api_delete_dto_reexports.t STORE_API_DELETE_DTO_INTERNAL_HITS_FILE="${TMP_DIR}/store_api_delete_dto_internal_hits.txt" STORE_API_LIFECYCLE_HELPER_DEFINITION_HITS_FILE="${TMP_DIR}/store_api_lifecycle_helper_definition_hits.txt" STORE_API_LIFECYCLE_HELPER_OLD_CONSUMER_HITS_FILE="${TMP_DIR}/store_api_lifecycle_helper_old_consumer_hits.txt" +LIFECYCLE_AUDIT_SINK_BYPASS_HITS_FILE="${TMP_DIR}/lifecycle_audit_sink_bypass_hits.txt" STORE_API_EXTERNAL_LIST_CONSUMER_HITS_FILE="${TMP_DIR}/store_api_external_list_consumer_hits.txt" STORE_API_EXTERNAL_OPERATION_CONSUMER_HITS_FILE="${TMP_DIR}/store_api_external_operation_consumer_hits.txt" STORE_API_OBJECT_OPERATION_LOCAL_METHOD_HITS_FILE="${TMP_DIR}/store_api_object_operation_local_method_hits.txt" @@ -777,6 +778,18 @@ if [[ -s "$STORE_API_LIFECYCLE_HELPER_OLD_CONSUMER_HITS_FILE" ]]; then report_failure "lifecycle helper DTO consumers must import rustfs-storage-api directly: $(paste -sd '; ' "$STORE_API_LIFECYCLE_HELPER_OLD_CONSUMER_HITS_FILE")" fi +( + cd "$ROOT_DIR" + rg -n --with-filename '\b(?:send_event|EventArgs)\b' \ + crates/ecstore/src/bucket/lifecycle \ + --glob '*.rs' | + rg -v '^crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_audit\.rs:' || true +) >"$LIFECYCLE_AUDIT_SINK_BYPASS_HITS_FILE" + +if [[ -s "$LIFECYCLE_AUDIT_SINK_BYPASS_HITS_FILE" ]]; then + report_failure "lifecycle audit/event emission must stay behind bucket_lifecycle_audit sink: $(paste -sd '; ' "$LIFECYCLE_AUDIT_SINK_BYPASS_HITS_FILE")" +fi + ( cd "$ROOT_DIR" rg -n --no-heading 'rustfs_ecstore::store_api(?:::\{[^}]*\b(?:ListObjectVersionsInfo|ListObjectsV2Info|ObjectInfoOrErr)\b|::(?:ListObjectVersionsInfo|ListObjectsV2Info|ObjectInfoOrErr)\b)' \