From 55396f13d48e87f58a3c5be11978a09cbd0c3298 Mon Sep 17 00:00:00 2001 From: GatewayJ <835269233@qq.com> Date: Fri, 27 Feb 2026 22:24:57 +0800 Subject: [PATCH] feat: policy add object tag (#1908) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: GatewayJ <8352692332qq.com> Co-authored-by: loverustfs Co-authored-by: 安正超 --- rustfs/src/storage/access.rs | 153 ++++++++++++++++++++++++++++++-- rustfs/src/storage/ecfs.rs | 55 +++++++++++- rustfs/src/storage/ecfs_test.rs | 42 +++++++++ 3 files changed, 243 insertions(+), 7 deletions(-) diff --git a/rustfs/src/storage/access.rs b/rustfs/src/storage/access.rs index e8f9ac2c7..9f7c68915 100644 --- a/rustfs/src/storage/access.rs +++ b/rustfs/src/storage/access.rs @@ -23,6 +23,7 @@ use crate::auth::{check_key_valid, get_condition_values, get_session_token}; use crate::error::ApiError; use crate::license::license_check; use crate::server::RemoteAddr; +use metrics::counter; use rustfs_ecstore::StorageAPI; use rustfs_ecstore::bucket::metadata_sys; use rustfs_ecstore::bucket::policy_sys::PolicySys; @@ -200,9 +201,45 @@ async fn check_acl_access(req: &S3Request, req_info: &ReqInfo, action: &Ac Ok(acl_allows(&acl, user_id, is_authenticated, permission, ignore_public_acls)) } +#[derive(Clone, Debug)] +pub(crate) struct ObjectTagConditions(pub HashMap>); + +/// Returns true if the bucket has a policy that uses `s3:ExistingObjectTag` (or +/// `ExistingObjectTag/...`) conditions. Used to skip fetching object tags when +/// no tag-based policy is in effect. +async fn bucket_policy_uses_existing_object_tag(bucket: &str) -> bool { + let Ok((policy_str, _)) = metadata_sys::get_bucket_policy_raw(bucket).await else { + return false; + }; + policy_str.contains("ExistingObjectTag") +} + +impl FS { + /// Fetches object tags (when the bucket policy requires them) and wraps them + /// as `ObjectTagConditions`. Returns `AccessDenied` on transient storage + /// errors to avoid fail-open on Deny policies. + async fn fetch_tag_conditions( + &self, + bucket: &str, + key: &str, + version_id: Option<&str>, + op: &'static str, + ) -> S3Result { + let tag_conditions = if bucket_policy_uses_existing_object_tag(bucket).await { + counter!("rustfs.object_tag_conditions.fetched", "op" => op).increment(1); + self.get_object_tag_conditions_for_policy(bucket, key, version_id).await? + } else { + counter!("rustfs.object_tag_conditions.skipped", "op" => op).increment(1); + HashMap::new() + }; + Ok(ObjectTagConditions(tag_conditions)) + } +} + /// Authorizes the request based on the action and credentials. pub async fn authorize_request(req: &mut S3Request, action: Action) -> S3Result<()> { let remote_addr = req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)); + let object_tag_conditions = req.extensions.get::().cloned(); let req_info = req.extensions.get::().expect("ReqInfo not found"); @@ -216,7 +253,16 @@ pub async fn authorize_request(req: &mut S3Request, action: Action) -> S3R let default_claims = HashMap::new(); let claims = cred.claims.as_ref().unwrap_or(&default_claims); - let conditions = get_condition_values(&req.headers, cred, req_info.version_id.as_deref(), None, remote_addr); + let mut conditions = get_condition_values(&req.headers, cred, req_info.version_id.as_deref(), None, remote_addr); + // Merge object tag conditions; extend existing values if the same key exists (e.g. from get_condition_values). + if let Some(ref tags) = object_tag_conditions { + for (k, v) in &tags.0 { + conditions + .entry(k.clone()) + .and_modify(|existing| existing.extend(v.iter().cloned())) + .or_insert_with(|| v.clone()); + } + } let bucket_name = req_info.bucket.as_deref().unwrap_or(""); if !bucket_name.is_empty() @@ -348,13 +394,22 @@ pub async fn authorize_request(req: &mut S3Request, action: Action) -> S3R } } } else { - let conditions = get_condition_values( + let mut conditions = get_condition_values( &req.headers, &rustfs_credentials::Credentials::default(), req_info.version_id.as_deref(), req.region.clone(), remote_addr, ); + // Merge object tag conditions; extend existing values if the same key exists. + if let Some(ref tags) = object_tag_conditions { + for (k, v) in &tags.0 { + conditions + .entry(k.clone()) + .and_modify(|existing| existing.extend(v.iter().cloned())) + .or_insert_with(|| v.clone()); + } + } let bucket_name = req_info.bucket.as_deref().unwrap_or(""); if !bucket_name.is_empty() @@ -526,7 +581,6 @@ impl S3Access for FS { /// This method returns `Ok(())` by default. async fn copy_object(&self, req: &mut S3Request) -> S3Result<()> { { - let req_info = req.extensions.get_mut::().expect("ReqInfo not found"); let (src_bucket, src_key, version_id) = match &req.input.copy_source { CopySource::AccessPoint { .. } => return Err(s3_error!(NotImplemented)), CopySource::Bucket { bucket, key, version_id } => { @@ -534,9 +588,15 @@ impl S3Access for FS { } }; - req_info.bucket = Some(src_bucket); - req_info.object = Some(src_key); - req_info.version_id = version_id; + let req_info = req.extensions.get_mut::().expect("ReqInfo not found"); + req_info.bucket = Some(src_bucket.clone()); + req_info.object = Some(src_key.clone()); + req_info.version_id = version_id.clone(); + + let tag_conds = self + .fetch_tag_conditions(&src_bucket, &src_key, version_id.as_deref(), "copy_object_src") + .await?; + req.extensions.insert(tag_conds); authorize_request(req, Action::S3Action(S3Action::GetObjectAction)).await?; } @@ -696,6 +756,11 @@ impl S3Access for FS { req_info.object = Some(req.input.key.clone()); req_info.version_id = req.input.version_id.clone(); + let tag_conds = self + .fetch_tag_conditions(&req.input.bucket, &req.input.key, req.input.version_id.as_deref(), "delete_object") + .await?; + req.extensions.insert(tag_conds); + authorize_request(req, Action::S3Action(S3Action::DeleteObjectAction)).await?; // S3 Standard: When bypass_governance header is set, must have s3:BypassGovernanceRetention permission @@ -715,6 +780,16 @@ impl S3Access for FS { req_info.object = Some(req.input.key.clone()); req_info.version_id = req.input.version_id.clone(); + let tag_conds = self + .fetch_tag_conditions( + &req.input.bucket, + &req.input.key, + req.input.version_id.as_deref(), + "delete_object_tagging", + ) + .await?; + req.extensions.insert(tag_conds); + authorize_request(req, Action::S3Action(S3Action::DeleteObjectTaggingAction)).await } @@ -947,6 +1022,11 @@ impl S3Access for FS { req_info.object = Some(req.input.key.clone()); req_info.version_id = req.input.version_id.clone(); + let tag_conds = self + .fetch_tag_conditions(&req.input.bucket, &req.input.key, req.input.version_id.as_deref(), "get_object") + .await?; + req.extensions.insert(tag_conds); + authorize_request(req, Action::S3Action(S3Action::GetObjectAction)).await } @@ -959,6 +1039,11 @@ impl S3Access for FS { req_info.object = Some(req.input.key.clone()); req_info.version_id = req.input.version_id.clone(); + let tag_conds = self + .fetch_tag_conditions(&req.input.bucket, &req.input.key, req.input.version_id.as_deref(), "get_object_acl") + .await?; + req.extensions.insert(tag_conds); + authorize_request(req, get_object_acl_authorize_action()).await } @@ -971,6 +1056,16 @@ impl S3Access for FS { req_info.object = Some(req.input.key.clone()); req_info.version_id = req.input.version_id.clone(); + let tag_conds = self + .fetch_tag_conditions( + &req.input.bucket, + &req.input.key, + req.input.version_id.as_deref(), + "get_object_attributes", + ) + .await?; + req.extensions.insert(tag_conds); + if req.input.version_id.is_some() { authorize_request(req, Action::S3Action(S3Action::GetObjectVersionAttributesAction)).await?; authorize_request(req, Action::S3Action(S3Action::GetObjectVersionAction)).await?; @@ -1025,6 +1120,11 @@ impl S3Access for FS { req_info.object = Some(req.input.key.clone()); req_info.version_id = req.input.version_id.clone(); + let tag_conds = self + .fetch_tag_conditions(&req.input.bucket, &req.input.key, req.input.version_id.as_deref(), "get_object_tagging") + .await?; + req.extensions.insert(tag_conds); + authorize_request(req, Action::S3Action(S3Action::GetObjectTaggingAction)).await } @@ -1064,6 +1164,11 @@ impl S3Access for FS { req_info.object = Some(req.input.key.clone()); req_info.version_id = req.input.version_id.clone(); + let tag_conds = self + .fetch_tag_conditions(&req.input.bucket, &req.input.key, req.input.version_id.as_deref(), "head_object") + .await?; + req.extensions.insert(tag_conds); + authorize_request(req, Action::S3Action(S3Action::GetObjectAction)).await } @@ -1356,6 +1461,11 @@ impl S3Access for FS { req_info.object = Some(req.input.key.clone()); req_info.version_id = req.input.version_id.clone(); + let tag_conds = self + .fetch_tag_conditions(&req.input.bucket, &req.input.key, req.input.version_id.as_deref(), "put_object_acl") + .await?; + req.extensions.insert(tag_conds); + authorize_request(req, put_object_acl_authorize_action()).await } @@ -1407,6 +1517,11 @@ impl S3Access for FS { req_info.object = Some(req.input.key.clone()); req_info.version_id = req.input.version_id.clone(); + let tag_conds = self + .fetch_tag_conditions(&req.input.bucket, &req.input.key, req.input.version_id.as_deref(), "put_object_tagging") + .await?; + req.extensions.insert(tag_conds); + authorize_request(req, Action::S3Action(S3Action::PutObjectTaggingAction)).await } @@ -1472,6 +1587,7 @@ impl S3Access for FS { #[cfg(test)] mod tests { use super::*; + use std::collections::HashMap; #[test] fn get_bucket_policy_uses_get_bucket_policy_action() { @@ -1492,4 +1608,29 @@ mod tests { fn put_object_acl_uses_put_object_acl_action() { assert_eq!(put_object_acl_authorize_action(), Action::S3Action(S3Action::PutObjectAclAction)); } + + /// Object tag conditions must use keys like ExistingObjectTag/ so that + /// bucket policy conditions (e.g. s3:ExistingObjectTag/security) are evaluated correctly. + #[test] + fn test_object_tag_conditions_key_format() { + let mut tags = HashMap::new(); + tags.insert("ExistingObjectTag/security".to_string(), vec!["public".to_string()]); + tags.insert("ExistingObjectTag/project".to_string(), vec!["webapp".to_string()]); + let object_tag_conditions = ObjectTagConditions(tags); + let mut conditions = HashMap::new(); + conditions.insert("delimiter".to_string(), vec!["/".to_string()]); + for (k, v) in &object_tag_conditions.0 { + conditions.insert(k.clone(), v.clone()); + } + assert_eq!(conditions.get("ExistingObjectTag/security"), Some(&vec!["public".to_string()])); + assert_eq!(conditions.get("ExistingObjectTag/project"), Some(&vec!["webapp".to_string()])); + assert_eq!(conditions.get("delimiter"), Some(&vec!["/".to_string()])); + } + + /// When bucket has no policy or policy fetch fails, tag-based check is skipped (returns false). + #[tokio::test] + async fn test_bucket_policy_uses_existing_object_tag_no_policy() { + let result = bucket_policy_uses_existing_object_tag("test-bucket-no-policy-xyz-absent").await; + assert!(!result, "bucket with no policy should not use ExistingObjectTag"); + } } diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 32f7e1bb1..2a7680cee 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -15,11 +15,18 @@ use crate::app::bucket_usecase::DefaultBucketUsecase; use crate::app::multipart_usecase::DefaultMultipartUsecase; use crate::app::object_usecase::DefaultObjectUsecase; +use rustfs_ecstore::{ + StorageAPI, + bucket::tagging::decode_tags_to_map, + error::{is_err_object_not_found, is_err_version_not_found}, + new_object_layer_fn, + store_api::ObjectOptions, +}; use s3s::{S3, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, dto::*, s3_error}; use serde::{Deserialize, Serialize}; use std::{fmt::Debug, sync::LazyLock}; use tokio::io::{AsyncRead, AsyncSeek}; -use tracing::{error, instrument}; +use tracing::{debug, error, instrument, warn}; use uuid::Uuid; const DEFAULT_OWNER_ID: &str = "rustfsadmin"; @@ -584,6 +591,52 @@ impl FS { Self {} } + pub async fn get_object_tag_conditions_for_policy( + &self, + bucket: &str, + object: &str, + version_id: Option<&str>, + ) -> S3Result>> { + let Some(store) = new_object_layer_fn() else { + return Ok(std::collections::HashMap::new()); + }; + let opts = ObjectOptions { + version_id: version_id.map(String::from), + ..Default::default() + }; + let tags = match store.get_object_tags(bucket, object, &opts).await { + Ok(t) => t, + Err(e) => { + if is_err_object_not_found(&e) || is_err_version_not_found(&e) { + debug!( + target: "rustfs::storage::ecfs", + bucket = %bucket, + object = %object, + version_id = ?version_id, + error = %e, + "object or version not found when fetching tags for policy; treating as no tags" + ); + return Ok(std::collections::HashMap::new()); + } + warn!( + target: "rustfs::storage::ecfs", + bucket = %bucket, + object = %object, + version_id = ?version_id, + error = %e, + "get_object_tags failed for policy conditions; denying request" + ); + return Err(s3_error!(AccessDenied, "Access Denied")); + } + }; + let map = decode_tags_to_map(&tags); + let mut out = std::collections::HashMap::new(); + for (k, v) in map { + out.insert(format!("ExistingObjectTag/{}", k), vec![v]); + } + Ok(out) + } + #[cfg(test)] pub(crate) fn normalize_delete_objects_version_id( &self, diff --git a/rustfs/src/storage/ecfs_test.rs b/rustfs/src/storage/ecfs_test.rs index 978bb4b34..6d4e82ef0 100644 --- a/rustfs/src/storage/ecfs_test.rs +++ b/rustfs/src/storage/ecfs_test.rs @@ -1256,4 +1256,46 @@ mod tests { assert!(result.is_ok(), "Should succeed with valid ARN"); assert_eq!(event_rules.len(), 1, "Should add one rule"); } + + // --- Object tag conditions for bucket policy (s3:ExistingObjectTag) --- + + /// Verifies that object tags are formatted as ExistingObjectTag/ condition keys + /// with a single-element vec value, matching the format expected by policy evaluation. + #[test] + fn test_object_tag_condition_key_format() { + use rustfs_ecstore::bucket::tagging::decode_tags_to_map; + use std::collections::HashMap; + + let tags_str = "security=public&project=webapp&env=prod"; + let map = decode_tags_to_map(tags_str); + let mut out: HashMap> = HashMap::new(); + for (k, v) in map { + out.insert(format!("ExistingObjectTag/{}", k), vec![v]); + } + + assert_eq!(out.get("ExistingObjectTag/security"), Some(&vec!["public".to_string()])); + assert_eq!(out.get("ExistingObjectTag/project"), Some(&vec!["webapp".to_string()])); + assert_eq!(out.get("ExistingObjectTag/env"), Some(&vec!["prod".to_string()])); + assert_eq!(out.len(), 3); + } + + /// When no object store is available (e.g. unit test env), get_object_tag_conditions_for_policy + /// returns Ok(empty map) so authorization can proceed without tag conditions. + #[tokio::test] + async fn test_get_object_tag_conditions_for_policy_returns_empty_without_store() { + let fs = FS::new(); + let out = fs.get_object_tag_conditions_for_policy("bucket", "key", None).await.unwrap(); + assert!(out.is_empty(), "without store should return empty tag conditions"); + } + + /// With version_id specified, the same no-store path returns Ok(empty) (versioned object path). + #[tokio::test] + async fn test_get_object_tag_conditions_for_policy_version_id_returns_empty_without_store() { + let fs = FS::new(); + let out = fs + .get_object_tag_conditions_for_policy("bucket", "key", Some("v1")) + .await + .unwrap(); + assert!(out.is_empty()); + } }