fix(audit): include deleted objects in bulk audit entries (#6592)

This commit is contained in:
RJ Regenold
2026-08-25 20:35:03 -05:00
committed by GitHub
parent ba629bdae0
commit 75a71fe6d7
2 changed files with 140 additions and 0 deletions
+79
View File
@@ -125,6 +125,7 @@ use http::{HeaderMap, HeaderValue, StatusCode};
use md5::{Digest as Md5Digest, Md5}; use md5::{Digest as Md5Digest, Md5};
use metrics::{counter, histogram}; use metrics::{counter, histogram};
use pin_project_lite::pin_project; use pin_project_lite::pin_project;
use rustfs_audit::ObjectVersion as AuditObjectVersion;
use rustfs_concurrency::GetObjectQueueSnapshot; use rustfs_concurrency::GetObjectQueueSnapshot;
use rustfs_config::MI_B; use rustfs_config::MI_B;
use rustfs_filemeta::{NULL_VERSION_ID, RestoreStatusOps, parse_restore_obj_status}; use rustfs_filemeta::{NULL_VERSION_ID, RestoreStatusOps, parse_restore_obj_status};
@@ -3266,6 +3267,20 @@ impl<T: Send + 'static> Drop for EagerPutCommitOwner<T> {
} }
} }
fn successful_delete_audit_objects(
delete: &s3s::dto::Delete,
successful_results: impl IntoIterator<Item = bool>,
) -> Vec<AuditObjectVersion> {
delete
.objects
.iter()
.zip(successful_results)
.filter_map(|(requested, successful)| {
successful.then(|| AuditObjectVersion::new(requested.key.clone(), requested.version_id.clone()))
})
.collect()
}
fn normalize_delete_objects_version_id( fn normalize_delete_objects_version_id(
version_id: Option<String>, version_id: Option<String>,
) -> std::result::Result<(Option<String>, Option<Uuid>), String> { ) -> std::result::Result<(Option<String>, Option<Uuid>), String> {
@@ -8695,6 +8710,13 @@ impl DefaultObjectUsecase {
errors: Some(errors), errors: Some(errors),
..Default::default() ..Default::default()
}; };
let helper = if helper.wants_audit_object_info() {
let audit_objects =
successful_delete_audit_objects(&delete, delete_results.iter().map(|result| result.delete_object.is_some()));
helper.audit_objects(audit_objects)
} else {
helper
};
let replication_deletes = if replicate_deletes { let replication_deletes = if replicate_deletes {
delete_results delete_results
@@ -18010,6 +18032,63 @@ mod tests {
assert_eq!(err.message(), Some("Not init")); assert_eq!(err.message(), Some("Not init"));
} }
#[test]
fn delete_objects_audit_details_include_only_successful_request_entries() {
let requested = vec![
ObjectIdentifier {
key: "first-key".to_string(),
version_id: None,
..Default::default()
},
ObjectIdentifier {
key: "denied-key".to_string(),
version_id: Some(Uuid::new_v4().to_string()),
..Default::default()
},
ObjectIdentifier {
key: "versioned-key".to_string(),
version_id: Some("requested-version".to_string()),
..Default::default()
},
];
let objects = successful_delete_audit_objects(
&Delete {
objects: requested,
quiet: Some(true),
},
[true, false, true],
);
assert_eq!(
objects,
vec![
AuditObjectVersion::new("first-key".to_string(), None),
AuditObjectVersion::new("versioned-key".to_string(), Some("requested-version".to_string())),
]
);
}
#[test]
fn delete_objects_audit_details_are_empty_when_every_entry_fails() {
let requested = vec![ObjectIdentifier {
key: "failed-key".to_string(),
version_id: None,
..Default::default()
}];
assert!(
successful_delete_audit_objects(
&Delete {
objects: requested,
quiet: None,
},
[false]
)
.is_empty()
);
}
#[test] #[test]
fn normalize_delete_objects_version_id_preserves_explicit_null_marker() { fn normalize_delete_objects_version_id_preserves_explicit_null_marker() {
let (wire_version_id, internal_version_id) = let (wire_version_id, internal_version_id) =
+61
View File
@@ -22,6 +22,7 @@ use hashbrown::HashMap;
use http::StatusCode; use http::StatusCode;
use metrics::counter; use metrics::counter;
use rustfs_audit::{ use rustfs_audit::{
ObjectVersion,
entity::{ApiDetails, ApiDetailsBuilder, AuditEntryBuilder}, entity::{ApiDetails, ApiDetailsBuilder, AuditEntryBuilder},
global::AuditLogger, global::AuditLogger,
}; };
@@ -246,6 +247,11 @@ impl OperationHelper {
matches!(self, Self::Enabled(state) if state.event_builder.is_some()) matches!(self, Self::Enabled(state) if state.event_builder.is_some())
} }
/// True when the audit entry can include object details.
pub fn wants_audit_object_info(&self) -> bool {
matches!(self, Self::Enabled(state) if state.audit_builder.is_some())
}
#[cfg(test)] #[cfg(test)]
pub(crate) fn event_args(&self) -> Option<rustfs_notify::EventArgs> { pub(crate) fn event_args(&self) -> Option<rustfs_notify::EventArgs> {
match self { match self {
@@ -264,6 +270,17 @@ impl OperationHelper {
self self
} }
/// Sets nonempty object details on the audit entry.
pub fn audit_objects(mut self, objects: Vec<ObjectVersion>) -> Self {
if !objects.is_empty()
&& let Self::Enabled(state) = &mut self
&& state.audit_builder.is_some()
{
state.api_builder = state.api_builder.clone().objects(objects);
}
self
}
/// Set the version ID for event notifications. /// Set the version ID for event notifications.
pub fn version_id(mut self, version_id: impl Into<String>) -> Self { pub fn version_id(mut self, version_id: impl Into<String>) -> Self {
if let Self::Enabled(state) = &mut self if let Self::Enabled(state) = &mut self
@@ -432,6 +449,7 @@ mod tests {
use base64::Engine as _; use base64::Engine as _;
use http::{Extensions, HeaderMap, HeaderValue, Method, Uri}; use http::{Extensions, HeaderMap, HeaderValue, Method, Uri};
use metrics::{Counter, CounterFn, Gauge, GaugeFn, Histogram, HistogramFn, Key, KeyName, Metadata, SharedString, Unit}; use metrics::{Counter, CounterFn, Gauge, GaugeFn, Histogram, HistogramFn, Key, KeyName, Metadata, SharedString, Unit};
use rustfs_audit::ObjectVersion;
use rustfs_credentials::Credentials; use rustfs_credentials::Credentials;
use rustfs_s3_ops::S3Operation; use rustfs_s3_ops::S3Operation;
use rustfs_s3_types::EventName; use rustfs_s3_types::EventName;
@@ -555,6 +573,49 @@ mod tests {
); );
} }
#[test]
fn operation_helper_adds_objects_to_audit_details() {
with_vars(
[
(rustfs_config::ENV_NOTIFY_ENABLE, Some("false")),
(rustfs_config::ENV_AUDIT_ENABLE, Some("true")),
],
|| {
refresh_notify_module_enabled();
refresh_audit_module_enabled();
let input = DeleteObjectTaggingInput::builder()
.bucket("test-bucket".to_string())
.key("test-key".to_string())
.build()
.expect("delete object tagging input should build");
let req = build_request(input, Method::DELETE, Uri::from_static("/test-bucket"));
let objects = vec![
ObjectVersion::new("first-key".to_string(), None),
ObjectVersion::new("second-key".to_string(), Some("version-123".to_string())),
];
let mut empty_helper = OperationHelper::new(&req, EventName::ObjectRemovedDelete, S3Operation::DeleteObjects)
.audit_objects(Vec::new());
let OperationHelper::Enabled(empty_state) = &mut empty_helper else {
panic!("helper should be enabled when the audit switch is on");
};
assert!(empty_state.api_builder.0.objects.is_none());
empty_state.audit_builder.take();
let result = Ok(S3Response::new(DeleteObjectTaggingOutput::default()));
let mut helper = OperationHelper::new(&req, EventName::ObjectRemovedDelete, S3Operation::DeleteObjects)
.audit_objects(objects.clone())
.complete(&result);
let OperationHelper::Enabled(state) = &mut helper else {
panic!("helper should be enabled when the audit switch is on");
};
let audit_entry = state.audit_builder.take().expect("audit builder should exist").build();
assert_eq!(audit_entry.api.objects, Some(objects));
},
);
}
#[test] #[test]
fn operation_helper_prioritizes_request_context_for_request_id() { fn operation_helper_prioritizes_request_context_for_request_id() {
with_vars( with_vars(