From e160633ac44299a59096e454e48b1ada625543e7 Mon Sep 17 00:00:00 2001 From: houseme Date: Sun, 12 Apr 2026 23:59:42 +0800 Subject: [PATCH] fix(admin): enforce event and audit target auth (#2508) --- rustfs/src/admin/handlers/audit.rs | 55 +++++++++++++++++++++---- rustfs/src/admin/handlers/event.rs | 64 +++++++++++++++++++++++++----- 2 files changed, 102 insertions(+), 17 deletions(-) diff --git a/rustfs/src/admin/handlers/audit.rs b/rustfs/src/admin/handlers/audit.rs index b6395071f..58416b870 100644 --- a/rustfs/src/admin/handlers/audit.rs +++ b/rustfs/src/admin/handlers/audit.rs @@ -12,9 +12,12 @@ // See the License for the specific language governing permissions and // limitations under the License. -use crate::admin::router::{AdminOperation, Operation, S3Router}; +use crate::admin::{ + auth::validate_admin_request, + router::{AdminOperation, Operation, S3Router}, +}; use crate::auth::{check_key_valid, get_session_token}; -use crate::server::ADMIN_PREFIX; +use crate::server::{ADMIN_PREFIX, RemoteAddr}; use futures::stream::{FuturesUnordered, StreamExt}; use hashbrown::HashSet as HbHashSet; use http::{HeaderMap, StatusCode}; @@ -24,6 +27,7 @@ use rustfs_audit::{audit_system, start_audit_system as start_global_audit_system use rustfs_config::audit::{AUDIT_MQTT_KEYS, AUDIT_MQTT_SUB_SYS, AUDIT_ROUTE_PREFIX, AUDIT_WEBHOOK_KEYS, AUDIT_WEBHOOK_SUB_SYS}; use rustfs_config::{DEFAULT_DELIMITER, ENABLE_KEY, ENV_PREFIX, EnableState, MAX_ADMIN_REQUEST_BODY_SIZE}; use rustfs_ecstore::config::Config; +use rustfs_policy::policy::action::{Action, AdminAction}; use rustfs_targets::{TargetError, check_mqtt_broker_available_with_tls, target::mqtt::MQTTTlsConfig}; use s3s::{Body, S3Request, S3Response, S3Result, header::CONTENT_TYPE, s3_error}; use serde::{Deserialize, Serialize}; @@ -98,12 +102,14 @@ fn normalized_endpoint_key(account_id: &str, service: &str) -> EndpointKey { (account_id.to_lowercase(), service.to_string()) } -async fn check_permissions(req: &S3Request) -> S3Result<()> { +async fn authorize_audit_admin_request(req: &S3Request, action: AdminAction) -> S3Result<()> { let Some(input_cred) = &req.credentials else { return Err(s3_error!(InvalidRequest, "credentials not found")); }; - check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; - Ok(()) + let (cred, owner) = + check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; + let remote_addr = req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)); + validate_admin_request(&req.headers, &cred, owner, false, vec![Action::AdminAction(action)], remote_addr).await } fn build_response(status: StatusCode, body: Body, request_id: Option<&http::HeaderValue>) -> S3Response<(StatusCode, Body)> { @@ -465,7 +471,7 @@ impl Operation for AuditTargetConfig { let _enter = span.enter(); let (target_type, target_name) = extract_target_params(¶ms)?; - check_permissions(&req).await?; + authorize_audit_admin_request(&req, AdminAction::SetBucketTargetAction).await?; let config_snapshot = load_server_config_from_store().await?; if let Some(reason) = audit_target_mutation_block_reason(&config_snapshot, target_type, target_name) { return Err(s3_error!(InvalidRequest, "{reason}")); @@ -579,7 +585,7 @@ impl Operation for ListAuditTargets { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { let span = Span::current(); let _enter = span.enter(); - check_permissions(&req).await?; + authorize_audit_admin_request(&req, AdminAction::GetBucketTargetAction).await?; let mut runtime_statuses = HashMap::new(); if let Some(system) = audit_system() { @@ -622,7 +628,7 @@ impl Operation for RemoveAuditTarget { let _enter = span.enter(); let (target_type, target_name) = extract_target_params(¶ms)?; - check_permissions(&req).await?; + authorize_audit_admin_request(&req, AdminAction::SetBucketTargetAction).await?; let config_snapshot = load_server_config_from_store().await?; if let Some(reason) = audit_target_mutation_block_reason(&config_snapshot, target_type, target_name) { return Err(s3_error!(InvalidRequest, "{reason}")); @@ -907,4 +913,37 @@ mod tests { assert!(audit_target_mutation_block_reason(&config, AUDIT_WEBHOOK_SUB_SYS, "primary").is_none()); }); } + + #[test] + fn audit_target_handlers_require_admin_authorization_contract() { + let src = include_str!("audit.rs"); + let put_block = extract_block_between_markers(src, "impl Operation for AuditTargetConfig", "pub struct ListAuditTargets"); + let list_block = + extract_block_between_markers(src, "impl Operation for ListAuditTargets", "pub struct RemoveAuditTarget"); + let delete_block = extract_block_between_markers(src, "impl Operation for RemoveAuditTarget", "#[cfg(test)]"); + + assert!( + put_block.contains("authorize_audit_admin_request(&req, AdminAction::SetBucketTargetAction).await?;"), + "audit target writes should require SetBucketTargetAction" + ); + assert!( + list_block.contains("authorize_audit_admin_request(&req, AdminAction::GetBucketTargetAction).await?;"), + "audit target list should require GetBucketTargetAction" + ); + assert!( + delete_block.contains("authorize_audit_admin_request(&req, AdminAction::SetBucketTargetAction).await?;"), + "audit target deletion should require SetBucketTargetAction" + ); + } + + fn extract_block_between_markers<'a>(src: &'a str, start_marker: &str, end_marker: &str) -> &'a str { + let start = src + .find(start_marker) + .unwrap_or_else(|| panic!("Expected marker `{start_marker}` in source")); + let after_start = &src[start..]; + let end = after_start + .find(end_marker) + .unwrap_or_else(|| panic!("Expected end marker `{end_marker}` in source")); + &after_start[..end] + } } diff --git a/rustfs/src/admin/handlers/event.rs b/rustfs/src/admin/handlers/event.rs index ffb60f572..595163fff 100644 --- a/rustfs/src/admin/handlers/event.rs +++ b/rustfs/src/admin/handlers/event.rs @@ -12,9 +12,12 @@ // See the License for the specific language governing permissions and // limitations under the License. -use crate::admin::router::{AdminOperation, Operation, S3Router}; +use crate::admin::{ + auth::validate_admin_request, + router::{AdminOperation, Operation, S3Router}, +}; use crate::auth::{check_key_valid, get_session_token}; -use crate::server::ADMIN_PREFIX; +use crate::server::{ADMIN_PREFIX, RemoteAddr}; use futures::stream::{FuturesUnordered, StreamExt}; use hashbrown::HashSet as HbHashSet; use http::{HeaderMap, StatusCode}; @@ -25,6 +28,7 @@ use rustfs_config::notify::{ }; use rustfs_config::{DEFAULT_DELIMITER, ENABLE_KEY, ENV_PREFIX, EnableState, MAX_ADMIN_REQUEST_BODY_SIZE}; use rustfs_ecstore::config::Config; +use rustfs_policy::policy::action::{Action, AdminAction}; use rustfs_targets::{TargetError, check_mqtt_broker_available_with_tls, target::mqtt::MQTTTlsConfig}; use s3s::{Body, S3Request, S3Response, S3Result, header::CONTENT_TYPE, s3_error}; use serde::{Deserialize, Serialize}; @@ -107,12 +111,14 @@ fn normalized_endpoint_key(account_id: &str, service: &str) -> EndpointKey { // --- Helper Functions --- -async fn check_permissions(req: &S3Request) -> S3Result<()> { +async fn authorize_notification_admin_request(req: &S3Request, action: AdminAction) -> S3Result<()> { let Some(input_cred) = &req.credentials else { return Err(s3_error!(InvalidRequest, "credentials not found")); }; - check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; - Ok(()) + let (cred, owner) = + check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; + let remote_addr = req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)); + validate_admin_request(&req.headers, &cred, owner, false, vec![Action::AdminAction(action)], remote_addr).await } fn get_notification_system() -> S3Result> { @@ -392,7 +398,7 @@ impl Operation for NotificationTarget { let _enter = span.enter(); let (target_type, target_name) = extract_target_params(¶ms)?; - check_permissions(&req).await?; + authorize_notification_admin_request(&req, AdminAction::SetBucketTargetAction).await?; let ns = get_notification_system()?; let config_snapshot = ns.config.read().await.clone(); if let Some(reason) = target_mutation_block_reason(&config_snapshot, target_type, target_name) { @@ -502,7 +508,7 @@ impl Operation for ListNotificationTargets { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { let span = Span::current(); let _enter = span.enter(); - check_permissions(&req).await?; + authorize_notification_admin_request(&req, AdminAction::GetBucketTargetAction).await?; let ns = get_notification_system()?; let targets = ns.get_target_values().await; @@ -541,7 +547,7 @@ impl Operation for ListTargetsArns { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { let span = Span::current(); let _enter = span.enter(); - check_permissions(&req).await?; + authorize_notification_admin_request(&req, AdminAction::GetBucketTargetAction).await?; let ns = get_notification_system()?; let targets = ns.get_target_values().await; @@ -586,7 +592,7 @@ impl Operation for RemoveNotificationTarget { let _enter = span.enter(); let (target_type, target_name) = extract_target_params(¶ms)?; - check_permissions(&req).await?; + authorize_notification_admin_request(&req, AdminAction::SetBucketTargetAction).await?; let ns = get_notification_system()?; let config_snapshot = ns.config.read().await.clone(); if let Some(reason) = target_mutation_block_reason(&config_snapshot, target_type, target_name) { @@ -896,4 +902,44 @@ mod tests { }, ); } + + #[test] + fn notification_target_handlers_require_admin_authorization_contract() { + let src = include_str!("event.rs"); + let put_block = + extract_block_between_markers(src, "impl Operation for NotificationTarget", "pub struct ListNotificationTargets"); + let list_block = + extract_block_between_markers(src, "impl Operation for ListNotificationTargets", "pub struct ListTargetsArns"); + let arns_block = + extract_block_between_markers(src, "impl Operation for ListTargetsArns", "pub struct RemoveNotificationTarget"); + let delete_block = extract_block_between_markers(src, "impl Operation for RemoveNotificationTarget", "fn extract_param"); + + assert!( + put_block.contains("authorize_notification_admin_request(&req, AdminAction::SetBucketTargetAction).await?;"), + "notification target writes should require SetBucketTargetAction" + ); + assert!( + list_block.contains("authorize_notification_admin_request(&req, AdminAction::GetBucketTargetAction).await?;"), + "notification target list should require GetBucketTargetAction" + ); + assert!( + arns_block.contains("authorize_notification_admin_request(&req, AdminAction::GetBucketTargetAction).await?;"), + "notification target arn listing should require GetBucketTargetAction" + ); + assert!( + delete_block.contains("authorize_notification_admin_request(&req, AdminAction::SetBucketTargetAction).await?;"), + "notification target deletion should require SetBucketTargetAction" + ); + } + + fn extract_block_between_markers<'a>(src: &'a str, start_marker: &str, end_marker: &str) -> &'a str { + let start = src + .find(start_marker) + .unwrap_or_else(|| panic!("Expected marker `{start_marker}` in source")); + let after_start = &src[start..]; + let end = after_start + .find(end_marker) + .unwrap_or_else(|| panic!("Expected end marker `{end_marker}` in source")); + &after_start[..end] + } }