fix: populate tagging notification principalId and object metadata (#2342)

This commit is contained in:
houseme
2026-03-30 14:32:39 +08:00
committed by GitHub
parent 387c385dfa
commit 16db18216d
11 changed files with 481 additions and 365 deletions
+3 -1
View File
@@ -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,
+103 -22
View File
@@ -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))]
+88 -4
View File
@@ -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::<ReqInfo>();
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<T>(input: T, method: Method, uri: Uri) -> S3Request<T> {
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"));
}
}