mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-31 09:18:28 +00:00
fix(admin): enforce event and audit target auth (#2508)
This commit is contained in:
@@ -12,9 +12,12 @@
|
|||||||
// See the License for the specific language governing permissions and
|
// See the License for the specific language governing permissions and
|
||||||
// limitations under the License.
|
// 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::auth::{check_key_valid, get_session_token};
|
||||||
use crate::server::ADMIN_PREFIX;
|
use crate::server::{ADMIN_PREFIX, RemoteAddr};
|
||||||
use futures::stream::{FuturesUnordered, StreamExt};
|
use futures::stream::{FuturesUnordered, StreamExt};
|
||||||
use hashbrown::HashSet as HbHashSet;
|
use hashbrown::HashSet as HbHashSet;
|
||||||
use http::{HeaderMap, StatusCode};
|
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::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_config::{DEFAULT_DELIMITER, ENABLE_KEY, ENV_PREFIX, EnableState, MAX_ADMIN_REQUEST_BODY_SIZE};
|
||||||
use rustfs_ecstore::config::Config;
|
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 rustfs_targets::{TargetError, check_mqtt_broker_available_with_tls, target::mqtt::MQTTTlsConfig};
|
||||||
use s3s::{Body, S3Request, S3Response, S3Result, header::CONTENT_TYPE, s3_error};
|
use s3s::{Body, S3Request, S3Response, S3Result, header::CONTENT_TYPE, s3_error};
|
||||||
use serde::{Deserialize, Serialize};
|
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())
|
(account_id.to_lowercase(), service.to_string())
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn check_permissions(req: &S3Request<Body>) -> S3Result<()> {
|
async fn authorize_audit_admin_request(req: &S3Request<Body>, action: AdminAction) -> S3Result<()> {
|
||||||
let Some(input_cred) = &req.credentials else {
|
let Some(input_cred) = &req.credentials else {
|
||||||
return Err(s3_error!(InvalidRequest, "credentials not found"));
|
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?;
|
let (cred, owner) =
|
||||||
Ok(())
|
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?;
|
||||||
|
let remote_addr = req.extensions.get::<Option<RemoteAddr>>().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)> {
|
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 _enter = span.enter();
|
||||||
let (target_type, target_name) = extract_target_params(¶ms)?;
|
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?;
|
let config_snapshot = load_server_config_from_store().await?;
|
||||||
if let Some(reason) = audit_target_mutation_block_reason(&config_snapshot, target_type, target_name) {
|
if let Some(reason) = audit_target_mutation_block_reason(&config_snapshot, target_type, target_name) {
|
||||||
return Err(s3_error!(InvalidRequest, "{reason}"));
|
return Err(s3_error!(InvalidRequest, "{reason}"));
|
||||||
@@ -579,7 +585,7 @@ impl Operation for ListAuditTargets {
|
|||||||
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||||
let span = Span::current();
|
let span = Span::current();
|
||||||
let _enter = span.enter();
|
let _enter = span.enter();
|
||||||
check_permissions(&req).await?;
|
authorize_audit_admin_request(&req, AdminAction::GetBucketTargetAction).await?;
|
||||||
|
|
||||||
let mut runtime_statuses = HashMap::new();
|
let mut runtime_statuses = HashMap::new();
|
||||||
if let Some(system) = audit_system() {
|
if let Some(system) = audit_system() {
|
||||||
@@ -622,7 +628,7 @@ impl Operation for RemoveAuditTarget {
|
|||||||
let _enter = span.enter();
|
let _enter = span.enter();
|
||||||
let (target_type, target_name) = extract_target_params(¶ms)?;
|
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?;
|
let config_snapshot = load_server_config_from_store().await?;
|
||||||
if let Some(reason) = audit_target_mutation_block_reason(&config_snapshot, target_type, target_name) {
|
if let Some(reason) = audit_target_mutation_block_reason(&config_snapshot, target_type, target_name) {
|
||||||
return Err(s3_error!(InvalidRequest, "{reason}"));
|
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());
|
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]
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,9 +12,12 @@
|
|||||||
// See the License for the specific language governing permissions and
|
// See the License for the specific language governing permissions and
|
||||||
// limitations under the License.
|
// 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::auth::{check_key_valid, get_session_token};
|
||||||
use crate::server::ADMIN_PREFIX;
|
use crate::server::{ADMIN_PREFIX, RemoteAddr};
|
||||||
use futures::stream::{FuturesUnordered, StreamExt};
|
use futures::stream::{FuturesUnordered, StreamExt};
|
||||||
use hashbrown::HashSet as HbHashSet;
|
use hashbrown::HashSet as HbHashSet;
|
||||||
use http::{HeaderMap, StatusCode};
|
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_config::{DEFAULT_DELIMITER, ENABLE_KEY, ENV_PREFIX, EnableState, MAX_ADMIN_REQUEST_BODY_SIZE};
|
||||||
use rustfs_ecstore::config::Config;
|
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 rustfs_targets::{TargetError, check_mqtt_broker_available_with_tls, target::mqtt::MQTTTlsConfig};
|
||||||
use s3s::{Body, S3Request, S3Response, S3Result, header::CONTENT_TYPE, s3_error};
|
use s3s::{Body, S3Request, S3Response, S3Result, header::CONTENT_TYPE, s3_error};
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
@@ -107,12 +111,14 @@ fn normalized_endpoint_key(account_id: &str, service: &str) -> EndpointKey {
|
|||||||
|
|
||||||
// --- Helper Functions ---
|
// --- Helper Functions ---
|
||||||
|
|
||||||
async fn check_permissions(req: &S3Request<Body>) -> S3Result<()> {
|
async fn authorize_notification_admin_request(req: &S3Request<Body>, action: AdminAction) -> S3Result<()> {
|
||||||
let Some(input_cred) = &req.credentials else {
|
let Some(input_cred) = &req.credentials else {
|
||||||
return Err(s3_error!(InvalidRequest, "credentials not found"));
|
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?;
|
let (cred, owner) =
|
||||||
Ok(())
|
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?;
|
||||||
|
let remote_addr = req.extensions.get::<Option<RemoteAddr>>().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<Arc<rustfs_notify::NotificationSystem>> {
|
fn get_notification_system() -> S3Result<Arc<rustfs_notify::NotificationSystem>> {
|
||||||
@@ -392,7 +398,7 @@ impl Operation for NotificationTarget {
|
|||||||
let _enter = span.enter();
|
let _enter = span.enter();
|
||||||
let (target_type, target_name) = extract_target_params(¶ms)?;
|
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 ns = get_notification_system()?;
|
||||||
let config_snapshot = ns.config.read().await.clone();
|
let config_snapshot = ns.config.read().await.clone();
|
||||||
if let Some(reason) = target_mutation_block_reason(&config_snapshot, target_type, target_name) {
|
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<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||||
let span = Span::current();
|
let span = Span::current();
|
||||||
let _enter = span.enter();
|
let _enter = span.enter();
|
||||||
check_permissions(&req).await?;
|
authorize_notification_admin_request(&req, AdminAction::GetBucketTargetAction).await?;
|
||||||
let ns = get_notification_system()?;
|
let ns = get_notification_system()?;
|
||||||
|
|
||||||
let targets = ns.get_target_values().await;
|
let targets = ns.get_target_values().await;
|
||||||
@@ -541,7 +547,7 @@ impl Operation for ListTargetsArns {
|
|||||||
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||||
let span = Span::current();
|
let span = Span::current();
|
||||||
let _enter = span.enter();
|
let _enter = span.enter();
|
||||||
check_permissions(&req).await?;
|
authorize_notification_admin_request(&req, AdminAction::GetBucketTargetAction).await?;
|
||||||
let ns = get_notification_system()?;
|
let ns = get_notification_system()?;
|
||||||
|
|
||||||
let targets = ns.get_target_values().await;
|
let targets = ns.get_target_values().await;
|
||||||
@@ -586,7 +592,7 @@ impl Operation for RemoveNotificationTarget {
|
|||||||
let _enter = span.enter();
|
let _enter = span.enter();
|
||||||
let (target_type, target_name) = extract_target_params(¶ms)?;
|
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 ns = get_notification_system()?;
|
||||||
let config_snapshot = ns.config.read().await.clone();
|
let config_snapshot = ns.config.read().await.clone();
|
||||||
if let Some(reason) = target_mutation_block_reason(&config_snapshot, target_type, target_name) {
|
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]
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user