// Copyright 2024 RustFS Team // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. // You may obtain a copy of the License at // // http://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // See the License for the specific language governing permissions and // limitations under the License. use crate::error::Error as IamError; use crate::error::is_err_no_such_account; use crate::error::is_err_no_such_temp_account; use crate::error::{Error, Result}; use crate::federation::OIDC_VIRTUAL_PARENT_CLAIM; use crate::manager::extract_jwt_claims; use crate::manager::get_default_policies; use crate::manager::{IamCache, IamSyncMetricsSnapshot}; use crate::store::GroupInfo; use crate::store::MappedPolicy; use crate::store::Store; use crate::store::UserType; use crate::utils::{extract_claims, extract_claims_allow_missing_exp}; use crate::{ notify_iam_delete_policy, notify_iam_delete_user, notify_iam_load_group, notify_iam_load_policy, notify_iam_load_policy_mapping, notify_iam_load_service_account, notify_iam_load_user, }; use rustfs_credentials::{Credentials, EMBEDDED_POLICY_TYPE, INHERITED_POLICY_TYPE}; use rustfs_madmin::AddOrUpdateUserReq; use rustfs_madmin::GroupDesc; use rustfs_policy::arn::ARN; use rustfs_policy::auth::{ ACCOUNT_ON, UserIdentity, contains_reserved_chars, create_new_credentials_with_metadata, generate_credentials, is_access_key_valid, is_secret_key_valid, }; use rustfs_policy::policy::Args; use rustfs_policy::policy::opa; use rustfs_policy::policy::{Policy, PolicyDoc, iam_policy_claim_name_sa, policy_needs_existing_object_tag_for_args}; use serde_json::Value; use std::collections::{HashMap, HashSet}; use std::sync::Arc; use std::sync::OnceLock; use time::OffsetDateTime; use tokio::sync::RwLock; use tracing::{error, info, warn}; pub const MAX_SVCSESSION_POLICY_SIZE: usize = 4096; pub const SITE_REPLICATOR_SERVICE_ACCOUNT: &str = "site-replicator-0"; pub const STATUS_ENABLED: &str = "enabled"; pub const STATUS_DISABLED: &str = "disabled"; pub const POLICYNAME: &str = "policy"; pub const SESSION_POLICY_NAME: &str = "sessionPolicy"; pub const SESSION_POLICY_NAME_EXTRACTED: &str = "sessionPolicy-extracted"; const VERIFIED_FEDERATED_POLICY_CLAIM: &str = "x-rustfs-internal-federated-policy"; const FEDERATED_POLICY_BOUNDARY_CLAIM: &str = "x-rustfs-internal-federated-boundary"; pub(crate) const SITE_REPLICATOR_CLAIM: &str = "site-replicator"; #[derive(Clone)] enum PolicyPluginState { Disabled, Initializing, Ready(opa::AuthZPlugin), Failed, } static POLICY_PLUGIN_STATE: OnceLock>> = OnceLock::new(); fn get_policy_plugin_state() -> Arc> { POLICY_PLUGIN_STATE .get_or_init(|| { let configured = opa::is_configured(); let state = Arc::new(RwLock::new(if configured { PolicyPluginState::Initializing } else { PolicyPluginState::Disabled })); if configured { let state = Arc::clone(&state); tokio::spawn(async move { let next_state = match opa::lookup_config().await { Ok(conf) if conf.enable() => { info!("OPA plugin enabled"); PolicyPluginState::Ready(opa::AuthZPlugin::new(conf)) } Ok(_) => PolicyPluginState::Failed, Err(e) => { error!("Error loading OPA configuration err:{}", e); PolicyPluginState::Failed } }; *state.write().await = next_state; }); } state }) .clone() } pub struct IamSys { store: Arc>, roles_map: HashMap, } #[derive(Clone)] enum PreparedSessionPolicy { None, DenyAll, Policy(Policy), } #[derive(Clone, Copy)] enum PreparedServicePolicyMode { Inherited, SessionBound, } #[derive(Clone)] enum PreparedIamMode { Opa, Owner, Deny, Regular { combined_policy: Policy, }, Sts { is_owner: bool, combined_policy: Policy, session_policy: PreparedSessionPolicy, }, ServiceAccount { is_owner: bool, bypass_parent_policy: bool, parent_user: String, combined_policy: Policy, federated_boundary: Option, mode: PreparedServicePolicyMode, session_policy: PreparedSessionPolicy, }, } #[derive(Clone)] pub struct PreparedIamAuth { pub needs_existing_object_tag: bool, mode: PreparedIamMode, } impl PreparedIamAuth { /// Evaluate whether the already-prepared IAM context needs ExistingObjectTag /// conditions for the provided request args. pub async fn needs_existing_object_tag_for_args(&self, args: &Args<'_>) -> bool { match &self.mode { PreparedIamMode::Opa | PreparedIamMode::Owner | PreparedIamMode::Deny => false, PreparedIamMode::Regular { combined_policy } => { policy_needs_existing_object_tag_for_args(combined_policy, args).await } PreparedIamMode::Sts { combined_policy, session_policy, .. } => { policy_needs_existing_object_tag_for_args(combined_policy, args).await || prepared_session_policy_needs_existing_object_tag_for_args(session_policy, args).await } PreparedIamMode::ServiceAccount { combined_policy, federated_boundary, mode, session_policy, .. } => { policy_needs_existing_object_tag_for_args(combined_policy, args).await || match federated_boundary { Some(policy) => policy_needs_existing_object_tag_for_args(policy, args).await, None => false, } || matches!(mode, PreparedServicePolicyMode::SessionBound) && prepared_session_policy_needs_existing_object_tag_for_args(session_policy, args).await } } } /// Returns the resolved identity policy prepared for the current auth mode. /// /// This is intended for read-only views (for example `/accountinfo`) so /// callers can reuse the same policy resolution path as authorization. pub fn combined_policy_for_view(&self) -> Option<&Policy> { match &self.mode { PreparedIamMode::Regular { combined_policy } => Some(combined_policy), PreparedIamMode::Sts { combined_policy, .. } => Some(combined_policy), PreparedIamMode::ServiceAccount { combined_policy, .. } => Some(combined_policy), PreparedIamMode::Opa | PreparedIamMode::Owner | PreparedIamMode::Deny => None, } } } impl IamSys { /// Create a new IamSys instance with the given IamCache store /// /// # Arguments /// * `store` - An Arc to the IamCache instance /// /// # Returns /// A new instance of IamSys pub fn new(store: Arc>) -> Self { get_policy_plugin_state(); Self { store, roles_map: HashMap::new(), } } /// Check if the IamSys has a watcher configured /// /// # Returns /// `true` if a watcher is configured, `false` otherwise pub fn has_watcher(&self) -> bool { self.store.api.has_watcher() } pub fn sync_metrics_snapshot(&self) -> IamSyncMetricsSnapshot { self.store.sync_metrics_snapshot() } pub async fn set_policy_plugin_client(client: rustfs_policy::policy::opa::AuthZPlugin) { let policy_plugin_state = get_policy_plugin_state(); let mut guard = policy_plugin_state.write().await; *guard = PolicyPluginState::Ready(client); } pub async fn get_policy_plugin_client() -> Option { let policy_plugin_state = get_policy_plugin_state(); let guard = policy_plugin_state.read().await; match &*guard { PolicyPluginState::Ready(client) => Some(client.clone()), PolicyPluginState::Disabled | PolicyPluginState::Initializing | PolicyPluginState::Failed => None, } } async fn policy_plugin_state() -> PolicyPluginState { get_policy_plugin_state().read().await.clone() } pub async fn load_group(&self, name: &str) -> Result<()> { self.store.group_notification_handler(name).await } pub async fn load_groups(&self, m: &mut HashMap) -> Result<()> { self.store.api.load_groups(m).await } pub async fn load_policy(&self, name: &str) -> Result<()> { self.store.policy_notification_handler(name).await } pub async fn load_policy_mapping(&self, name: &str, user_type: UserType, is_group: bool) -> Result<()> { self.store .policy_mapping_notification_handler(name, user_type, is_group) .await } pub async fn load_user(&self, name: &str, user_type: UserType) -> Result<()> { self.store.user_notification_handler(name, user_type).await } pub async fn load_users(&self, user_type: UserType, m: &mut HashMap) -> Result<()> { self.store.api.load_users(user_type, m).await?; Ok(()) } pub async fn load_service_account(&self, name: &str) -> Result<()> { self.store.user_notification_handler(name, UserType::Svc).await } pub async fn delete_policy(&self, name: &str, notify: bool) -> Result<()> { for k in get_default_policies().keys() { if k == name { return Err(Error::other("system policy can not be deleted")); } } self.store.delete_policy(name, notify).await?; if !notify || self.has_watcher() { return Ok(()); } for r in notify_iam_delete_policy(name).await { if let Some(err) = r.err { warn!("notify delete_policy failed: {}", err); } } Ok(()) } pub async fn info_policy(&self, name: &str) -> Result { let d = self.store.get_policy_doc(name).await?; let pdata = serde_json::to_value(&d.policy)?; Ok(rustfs_madmin::PolicyInfo { policy_name: name.to_string(), policy: pdata, create_date: d.create_date, update_date: d.update_date, }) } pub async fn load_mapped_policies( &self, user_type: UserType, is_group: bool, m: &mut HashMap, ) -> Result<()> { self.store.api.load_mapped_policies(user_type, is_group, m).await } pub async fn list_policies(&self, bucket_name: &str) -> Result> { self.store.list_policies(bucket_name).await } /// Backward-compatible misspelling retained until the next breaking release. #[deprecated( since = "1.0.0", note = "use list_policies instead; this alias will be removed in the next breaking release" )] pub async fn list_polices(&self, bucket_name: &str) -> Result> { self.list_policies(bucket_name).await } pub async fn list_policy_docs(&self, bucket_name: &str) -> Result> { self.store.list_policy_docs(bucket_name).await } pub async fn set_policy(&self, name: &str, policy: Policy) -> Result { let updated_at = self.store.set_policy(name, policy).await?; if !self.has_watcher() { for r in notify_iam_load_policy(name).await { if let Some(err) = r.err { warn!("notify load_policy failed: {}", err); } } } Ok(updated_at) } pub async fn get_role_policy(&self, arn_str: &str) -> Result<(ARN, String)> { let Some(arn) = ARN::parse(arn_str).ok() else { return Err(Error::other("Invalid ARN")); }; let Some(policy) = self.roles_map.get(&arn) else { return Err(Error::other("No such role")); }; Ok((arn, policy.clone())) } pub async fn delete_user(&self, name: &str, notify: bool) -> Result<()> { self.store.delete_user(name, UserType::Reg).await?; if notify && !self.has_watcher() { for r in notify_iam_delete_user(name).await { if let Some(err) = r.err { warn!("notify delete_user failed: {}", err); } } } Ok(()) } /// Delete a single temporary/STS account by its access key. /// /// Unlike [`Self::delete_user`], which operates on regular (`UserType::Reg`) /// identities, this removes an STS credential (`UserType::Sts`) from both the /// in-memory cache and the backing store, immediately invalidating the /// associated session token. This is the primitive used by the admin /// `revoke-tokens` endpoint to revoke STS credentials for a parent user. pub async fn delete_temp_account(&self, access_key: &str, notify: bool) -> Result<()> { self.store.delete_user(access_key, UserType::Sts).await?; if notify && !self.has_watcher() { for r in notify_iam_delete_user(access_key).await { if let Some(err) = r.err { warn!("notify delete_temp_account failed: {}", err); } } } Ok(()) } async fn notify_for_user(&self, name: &str, is_temp: bool) { if self.has_watcher() { return; } // Fire-and-forget notification to peers - don't block auth operations // This is critical for cluster recovery: login should not wait for dead peers let name = name.to_string(); tokio::spawn(async move { for r in notify_iam_load_user(&name, is_temp).await { if let Some(err) = r.err { warn!("notify load_user failed (non-blocking): {}", err); } } }); } async fn notify_for_service_account(&self, name: &str) { if self.has_watcher() { return; } // Fire-and-forget notification to peers - don't block service account operations let name = name.to_string(); tokio::spawn(async move { for r in notify_iam_load_service_account(&name).await { if let Some(err) = r.err { warn!("notify load_service_account failed (non-blocking): {}", err); } } }); } pub async fn current_policies(&self, name: &str) -> String { self.store.merge_policies(name).await.0 } pub async fn list_bucket_users(&self, bucket_name: &str) -> Result> { self.store.get_bucket_users(bucket_name).await } pub async fn list_users(&self) -> Result> { self.store.get_users().await } pub async fn set_temp_user(&self, name: &str, cred: &Credentials, policy_name: Option<&str>) -> Result { let updated_at = self.store.set_temp_user(name, cred, policy_name).await?; self.notify_for_user(&cred.access_key, true).await; Ok(updated_at) } pub async fn is_temp_user(&self, name: &str) -> Result<(bool, String)> { let Some(u) = self.store.get_user(name).await else { return Err(IamError::NoSuchUser(name.to_string())); }; if u.credentials.is_temp() { Ok((true, u.credentials.parent_user)) } else { Ok((false, "".to_string())) } } pub async fn is_service_account(&self, name: &str) -> Result<(bool, String)> { let Some(u) = self.store.get_user(name).await else { return Err(IamError::NoSuchUser(name.to_string())); }; if u.credentials.is_service_account() { Ok((true, u.credentials.parent_user)) } else { Ok((false, "".to_string())) } } pub async fn get_user_info(&self, name: &str) -> Result { self.store.get_user_info(name).await } pub async fn set_user_status(&self, name: &str, status: rustfs_madmin::AccountStatus) -> Result { let updated_at = self.store.set_user_status(name, status).await?; self.notify_for_user(name, false).await; Ok(updated_at) } pub async fn new_service_account( &self, parent_user: &str, groups: Option>, opts: NewServiceAccountOpts, ) -> Result<(Credentials, OffsetDateTime)> { if parent_user.is_empty() { return Err(IamError::InvalidArgument); } if !opts.access_key.is_empty() && opts.secret_key.is_empty() { return Err(IamError::NoSecretKeyWithAccessKey); } if !opts.secret_key.is_empty() && opts.access_key.is_empty() { return Err(IamError::NoAccessKeyWithSecretKey); } if parent_user == opts.access_key { return Err(IamError::AccessKeyAlreadyExists); } if opts.access_key == SITE_REPLICATOR_SERVICE_ACCOUNT && !opts.allow_site_replicator_account { return Err(IamError::IAMActionNotAllowed); } let policy_buf = if let Some(policy) = opts.session_policy { policy.validate()?; let buf = serde_json::to_vec(&policy)?; if buf.len() > MAX_SVCSESSION_POLICY_SIZE { return Err(IamError::PolicyTooLarge); } buf } else { Vec::new() }; let mut m: HashMap = HashMap::new(); m.insert("parent".to_owned(), Value::String(parent_user.to_owned())); if opts.access_key == SITE_REPLICATOR_SERVICE_ACCOUNT && opts.allow_site_replicator_account { m.insert(SITE_REPLICATOR_CLAIM.to_owned(), Value::Bool(true)); } if !policy_buf.is_empty() { m.insert( SESSION_POLICY_NAME.to_owned(), Value::String(base64_simd::URL_SAFE_NO_PAD.encode_to_string(&policy_buf)), ); m.insert(iam_policy_claim_name_sa(), Value::String(EMBEDDED_POLICY_TYPE.to_owned())); } else { m.insert(iam_policy_claim_name_sa(), Value::String(INHERITED_POLICY_TYPE.to_owned())); } if let Some(claims) = opts.claims { for (k, v) in claims.iter() { if !m.contains_key(k) { m.insert(k.to_owned(), v.to_owned()); } } } if let Some(expiration) = opts.expiration { m.insert("exp".to_string(), Value::Number(serde_json::Number::from(expiration.unix_timestamp()))); } let (access_key, secret_key) = if !opts.access_key.is_empty() || !opts.secret_key.is_empty() { (opts.access_key, opts.secret_key) } else { generate_credentials()? }; let mut cred = create_new_credentials_with_metadata(&access_key, &secret_key, &m, &secret_key)?; cred.parent_user = parent_user.to_owned(); cred.groups = groups; cred.status = ACCOUNT_ON.to_owned(); cred.name = opts.name; cred.description = opts.description; let create_at = self.store.add_service_account(cred.clone()).await?; self.notify_for_service_account(&cred.access_key).await; Ok((cred, create_at)) } pub async fn update_service_account(&self, name: &str, opts: UpdateServiceAccountOpts) -> Result { if name == SITE_REPLICATOR_SERVICE_ACCOUNT && !opts.allow_site_replicator_account { return Err(IamError::IAMActionNotAllowed); } let updated_at = self.store.update_service_account(name, opts).await?; self.notify_for_service_account(name).await; Ok(updated_at) } pub async fn list_service_accounts(&self, access_key: &str) -> Result> { self.store.list_service_accounts(access_key).await } pub async fn list_temp_accounts(&self, access_key: &str) -> Result> { self.store.list_temp_accounts(access_key).await } pub async fn list_sts_accounts(&self, access_key: &str) -> Result> { self.store.list_sts_accounts(access_key).await } pub async fn get_service_account(&self, access_key: &str) -> Result<(Credentials, Option)> { let (mut da, policy) = self.get_service_account_internal(access_key).await?; da.credentials.secret_key.clear(); da.credentials.session_token.clear(); Ok((da.credentials, policy)) } pub async fn get_site_replicator_service_account_secret(&self, access_key: &str) -> Result { if access_key != SITE_REPLICATOR_SERVICE_ACCOUNT { return Err(IamError::NoSuchServiceAccount(access_key.to_string())); } let (user, claims) = self.get_account_with_claims_allow_missing_exp(access_key).await?; if !user.credentials.is_service_account() || !claims.get(SITE_REPLICATOR_CLAIM).and_then(Value::as_bool).unwrap_or(false) { return Err(IamError::NoSuchServiceAccount(access_key.to_string())); } Ok(user.credentials.secret_key) } async fn get_service_account_internal(&self, access_key: &str) -> Result<(UserIdentity, Option)> { let (sa, claims) = match self.get_account_with_claims_allow_missing_exp(access_key).await { Ok(res) => res, Err(err) => { if is_err_no_such_account(&err) { return Err(IamError::NoSuchServiceAccount(access_key.to_string())); } return Err(err); } }; if !sa.credentials.is_service_account() { return Err(IamError::NoSuchServiceAccount(access_key.to_string())); } let op_pt = claims.get(&iam_policy_claim_name_sa()); let op_sp = claims.get(SESSION_POLICY_NAME); if let (Some(pt), Some(sp)) = (op_pt, op_sp) && pt == EMBEDDED_POLICY_TYPE { let policy = serde_json::from_slice(&base64_simd::URL_SAFE_NO_PAD.decode_to_vec(sp.as_str().unwrap_or_default().as_bytes())?)?; return Ok((sa, Some(policy))); } Ok((sa, None)) } async fn get_account_with_claims(&self, access_key: &str) -> Result<(UserIdentity, HashMap)> { let Some(acc) = self.store.get_user(access_key).await else { return Err(IamError::NoSuchAccount(access_key.to_string())); }; let m = extract_jwt_claims(&acc)?; Ok((acc, m)) } async fn get_account_with_claims_allow_missing_exp( &self, access_key: &str, ) -> Result<(UserIdentity, HashMap)> { let Some(acc) = self.store.get_user(access_key).await else { return Err(IamError::NoSuchAccount(access_key.to_string())); }; let m = crate::manager::extract_jwt_claims_allow_missing_exp(&acc)?; Ok((acc, m)) } pub async fn get_temporary_account(&self, access_key: &str) -> Result<(Credentials, Option)> { let (mut sa, policy) = match self.get_temp_account(access_key).await { Ok(res) => res, Err(err) => { if is_err_no_such_temp_account(&err) { // TODO: load_user match self.get_temp_account(access_key).await { Ok(res) => res, Err(err) => return Err(err), }; } return Err(err); } }; sa.credentials.secret_key.clear(); sa.credentials.session_token.clear(); Ok((sa.credentials, policy)) } async fn get_temp_account(&self, access_key: &str) -> Result<(UserIdentity, Option)> { let (sa, claims) = match self.get_account_with_claims(access_key).await { Ok(res) => res, Err(err) => { if is_err_no_such_account(&err) { return Err(IamError::NoSuchTempAccount(access_key.to_string())); } return Err(err); } }; if !sa.credentials.is_temp() { return Err(IamError::NoSuchTempAccount(access_key.to_string())); } let op_pt = claims.get(&iam_policy_claim_name_sa()); let op_sp = claims.get(SESSION_POLICY_NAME); if let (Some(pt), Some(sp)) = (op_pt, op_sp) && pt == EMBEDDED_POLICY_TYPE { let policy = serde_json::from_slice(&base64_simd::URL_SAFE_NO_PAD.decode_to_vec(sp.as_str().unwrap_or_default().as_bytes())?)?; return Ok((sa, Some(policy))); } Ok((sa, None)) } pub async fn get_claims_for_svc_acc(&self, access_key: &str) -> Result> { let Some(u) = self.store.get_user(access_key).await else { return Err(IamError::NoSuchServiceAccount(access_key.to_string())); }; if !u.credentials.is_service_account() { return Err(IamError::NoSuchServiceAccount(access_key.to_string())); } crate::manager::extract_jwt_claims_allow_missing_exp(&u) } pub async fn delete_service_account(&self, access_key: &str, notify: bool) -> Result<()> { self.delete_service_account_with_revision(access_key, notify).await?; Ok(()) } pub async fn delete_service_account_with_revision(&self, access_key: &str, notify: bool) -> Result> { let Some(u) = self.store.get_user(access_key).await else { return Ok(None); }; if !u.credentials.is_service_account() { return Ok(None); } let deleted_at = self.store.delete_service_account_with_revision(access_key).await?; if notify && !self.has_watcher() { for result in notify_iam_load_service_account(access_key).await { if let Some(err) = result.err { warn!("notify deleted service account failed: {}", err); } } } Ok(Some(deleted_at)) } async fn notify_for_group(&self, group: &str) { if self.has_watcher() { return; } // Fire-and-forget notification to peers - don't block group operations let group = group.to_string(); tokio::spawn(async move { for r in notify_iam_load_group(&group).await { if let Some(err) = r.err { warn!("notify load_group failed (non-blocking): {}", err); } } }); } pub async fn create_user(&self, access_key: &str, args: &AddOrUpdateUserReq) -> Result { if !is_access_key_valid(access_key) { return Err(IamError::InvalidAccessKeyLength); } if contains_reserved_chars(access_key) { return Err(IamError::ContainsReservedChars); } if !is_secret_key_valid(&args.secret_key) { return Err(IamError::InvalidSecretKeyLength); } let updated_at = self.store.add_user(access_key, args).await?; self.load_user(access_key, UserType::Reg).await?; self.notify_for_user(access_key, false).await; Ok(updated_at) } pub async fn set_user_secret_key(&self, access_key: &str, secret_key: &str) -> Result<()> { if !is_access_key_valid(access_key) { return Err(IamError::InvalidAccessKeyLength); } if !is_secret_key_valid(secret_key) { return Err(IamError::InvalidSecretKeyLength); } self.store.update_user_secret_key(access_key, secret_key).await } /// Add SSH public key for a user (for SFTP authentication) pub async fn add_user_ssh_public_key(&self, access_key: &str, public_key: &str) -> Result<()> { if !is_access_key_valid(access_key) { return Err(IamError::InvalidAccessKeyLength); } if public_key.is_empty() { return Err(IamError::InvalidArgument); } self.store.add_user_ssh_public_key(access_key, public_key).await } pub async fn check_key(&self, access_key: &str) -> Result<(Option, bool)> { if let Some(sys_cred) = crate::root_credentials::credentials() && sys_cred.access_key == access_key { return Ok((Some(UserIdentity::new(sys_cred)), true)); } match self.store.get_user(access_key).await { Some(res) => { let ok = res.credentials.is_valid(); Ok((Some(res), ok)) } None => { self.store.load_user(access_key).await?; if let Some(res) = self.store.get_user(access_key).await { let ok = res.credentials.is_valid(); Ok((Some(res), ok)) } else { Ok((None, false)) } } } } pub async fn get_user(&self, access_key: &str) -> Option { match self.check_key(access_key).await { Ok((u, _)) => u, _ => None, } } pub async fn add_users_to_group(&self, group: &str, users: Vec) -> Result { if contains_reserved_chars(group) { return Err(IamError::GroupNameContainsReservedChars); } let updated_at = self.store.add_users_to_group(group, users).await?; self.notify_for_group(group).await; Ok(updated_at) } pub async fn remove_users_from_group(&self, group: &str, users: Vec) -> Result { let updated_at = self.store.remove_users_from_group(group, users).await?; self.notify_for_group(group).await; Ok(updated_at) } pub async fn set_group_status(&self, group: &str, enable: bool) -> Result { let updated_at = self.store.set_group_status(group, enable).await?; self.notify_for_group(group).await; Ok(updated_at) } pub async fn get_group_description(&self, group: &str) -> Result { self.store.get_group_description(group).await } pub async fn list_groups_load(&self) -> Result> { self.store.update_groups().await } pub async fn list_groups(&self) -> Result> { self.store.list_groups().await } pub async fn policy_db_set(&self, name: &str, user_type: UserType, is_group: bool, policy: &str) -> Result { let updated_at = self.store.policy_db_set(name, user_type, is_group, policy).await?; if !self.has_watcher() { for r in notify_iam_load_policy_mapping(name, user_type.to_u64(), is_group).await { if let Some(err) = r.err { warn!("notify load_policy failed: {}", err); } } } Ok(updated_at) } pub async fn policy_db_get(&self, name: &str, groups: &Option>) -> Result> { self.store.policy_db_get(name, groups).await } pub async fn sts_policy_db_get(&self, name: &str, groups: &Option>) -> Result> { self.store.sts_policy_db_get(name, groups).await } /// Check whether a policy name from a JWT claim is safe to resolve against the IAM store. /// /// Allowed characters: `[a-zA-Z0-9_:.-]` /// - Colons: K8s service account `sub` claims (`system:serviceaccount:ns:sa`) /// - Dots: OIDC group names from providers like Okta/Entra (`org.team.role`) /// /// Requirements: /// - At least one ASCII alphanumeric character (prevents meaningless names /// like `.`, `..`, `-`, `___`, or `:.:`). /// - No characters outside the allowed set (helps mitigate path traversal /// and injection when names are used in storage keys or log output). fn is_safe_claim_policy_name(policy: &str) -> bool { is_safe_claim_policy_name(policy) } // JWT policy claims carry canned policy names only; policy documents are resolved by IAM store. fn safe_claim_policy_names(claims: &HashMap, parent_user: &str) -> Vec { let Some(claim_policies) = claims.get(POLICYNAME).and_then(|v| v.as_str()) else { return Vec::new(); }; claim_policies .split(',') .map(str::trim) .filter(|policy_name| { if policy_name.is_empty() { return false; } if !Self::is_safe_claim_policy_name(policy_name) { tracing::debug!( parent_user = %parent_user, "prepare_sts_auth: ignoring unsafe policy name in STS policy claim" ); return false; } true }) .map(ToOwned::to_owned) .collect() } async fn verified_federated_service_account_claims( &self, claims: &HashMap, parent_user: &str, groups: &Option>, ) -> Result>> { let Some(policy_names) = self.verified_federated_policy_names(claims, parent_user, groups).await? else { return Ok(None); }; let (resolved, _) = self.store.merge_policies(&policy_names.join(",")).await; let resolved_names = MappedPolicy::new(&resolved).to_slice(); if policy_names.iter().any(|policy_name| !resolved_names.contains(policy_name)) { return Err(IamError::InvalidArgument); } let boundary = match claims.get(SESSION_POLICY_NAME) { Some(Value::String(encoded)) => { decode_session_policy(encoded)?; Some(encoded.clone()) } Some(_) => return Err(IamError::InvalidArgument), None => None, }; let mut verified_claims = HashMap::from([ (VERIFIED_FEDERATED_POLICY_CLAIM.to_string(), Value::Bool(true)), (POLICYNAME.to_string(), Value::String(resolved)), ]); if claims.get(OIDC_VIRTUAL_PARENT_CLAIM).and_then(Value::as_str) == Some(parent_user) { verified_claims.insert(OIDC_VIRTUAL_PARENT_CLAIM.to_string(), Value::String(parent_user.to_string())); } if let Some(boundary) = boundary { verified_claims.insert(FEDERATED_POLICY_BOUNDARY_CLAIM.to_string(), Value::String(boundary)); } Ok(Some(verified_claims)) } pub async fn new_service_account_from_caller( &self, parent_user: &str, groups: Option>, mut opts: NewServiceAccountOpts, caller: &Credentials, ) -> Result<(Credentials, OffsetDateTime, HashMap)> { if parent_user != caller.access_key && parent_user != caller.parent_user { return Err(IamError::IAMActionNotAllowed); } let verified_federated_policy = if caller.is_service_account() { None } else { self.verified_federated_service_account_claims(caller.claims_or_empty(), &caller.parent_user, &caller.groups) .await .map_err(|_| IamError::IAMActionNotAllowed)? }; if let Some(source_claims) = caller.claims.as_ref() { merge_derived_service_account_claims(opts.claims.get_or_insert_default(), source_claims, verified_federated_policy)?; } let replication_claims = opts.claims.clone().unwrap_or_default(); let (credentials, updated_at) = self.new_service_account(parent_user, groups, opts).await?; Ok((credentials, updated_at, replication_claims)) } async fn oidc_policy_names( &self, claims: &HashMap, parent_user: &str, groups: &Option>, ) -> Result> { let mapped_policy_names = self.policy_db_get(parent_user, groups).await?; let policy_names = if mapped_policy_names.is_empty() { Self::safe_claim_policy_names(claims, parent_user) } else { mapped_policy_names }; if policy_names.is_empty() { return Err(IamError::InvalidArgument); } Ok(policy_names) } fn signed_oidc_policy_names(claims: &HashMap, parent_user: &str) -> Result> { let policy_names = Self::safe_claim_policy_names(claims, parent_user); if policy_names.is_empty() { return Err(IamError::InvalidArgument); } Ok(policy_names) } async fn verified_federated_policy_names( &self, claims: &HashMap, parent_user: &str, groups: &Option>, ) -> Result>> { if !is_rustfs_oidc_claims(claims) { return Ok(None); } if !is_verified_or_legacy_federated_session(claims, parent_user) { return Err(IamError::InvalidArgument); } let policy_names = if is_verified_federated_session(claims, parent_user) { Self::signed_oidc_policy_names(claims, parent_user)? } else { self.oidc_policy_names(claims, parent_user, groups).await? }; Ok(Some(policy_names)) } /// Compatibility wrapper for service-account authorization entry points. /// The canonical evaluation path is `prepare_service_account_auth + eval_prepared`. pub async fn is_allowed_service_account(&self, args: &Args<'_>, parent_user: &str) -> bool { let prepared = self.prepare_service_account_auth(args, parent_user).await; self.eval_prepared(&prepared, args).await } pub async fn get_combined_policy(&self, policies: &[String]) -> Policy { self.store.merge_policies(&policies.join(",")).await.1 } /// Prepare IAM authorization context once so callers can: /// 1) know whether policy evaluation may need `s3:ExistingObjectTag`, and /// 2) evaluate with final conditions without re-merging identity policies. pub async fn prepare_auth(&self, args: &Args<'_>) -> PreparedIamAuth { if args.is_owner { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Owner, }; } match Self::policy_plugin_state().await { PolicyPluginState::Ready(_) => { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Opa, }; } PolicyPluginState::Initializing | PolicyPluginState::Failed => { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } PolicyPluginState::Disabled => {} } let Ok((is_svc, parent_user)) = self.is_service_account(args.account).await else { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; }; if is_svc { return self.prepare_service_account_auth(args, &parent_user).await; } let Ok((is_temp, parent_user)) = self.is_temp_user(args.account).await else { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; }; if is_temp { return self.prepare_sts_auth(args, &parent_user).await; } self.prepare_regular_auth(args).await } pub async fn eval_prepared(&self, prepared: &PreparedIamAuth, args: &Args<'_>) -> bool { match &prepared.mode { PreparedIamMode::Opa => { let Some(opa_enable) = Self::get_policy_plugin_client().await else { tracing::warn!("eval_prepared: OPA mode requested but plugin is unavailable"); return false; }; opa_enable.is_allowed(args).await } PreparedIamMode::Owner => true, PreparedIamMode::Deny => false, PreparedIamMode::Regular { combined_policy } => combined_policy.is_allowed(args).await, PreparedIamMode::Sts { is_owner, combined_policy, session_policy, } => { let session_ok = evaluate_prepared_session_policy(session_policy, args).await; if let Some(ok) = session_ok { return ok && (*is_owner || combined_policy.is_allowed(args).await); } *is_owner || combined_policy.is_allowed(args).await } PreparedIamMode::ServiceAccount { is_owner, bypass_parent_policy, parent_user, combined_policy, federated_boundary, mode, session_policy, } => { let mut parent_args = args.clone(); parent_args.account = parent_user; let boundary_allowed = if let Some(policy) = federated_boundary { let mut boundary_args = args.clone(); boundary_args.is_owner = false; policy.is_allowed(&boundary_args).await } else { true }; let parent_allowed = (*bypass_parent_policy || *is_owner || combined_policy.is_allowed(&parent_args).await) && boundary_allowed; match mode { PreparedServicePolicyMode::Inherited => parent_allowed, PreparedServicePolicyMode::SessionBound => { let session_ok = evaluate_prepared_session_policy(session_policy, args).await; if let Some(ok) = session_ok { return ok && parent_allowed; } parent_allowed } } } } } async fn prepare_regular_auth(&self, args: &Args<'_>) -> PreparedIamAuth { let Ok(policies) = self.policy_db_get(args.account, args.groups).await else { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; }; if policies.is_empty() { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } let (resolved_policies, combined_policy) = self.store.merge_policies(&policies.join(",")).await; if args.deny_only { let resolved_policy_names: HashSet<&str> = resolved_policies .split(',') .filter(|policy_name| !policy_name.trim().is_empty()) .collect(); if policies .iter() .any(|policy_name| !resolved_policy_names.contains(policy_name.as_str())) { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } } PreparedIamAuth { needs_existing_object_tag: policy_needs_existing_object_tag_for_args(&combined_policy, args).await, mode: PreparedIamMode::Regular { combined_policy }, } } pub(crate) async fn prepare_sts_auth(&self, args: &Args<'_>, parent_user: &str) -> PreparedIamAuth { let federated_policy_names = match self .verified_federated_policy_names(args.claims, parent_user, args.groups) .await { Ok(policy_names) => policy_names, Err(_) => { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } }; let federated_session = federated_policy_names.is_some(); let verified_federated_session = is_verified_federated_session(args.claims, parent_user); let is_owner = !is_rustfs_oidc_claims(args.claims) && crate::is_root_access_key(parent_user); let role_arn = args.get_role_arn(); let (effective_groups, groups_source, policies) = if federated_session { (args.groups.clone(), "oidc_claims", Vec::new()) } else if is_owner { (None, "owner", Vec::new()) } else if let Some(arn_str) = role_arn { let Ok(arn) = ARN::parse(arn_str) else { tracing::warn!( parent_user = %parent_user, role_arn = %arn_str, "prepare_sts_auth: invalid role ARN in STS claims" ); return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; }; let p = MappedPolicy::new(self.roles_map.get(&arn).map_or_else(String::default, |v| v.clone()).as_str()).to_slice(); (None, "role", p) } else { let (effective_groups, groups_source) = match args.groups.as_ref() { Some(g) if !g.is_empty() => (args.groups.clone(), "args"), _ => match self.store.get_user(parent_user).await { Some(u) => (u.credentials.groups, "parent_user_credentials"), None => { tracing::warn!( parent_user = %parent_user, "prepare_sts_auth: groups fallback failed, parent user not found" ); (None, "parent_user_credentials") } }, }; let p = self.policy_db_get(parent_user, &effective_groups).await.unwrap_or_default(); (effective_groups, groups_source, p) }; let mut policy_names = if let Some(policy_names) = federated_policy_names { policy_names } else { policies }; if !is_owner && policy_names.is_empty() && !is_rustfs_oidc_claims(args.claims) { policy_names = Self::safe_claim_policy_names(args.claims, parent_user); } let mut combined_policy = Policy::default(); if !is_owner { if policy_names.is_empty() { if args.deny_only { combined_policy = Policy::default(); } else { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } } else { let requested_policies = policy_names.join(","); let (resolved_policies, c) = self.store.merge_policies(&requested_policies).await; if resolved_policies.is_empty() { tracing::warn!( parent_user = %parent_user, requested_policies = %requested_policies, "prepare_sts_auth: no STS policy names resolved" ); if args.deny_only && !verified_federated_session { combined_policy = Policy::default(); } else { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } } else { let resolved_policy_names = MappedPolicy::new(&resolved_policies).to_slice(); let has_unresolved_policy_names = policy_names .iter() .any(|policy_name| !resolved_policy_names.iter().any(|resolved| resolved == policy_name)); if has_unresolved_policy_names { if verified_federated_session { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } tracing::debug!( parent_user = %parent_user, requested_policies = %requested_policies, resolved_policies = %resolved_policies, "prepare_sts_auth: some STS policy names were not resolved" ); } combined_policy = c; } } } let session_policy = prepare_session_policy(args, false); tracing::debug!( "prepare_sts_auth: action={:?}, is_owner={}, parent_user={}, groups_source={}, effective_groups={:?}", args.action, is_owner, parent_user, groups_source, effective_groups ); PreparedIamAuth { needs_existing_object_tag: policy_needs_existing_object_tag_for_args(&combined_policy, args).await || prepared_session_policy_needs_existing_object_tag_for_args(&session_policy, args).await, mode: PreparedIamMode::Sts { is_owner, combined_policy, session_policy, }, } } async fn prepare_service_account_auth(&self, args: &Args<'_>, parent_user: &str) -> PreparedIamAuth { let Some(p) = args.claims.get("parent") else { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; }; if p.as_str() != Some(parent_user) { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } let Some(sa) = args.claims.get(&iam_policy_claim_name_sa()) else { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; }; let Some(sa_str) = sa.as_str() else { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; }; let mode = if sa_str == INHERITED_POLICY_TYPE { PreparedServicePolicyMode::Inherited } else { PreparedServicePolicyMode::SessionBound }; let session_policy = prepare_session_policy(args, true); let has_verified_federated_policy = args.claims.contains_key(VERIFIED_FEDERATED_POLICY_CLAIM); if is_rustfs_oidc_claims(args.claims) && !has_verified_federated_policy { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } let federated_boundary = if has_verified_federated_policy { if !is_valid_verified_federated_policy(args.claims, parent_user) { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } match federated_policy_boundary(args.claims) { Ok(boundary) => boundary, Err(_) => { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } } } else { None }; let bypass_parent_policy = args.account == SITE_REPLICATOR_SERVICE_ACCOUNT && args .claims .get(SITE_REPLICATOR_CLAIM) .and_then(Value::as_bool) .unwrap_or(false) && matches!(mode, PreparedServicePolicyMode::SessionBound) && matches!(session_policy, PreparedSessionPolicy::Policy(_)); let is_owner = !is_rustfs_oidc_claims(args.claims) && crate::is_root_access_key(parent_user); let role_arn = args.get_role_arn(); let svc_policies = if is_owner || bypass_parent_policy { Vec::new() } else if has_verified_federated_policy { Self::safe_claim_policy_names(args.claims, parent_user) } else if role_arn.is_some() { let Ok(arn) = ARN::parse(role_arn.unwrap_or_default()) else { tracing::warn!( parent_user = %parent_user, role_arn = ?role_arn, "prepare_service_account_auth: invalid role ARN in service account claims" ); return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; }; MappedPolicy::new(self.roles_map.get(&arn).map_or_else(String::default, |v| v.clone()).as_str()).to_slice() } else { let Ok(policies) = self.policy_db_get(parent_user, args.groups).await else { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; }; policies }; if !is_owner && !bypass_parent_policy && svc_policies.is_empty() { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } let combined_policy = if is_owner || bypass_parent_policy { Policy::default() } else { let (a, c) = self.store.merge_policies(&svc_policies.join(",")).await; let resolved = MappedPolicy::new(&a).to_slice(); if a.is_empty() || has_verified_federated_policy && svc_policies.iter().any(|policy_name| !resolved.contains(policy_name)) { return PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; } c }; let needs_existing_object_tag = policy_needs_existing_object_tag_for_args(&combined_policy, args).await || match &federated_boundary { Some(policy) => policy_needs_existing_object_tag_for_args(policy, args).await, None => false, } || matches!(mode, PreparedServicePolicyMode::SessionBound) && prepared_session_policy_needs_existing_object_tag_for_args(&session_policy, args).await; PreparedIamAuth { needs_existing_object_tag, mode: PreparedIamMode::ServiceAccount { is_owner, bypass_parent_policy, parent_user: parent_user.to_string(), combined_policy, federated_boundary, mode, session_policy, }, } } pub async fn is_allowed(&self, args: &Args<'_>) -> bool { let prepared = self.prepare_auth(args).await; self.eval_prepared(&prepared, args).await } /// Check if the underlying store is ready pub fn is_ready(&self) -> bool { self.store.is_ready() } } async fn prepared_session_policy_needs_existing_object_tag_for_args(policy: &PreparedSessionPolicy, args: &Args<'_>) -> bool { match policy { PreparedSessionPolicy::Policy(p) => policy_needs_existing_object_tag_for_args(p, args).await, PreparedSessionPolicy::None | PreparedSessionPolicy::DenyAll => false, } } fn prepare_session_policy(args: &Args<'_>, empty_is_none: bool) -> PreparedSessionPolicy { let Some(policy_str) = extract_session_policy_text(args.claims) else { return PreparedSessionPolicy::None; }; let Ok(sub_policy) = Policy::parse_config(policy_str.as_bytes()) else { return PreparedSessionPolicy::DenyAll; }; if empty_is_none { if sub_policy.version.is_empty() && sub_policy.statements.is_empty() && sub_policy.id.is_empty() { return PreparedSessionPolicy::None; } return PreparedSessionPolicy::Policy(sub_policy); } if sub_policy.version.is_empty() { return PreparedSessionPolicy::DenyAll; } PreparedSessionPolicy::Policy(sub_policy) } fn decode_session_policy(encoded: &str) -> Result { let bytes = base64_simd::URL_SAFE_NO_PAD.decode_to_vec(encoded.as_bytes())?; if bytes.len() > MAX_SVCSESSION_POLICY_SIZE { return Err(IamError::PolicyTooLarge); } let policy = Policy::parse_config(&bytes)?; policy.validate()?; if policy.version.is_empty() { return Err(IamError::InvalidArgument); } Ok(policy) } pub fn is_safe_claim_policy_name(policy: &str) -> bool { let mut has_alnum = false; for c in policy.bytes() { if c.is_ascii_alphanumeric() { has_alnum = true; } else if !matches!(c, b'_' | b'-' | b':' | b'.') { return false; } } has_alnum } pub fn is_rustfs_oidc_claims(claims: &HashMap) -> bool { claims.get("iss").and_then(Value::as_str) == Some("rustfs-oidc") && claims .get("oidc_provider") .and_then(Value::as_str) .is_some_and(|provider| !provider.is_empty()) && claims .get("sub") .and_then(Value::as_str) .is_some_and(|subject| !subject.is_empty()) } fn has_verified_federated_policy(claims: &HashMap) -> bool { claims .get(VERIFIED_FEDERATED_POLICY_CLAIM) .and_then(Value::as_bool) .unwrap_or(false) } pub fn remove_verified_federated_policy(claims: &mut HashMap) -> bool { let removed = claims.remove(VERIFIED_FEDERATED_POLICY_CLAIM).is_some(); claims.remove(FEDERATED_POLICY_BOUNDARY_CLAIM); claims.remove(OIDC_VIRTUAL_PARENT_CLAIM); removed } fn merge_derived_service_account_claims( target_claims: &mut HashMap, source_claims: &HashMap, verified_federated_policy: Option>, ) -> Result<()> { for (key, value) in source_claims { if key != "exp" { target_claims.insert(key.clone(), value.clone()); } } let inherited_verified_policy = remove_verified_federated_policy(target_claims); let Some(verified_policy) = verified_federated_policy else { return if inherited_verified_policy { Err(IamError::IAMActionNotAllowed) } else { Ok(()) }; }; target_claims.remove(SESSION_POLICY_NAME); target_claims.remove(SESSION_POLICY_NAME_EXTRACTED); target_claims.extend(verified_policy); Ok(()) } fn is_verified_federated_session(claims: &HashMap, parent_user: &str) -> bool { !parent_user.is_empty() && is_rustfs_oidc_claims(claims) && claims.get("parent").and_then(Value::as_str) == Some(parent_user) && claims.get(OIDC_VIRTUAL_PARENT_CLAIM).and_then(Value::as_str) == Some(parent_user) } fn is_verified_or_legacy_federated_session(claims: &HashMap, parent_user: &str) -> bool { is_verified_federated_session(claims, parent_user) || (!claims.contains_key(OIDC_VIRTUAL_PARENT_CLAIM) && !parent_user.is_empty() && is_rustfs_oidc_claims(claims) && claims.get("parent").and_then(Value::as_str) == Some(parent_user)) } fn is_valid_verified_federated_policy(claims: &HashMap, parent_user: &str) -> bool { has_verified_federated_policy(claims) && is_verified_or_legacy_federated_session(claims, parent_user) && claims.get(POLICYNAME).and_then(Value::as_str).is_some_and(|policies| { policies .split(',') .map(str::trim) .all(|policy| !policy.is_empty() && is_safe_claim_policy_name(policy)) }) } fn federated_policy_boundary(claims: &HashMap) -> Result> { let Some(boundary) = claims.get(FEDERATED_POLICY_BOUNDARY_CLAIM) else { return Ok(None); }; boundary .as_str() .ok_or(IamError::InvalidArgument) .and_then(decode_session_policy) .map(Some) } fn extract_session_policy_text(claims: &HashMap) -> Option { if let Some(policy_str) = claims.get(SESSION_POLICY_NAME_EXTRACTED).and_then(|v| v.as_str()) { return Some(policy_str.to_string()); } let encoded = claims.get(SESSION_POLICY_NAME).and_then(|v| v.as_str())?; let bytes = base64_simd::URL_SAFE_NO_PAD.decode_to_vec(encoded.as_bytes()).ok()?; String::from_utf8(bytes).ok() } async fn evaluate_prepared_session_policy(policy: &PreparedSessionPolicy, args: &Args<'_>) -> Option { match policy { PreparedSessionPolicy::None => None, PreparedSessionPolicy::DenyAll => Some(false), PreparedSessionPolicy::Policy(p) => { let mut session_policy_args = args.clone(); session_policy_args.is_owner = false; Some(p.is_allowed(&session_policy_args).await) } } } #[derive(Debug, Clone, Default)] pub struct NewServiceAccountOpts { pub session_policy: Option, pub access_key: String, pub secret_key: String, pub name: Option, pub description: Option, pub expiration: Option, pub allow_site_replicator_account: bool, pub claims: Option>, } pub struct UpdateServiceAccountOpts { pub session_policy: Option, pub secret_key: Option, pub name: Option, pub description: Option, pub expiration: Option, pub status: Option, /// Rebind the account to a different parent. /// /// Only site replication sets this, and only to repair an account left pointing at a /// parent that a root-credential change invalidated. Rebinding an arbitrary service /// account would let it inherit another user's policies, so it is gated behind /// `allow_site_replicator_account` in the same way the account itself is. pub parent_user: Option, pub allow_site_replicator_account: bool, } pub fn get_claims_from_token_with_secret(token: &str, secret: &str) -> Result> { let mut ms = extract_claims::>(token, secret).map_err(|e| Error::other(format!("extract claims err {e}")))?; if let Some(session_policy) = ms.claims.get(SESSION_POLICY_NAME) { let policy_str = session_policy.as_str().unwrap_or_default(); let policy = base64_simd::URL_SAFE_NO_PAD .decode_to_vec(policy_str.as_bytes()) .map_err(|e| Error::other(format!("base64 decode err {e}")))?; ms.claims.insert( SESSION_POLICY_NAME_EXTRACTED.to_string(), Value::String(String::from_utf8(policy).map_err(|e| Error::other(format!("utf8 decode err {e}")))?), ); } Ok(ms.claims) } pub fn get_claims_from_token_with_secret_allow_missing_exp(token: &str, secret: &str) -> Result> { let mut ms = extract_claims_allow_missing_exp::>(token, secret) .map_err(|e| Error::other(format!("extract claims err {e}")))?; if let Some(session_policy) = ms.claims.get(SESSION_POLICY_NAME) { let policy_str = session_policy.as_str().unwrap_or_default(); let policy = base64_simd::URL_SAFE_NO_PAD .decode_to_vec(policy_str.as_bytes()) .map_err(|e| Error::other(format!("base64 decode err {e}")))?; ms.claims.insert( SESSION_POLICY_NAME_EXTRACTED.to_string(), Value::String(String::from_utf8(policy).map_err(|e| Error::other(format!("utf8 decode err {e}")))?), ); } Ok(ms.claims) } #[cfg(test)] mod tests { use super::*; use crate::cache::{Cache, CacheEntity}; use crate::error::Error; use crate::manager::get_default_policies; use crate::store::{GroupInfo, MappedPolicy, Store, UserType}; use rustfs_credentials::{Credentials, init_global_action_credentials}; use rustfs_policy::auth::{UserIdentity, get_new_credentials_with_metadata}; use rustfs_policy::policy::Args; #[test] #[allow(deprecated)] fn deprecated_list_polices_api_is_available() { let _ = IamSys::::list_polices; } use rustfs_policy::policy::action::{Action, AdminAction, S3Action, StsAction}; use rustfs_policy::policy::policy_uses_existing_object_tag_conditions; use serde_json::Value; use std::{ collections::{HashMap, HashSet}, sync::{Arc, Mutex}, }; use time::OffsetDateTime; #[test] fn test_combined_policy_for_view_returns_regular_policy() { let policy = Policy { version: "2012-10-17".to_string(), ..Default::default() }; let prepared = PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Regular { combined_policy: policy }, }; let resolved = prepared.combined_policy_for_view(); assert_eq!(resolved.map(|p| p.version.as_str()), Some("2012-10-17")); } #[test] fn test_combined_policy_for_view_returns_none_for_deny() { let prepared = PreparedIamAuth { needs_existing_object_tag: false, mode: PreparedIamMode::Deny, }; assert!(prepared.combined_policy_for_view().is_none()); } const CUSTOM_STS_CLAIM_POLICY: &str = "custom-sts-claim-getobject"; const CUSTOM_STS_CLAIM_BUCKET: &str = "claim-bucket"; const CUSTOM_STS_CLAIM_POLICY_JSON: &str = r#"{ "Version": "2012-10-17", "Statement": [ { "Effect": "Allow", "Action": ["s3:GetObject"], "Resource": ["arn:aws:s3:::claim-bucket/allowed/*"] } ] }"#; /// Mock Store for STS tests: either group-attached policies via parent user, or no IAM policies. #[derive(Clone)] struct StsTestMockStore { /// When true, parent user has no groups and no mapped policies (empty `policy_db_get`). empty_policies: bool, saved_sts_users: Arc>>, saved_service_account_count: Arc>, } impl StsTestMockStore { fn new(empty_policies: bool) -> Self { Self { empty_policies, saved_sts_users: Arc::new(Mutex::new(HashMap::new())), saved_service_account_count: Arc::new(Mutex::new(0)), } } fn saved_service_account_count(&self) -> usize { *self .saved_service_account_count .lock() .expect("saved_service_account_count mutex poisoned") } } #[async_trait::async_trait] impl Store for StsTestMockStore { fn has_watcher(&self) -> bool { false } async fn save_iam_config(&self, _item: Item, _path: impl AsRef + Send) -> Result<()> { Ok(()) } async fn load_iam_config(&self, _path: impl AsRef + Send) -> Result { Err(Error::ConfigNotFound) } async fn delete_iam_config(&self, _path: impl AsRef + Send) -> Result<()> { Err(Error::InvalidArgument) } async fn save_user_identity( &self, name: &str, user_type: UserType, item: UserIdentity, _ttl: Option, ) -> Result<()> { if user_type == UserType::Svc { *self .saved_service_account_count .lock() .expect("saved_service_account_count mutex poisoned") += 1; } self.saved_sts_users .lock() .expect("saved_sts_users mutex poisoned") .insert(name.to_string(), item); Ok(()) } async fn delete_user_identity(&self, name: &str, _user_type: UserType) -> Result<()> { self.saved_sts_users .lock() .expect("saved_sts_users mutex poisoned") .remove(name); Ok(()) } async fn load_user_identity(&self, name: &str, _user_type: UserType) -> Result { self.saved_sts_users .lock() .expect("saved_sts_users mutex poisoned") .get(name) .cloned() .ok_or_else(|| Error::NoSuchUser(name.to_string())) } async fn load_user(&self, name: &str, user_type: UserType, m: &mut HashMap) -> Result<()> { if name == "deleted-notify-user" { return Err(Error::NoSuchUser(name.to_string())); } if user_type == UserType::Reg && name == "load-failure-user" { return Err(Error::Io(std::io::Error::other("load user failed"))); } if user_type == UserType::Reg && name == "notify-user" { let user = UserIdentity::from(Credentials { access_key: name.to_string(), secret_key: "notify-user-secret".to_string(), status: ACCOUNT_ON.to_string(), ..Default::default() }); m.insert(name.to_string(), user); } Ok(()) } async fn load_users(&self, _user_type: UserType, _m: &mut HashMap) -> Result<()> { Ok(()) } async fn load_secret_key(&self, _name: &str, _user_type: UserType) -> Result { Err(Error::InvalidArgument) } async fn save_group_info(&self, _name: &str, _item: GroupInfo) -> Result<()> { Err(Error::InvalidArgument) } async fn delete_group_info(&self, _name: &str) -> Result<()> { Err(Error::InvalidArgument) } async fn load_group(&self, name: &str, m: &mut HashMap) -> Result<()> { if name == "notify-group" { m.insert(name.to_string(), GroupInfo::new(vec!["notify-user".to_string()])); } Ok(()) } async fn load_groups(&self, _m: &mut HashMap) -> Result<()> { Err(Error::InvalidArgument) } async fn save_policy_doc(&self, _name: &str, _item: rustfs_policy::policy::PolicyDoc) -> Result<()> { Err(Error::InvalidArgument) } async fn delete_policy_doc(&self, _name: &str) -> Result<()> { Err(Error::InvalidArgument) } async fn load_policy(&self, _name: &str) -> Result { Err(Error::InvalidArgument) } async fn load_policy_doc(&self, _name: &str, _m: &mut HashMap) -> Result<()> { Err(Error::InvalidArgument) } async fn load_policy_docs(&self, _m: &mut HashMap) -> Result<()> { Err(Error::InvalidArgument) } async fn save_mapped_policy( &self, _name: &str, _user_type: UserType, _is_group: bool, _item: MappedPolicy, _ttl: Option, ) -> Result<()> { Err(Error::InvalidArgument) } async fn delete_mapped_policy(&self, _name: &str, _user_type: UserType, _is_group: bool) -> Result<()> { Err(Error::InvalidArgument) } async fn load_mapped_policy( &self, name: &str, user_type: UserType, is_group: bool, m: &mut HashMap, ) -> Result<()> { if user_type == UserType::Reg && !is_group && name == "notify-user" { m.insert(name.to_string(), MappedPolicy::new("readwrite")); } if user_type == UserType::Sts && !is_group && name == "notify-sts-parent" { m.insert(name.to_string(), MappedPolicy::new("readwrite")); } Ok(()) } async fn load_mapped_policies( &self, _user_type: UserType, _is_group: bool, _m: &mut HashMap, ) -> Result<()> { Ok(()) } async fn load_all(&self, cache: &Cache) -> Result<()> { let mut policy_docs = get_default_policies(); let custom_claim_policy = Policy::parse_config(CUSTOM_STS_CLAIM_POLICY_JSON.as_bytes()).expect("custom STS claim policy should parse"); policy_docs.insert(CUSTOM_STS_CLAIM_POLICY.to_string(), PolicyDoc::new(custom_claim_policy)); if self.empty_policies { const PARENT_USER: &str = "sts-empty-parent-policy-test"; let creds = Credentials { access_key: PARENT_USER.to_string(), secret_key: "longenoughsecret".to_string(), session_token: String::new(), expiration: None, status: "on".to_string(), parent_user: String::new(), groups: None, claims: None, name: None, description: None, }; let parent_identity = UserIdentity { version: 1, credentials: creds, update_at: Some(OffsetDateTime::now_utc()), }; let mut users = HashMap::new(); users.insert(PARENT_USER.to_string(), parent_identity); cache.with_write_lock(|cache| { cache.replace_policy_docs(CacheEntity::new(policy_docs)); cache.replace_users(CacheEntity::new(users)); cache.replace_groups(CacheEntity::default()); cache.replace_group_policies(CacheEntity::default()); cache.replace_user_policies(CacheEntity::default()); cache.replace_sts_accounts(CacheEntity::default()); cache.replace_sts_policies(CacheEntity::default()); cache.build_user_group_memberships(); }); return Ok(()); } const PARENT_USER: &str = "sts-fallback-test-parent"; const GROUP_NAME: &str = "testgroup"; let creds = Credentials { access_key: PARENT_USER.to_string(), secret_key: "longenoughsecret".to_string(), session_token: String::new(), expiration: None, status: "on".to_string(), parent_user: String::new(), groups: Some(vec![GROUP_NAME.to_string()]), claims: None, name: None, description: None, }; let parent_identity = UserIdentity { version: 1, credentials: creds, update_at: Some(OffsetDateTime::now_utc()), }; let mut users = HashMap::new(); users.insert(PARENT_USER.to_string(), parent_identity); let group = GroupInfo::new(vec![PARENT_USER.to_string()]); let mut groups = HashMap::new(); groups.insert(GROUP_NAME.to_string(), group); let group_policy = MappedPolicy::new("readwrite"); let mut group_policies = HashMap::new(); group_policies.insert(GROUP_NAME.to_string(), group_policy); cache.with_write_lock(|cache| { cache.replace_policy_docs(CacheEntity::new(policy_docs)); cache.replace_users(CacheEntity::new(users)); cache.replace_groups(CacheEntity::new(groups)); cache.replace_group_policies(CacheEntity::new(group_policies)); cache.replace_user_policies(CacheEntity::default()); cache.replace_sts_accounts(CacheEntity::default()); cache.replace_sts_policies(CacheEntity::default()); cache.build_user_group_memberships(); }); Ok(()) } } fn ensure_test_global_credentials() { if crate::root_credentials::credentials().is_none() { let _ = init_global_action_credentials(Some("TESTROOTACCESSKEY".to_string()), Some("TESTROOTSECRET123".to_string())); } } async fn test_iam_sys() -> IamSys { let store = StsTestMockStore::new(false); let cache = IamCache::new(store).await.expect("IAM cache should initialize"); IamSys::new(cache) } fn service_account_opts(access_key: &str, secret_key: &str) -> NewServiceAccountOpts { NewServiceAccountOpts { access_key: access_key.to_string(), secret_key: secret_key.to_string(), ..Default::default() } } #[tokio::test] async fn new_service_account_rejects_parent_access_key() { ensure_test_global_credentials(); let iam_sys = test_iam_sys().await; let access_key = "PARENTACCESSKEY123"; let err = iam_sys .new_service_account(access_key, None, service_account_opts(access_key, "parentAccessKeySecret123")) .await .expect_err("a service account must not reuse its parent access key"); assert_eq!(err, Error::AccessKeyAlreadyExists); } #[tokio::test] async fn service_account_mutations_return_the_persisted_replication_state() { ensure_test_global_credentials(); let iam_sys = test_iam_sys().await; let access_key = "REPLICATEDSTATE123"; iam_sys .new_service_account("parent-user", None, service_account_opts(access_key, "replicatedStateSecret123")) .await .expect("create service account"); let session_policy = Policy::parse_config(CUSTOM_STS_CLAIM_POLICY_JSON.as_bytes()).expect("parse service-account policy"); iam_sys .update_service_account( access_key, UpdateServiceAccountOpts { session_policy: Some(session_policy), secret_key: None, name: None, description: None, expiration: None, status: None, parent_user: None, allow_site_replicator_account: false, }, ) .await .expect("update service account"); let (updated, session_policy) = iam_sys .get_service_account(access_key) .await .expect("updated service account"); assert!(session_policy.is_some_and(|policy| !policy.statements.is_empty())); assert_eq!( updated .claims_or_empty() .get(&iam_policy_claim_name_sa()) .and_then(Value::as_str), Some(EMBEDDED_POLICY_TYPE) ); iam_sys .update_service_account( access_key, UpdateServiceAccountOpts { session_policy: Some(Policy::default()), secret_key: None, name: None, description: None, expiration: None, status: None, parent_user: None, allow_site_replicator_account: false, }, ) .await .expect("clear service account session policy"); let (cleared, session_policy) = iam_sys .get_service_account(access_key) .await .expect("cleared service account"); assert!(session_policy.is_none()); assert_eq!( cleared .claims_or_empty() .get(&iam_policy_claim_name_sa()) .and_then(Value::as_str), Some(INHERITED_POLICY_TYPE) ); let deleted_at = iam_sys .delete_service_account_with_revision(access_key, false) .await .expect("delete service account") .expect("deleted revision"); let (_, recreated_at) = iam_sys .new_service_account("parent-user", None, service_account_opts(access_key, "recreatedStateSecret123")) .await .expect("recreate service account"); assert!(deleted_at <= recreated_at); } #[tokio::test] async fn duplicate_service_account_returns_access_key_already_exists() { ensure_test_global_credentials(); let iam_sys = test_iam_sys().await; let access_key = "DUPLICATESERVICEKEY1"; iam_sys .new_service_account("svc-parent-user", None, service_account_opts(access_key, "firstServiceSecret123")) .await .expect("initial service account should be created"); let err = iam_sys .new_service_account("svc-parent-user", None, service_account_opts(access_key, "secondServiceSecret123")) .await .expect_err("duplicate service account should fail"); assert_eq!(err, Error::AccessKeyAlreadyExists); } #[tokio::test] async fn service_account_creation_rejects_regular_user_access_key() { ensure_test_global_credentials(); let iam_sys = test_iam_sys().await; let access_key = "REGULARUSERACCESSKEY1"; let user = AddOrUpdateUserReq { secret_key: "regularUserSecret123".to_string(), policy: None, status: rustfs_madmin::AccountStatus::Enabled, }; iam_sys .store .add_user(access_key, &user) .await .expect("regular user should be created"); let err = iam_sys .new_service_account("svc-parent-user", None, service_account_opts(access_key, "serviceAccountSecret123")) .await .expect_err("service-account creation must reject a regular-user access key"); assert_eq!(err, Error::AccessKeyAlreadyExists); assert_eq!(iam_sys.store.api.saved_service_account_count(), 0); let existing = iam_sys.get_user(access_key).await.expect("regular user should remain cached"); assert!(!existing.credentials.is_service_account()); assert_eq!(existing.credentials.secret_key, user.secret_key); } #[tokio::test] async fn service_account_creation_rejects_sts_access_key() { ensure_test_global_credentials(); let iam_sys = test_iam_sys().await; let access_key = "TEMPORARYSERVICEKEY12"; let temp_cred = Credentials { access_key: access_key.to_string(), secret_key: "temporarySecretKey123".to_string(), session_token: "temporary-session-token".to_string(), expiration: Some(OffsetDateTime::now_utc() + time::Duration::hours(1)), status: ACCOUNT_ON.to_string(), parent_user: "sts-parent-user".to_string(), ..Default::default() }; iam_sys .set_temp_user(access_key, &temp_cred, None) .await .expect("temporary credentials should be created"); let err = iam_sys .new_service_account("svc-parent-user", None, service_account_opts(access_key, "serviceAccountSecret123")) .await .expect_err("service-account creation must reject a temporary access key"); assert_eq!(err, Error::AccessKeyAlreadyExists); assert_eq!(iam_sys.store.api.saved_service_account_count(), 0); let existing = iam_sys .get_user(access_key) .await .expect("temporary account should remain cached"); assert!(existing.credentials.is_temp()); assert_eq!(existing.credentials.secret_key, temp_cred.secret_key); } #[tokio::test] async fn user_creation_rejects_service_account_access_key() { ensure_test_global_credentials(); let iam_sys = test_iam_sys().await; let access_key = "SERVICEACCOUNTKEY123"; iam_sys .new_service_account("svc-parent-user", None, service_account_opts(access_key, "serviceAccountSecret123")) .await .expect("service account should be created"); let user = AddOrUpdateUserReq { secret_key: "regularUserSecret123".to_string(), policy: None, status: rustfs_madmin::AccountStatus::Enabled, }; let err = iam_sys .create_user(access_key, &user) .await .expect_err("user creation must reject a service-account access key"); assert_eq!(err, Error::AccessKeyAlreadyExists); } #[tokio::test] async fn test_new_service_account_without_expiration_omits_exp_claim() { ensure_test_global_credentials(); let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let (cred, _) = iam_sys .new_service_account("svc-parent-user", None, NewServiceAccountOpts::default()) .await .expect("service account should be created without expiration"); assert!(cred.expiration.is_none()); let claims = get_claims_from_token_with_secret_allow_missing_exp(&cred.session_token, &cred.secret_key) .expect("service account JWT without expiration should decode"); assert!( !claims.contains_key("exp"), "service account without explicit expiration should not get a default JWT exp" ); } #[tokio::test] async fn test_update_service_account_updates_exp_claim() { ensure_test_global_credentials(); let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let initial_expiration = OffsetDateTime::now_utc() + time::Duration::hours(2); let (cred, _) = iam_sys .new_service_account( "svc-parent-user", None, NewServiceAccountOpts { expiration: Some(initial_expiration), ..Default::default() }, ) .await .expect("service account with explicit expiration should be created"); let updated_expiration = OffsetDateTime::now_utc() + time::Duration::hours(4); iam_sys .update_service_account( &cred.access_key, UpdateServiceAccountOpts { session_policy: None, secret_key: None, name: None, description: None, expiration: Some(updated_expiration), status: None, parent_user: None, allow_site_replicator_account: false, }, ) .await .expect("service account expiration should update"); let updated_user = iam_sys .get_user(&cred.access_key) .await .expect("updated service account should exist"); assert_eq!(updated_user.credentials.expiration, Some(updated_expiration)); let claims = get_claims_from_token_with_secret(&updated_user.credentials.session_token, &updated_user.credentials.secret_key) .expect("updated service account JWT should decode"); assert_eq!( claims.get("exp").and_then(|v| v.as_i64()), Some(updated_expiration.unix_timestamp()), "updating service account expiration must rewrite the JWT exp claim" ); } #[tokio::test] async fn test_update_service_account_adds_exp_claim_to_non_expiring_account() { ensure_test_global_credentials(); let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let (cred, _) = iam_sys .new_service_account("svc-parent-user", None, NewServiceAccountOpts::default()) .await .expect("service account without explicit expiration should be created"); let updated_expiration = OffsetDateTime::now_utc() + time::Duration::hours(3); iam_sys .update_service_account( &cred.access_key, UpdateServiceAccountOpts { session_policy: None, secret_key: None, name: None, description: None, expiration: Some(updated_expiration), status: None, parent_user: None, allow_site_replicator_account: false, }, ) .await .expect("service account without expiration should accept a new expiration"); let updated_user = iam_sys .get_user(&cred.access_key) .await .expect("updated service account should exist"); assert_eq!(updated_user.credentials.expiration, Some(updated_expiration)); let claims = get_claims_from_token_with_secret(&updated_user.credentials.session_token, &updated_user.credentials.secret_key) .expect("updated service account JWT should decode after adding expiration"); assert_eq!(claims.get("exp").and_then(|v| v.as_i64()), Some(updated_expiration.unix_timestamp())); } #[tokio::test] async fn test_site_replicator_update_requires_explicit_internal_allowance() { ensure_test_global_credentials(); let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let (cred, _) = iam_sys .new_service_account( "svc-parent-user", None, NewServiceAccountOpts { access_key: SITE_REPLICATOR_SERVICE_ACCOUNT.to_string(), secret_key: "siteReplicatorSecretKeyForTest1234567890".to_string(), allow_site_replicator_account: true, ..Default::default() }, ) .await .expect("internal site replicator account should be created with explicit allowance"); assert!( iam_sys .update_service_account( &cred.access_key, UpdateServiceAccountOpts { session_policy: None, secret_key: None, name: None, description: None, expiration: None, status: Some(STATUS_ENABLED.to_string()), parent_user: None, allow_site_replicator_account: false, }, ) .await .is_err(), "ordinary update must not mutate the internal site replicator account" ); iam_sys .update_service_account( &cred.access_key, UpdateServiceAccountOpts { session_policy: None, secret_key: None, name: None, description: None, expiration: None, status: Some(STATUS_ENABLED.to_string()), parent_user: None, allow_site_replicator_account: true, }, ) .await .expect("internal update should be allowed explicitly"); assert_eq!( iam_sys .get_site_replicator_service_account_secret(&cred.access_key) .await .expect("internal secret should be readable for canonical account"), cred.secret_key ); } /// A root-credential change can leave `site-replicator-0` bound to a parent that no /// longer exists, and the repair must rebind in place: deleting first would leave the /// site with no replication account at all if the recreate failed. The parent also lives /// in the session token, so both copies have to move together or authorization denies /// the account. #[tokio::test] async fn test_site_replicator_parent_rebind_updates_credential_and_claim() { ensure_test_global_credentials(); let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let (cred, _) = iam_sys .new_service_account( "stale-parent-user", None, NewServiceAccountOpts { access_key: SITE_REPLICATOR_SERVICE_ACCOUNT.to_string(), secret_key: "siteReplicatorSecretKeyForTest1234567890".to_string(), allow_site_replicator_account: true, ..Default::default() }, ) .await .expect("site replicator account should be created"); iam_sys .update_service_account( &cred.access_key, UpdateServiceAccountOpts { session_policy: None, secret_key: None, name: None, description: None, expiration: None, status: None, parent_user: Some("current-parent-user".to_string()), allow_site_replicator_account: true, }, ) .await .expect("site replication repair may rebind the parent"); let (identity, claims) = iam_sys .get_account_with_claims_allow_missing_exp(&cred.access_key) .await .expect("rebound account should still be readable"); assert_eq!(identity.credentials.parent_user, "current-parent-user"); assert_eq!( claims.get("parent").and_then(Value::as_str), Some("current-parent-user"), "the token claim must follow the credential or authorization denies the account" ); assert_eq!( iam_sys .get_site_replicator_service_account_secret(&cred.access_key) .await .expect("secret survives a rebind"), cred.secret_key, "peers keep using the same secret, so a rebind must not rotate it" ); } /// The rebind is a site-replication repair primitive, not a general capability: letting /// any caller re-parent a service account would let it inherit another user's policies. #[tokio::test] async fn test_parent_rebind_is_rejected_for_ordinary_service_accounts() { ensure_test_global_credentials(); let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let (cred, _) = iam_sys .new_service_account( "ordinary-parent", None, NewServiceAccountOpts { access_key: "ordinary-service-account".to_string(), secret_key: "ordinaryServiceAccountSecret1234567890".to_string(), ..Default::default() }, ) .await .expect("ordinary service account should be created"); for allow in [false, true] { assert!( iam_sys .update_service_account( &cred.access_key, UpdateServiceAccountOpts { session_policy: None, secret_key: None, name: None, description: None, expiration: None, status: None, parent_user: Some("victim-user".to_string()), allow_site_replicator_account: allow, }, ) .await .is_err(), "re-parenting an ordinary service account must be rejected (allow={allow})" ); } } #[tokio::test] async fn test_created_access_token_authorizes_with_parent_policy() { ensure_test_global_credentials(); let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let parent_user = "sts-fallback-test-parent"; let groups = Some(vec!["testgroup".to_string()]); let (cred, _) = iam_sys .new_service_account( parent_user, groups.clone(), NewServiceAccountOpts { access_key: "ACCESSTOKENTESTUSER".to_string(), secret_key: "accessTokenTestSecret".to_string(), ..Default::default() }, ) .await .expect("access token should be created"); let stored = iam_sys .get_user(&cred.access_key) .await .expect("created access token should be cached"); assert!(stored.credentials.is_service_account()); assert_eq!(stored.credentials.parent_user, parent_user); let claims = stored .credentials .claims .as_ref() .expect("created access token should have decoded JWT claims"); assert_eq!(claims.get("parent").and_then(Value::as_str), Some(parent_user)); assert_eq!( claims.get(&iam_policy_claim_name_sa()).and_then(Value::as_str), Some(INHERITED_POLICY_TYPE) ); let (is_service_account, resolved_parent) = iam_sys .is_service_account(&cred.access_key) .await .expect("created access token should be recognized as a service account"); assert!(is_service_account); assert_eq!(resolved_parent, parent_user); let (redacted, policy) = iam_sys .get_service_account(&cred.access_key) .await .expect("created access token should be readable"); assert_eq!(redacted.access_key, cred.access_key); assert_eq!(redacted.parent_user, parent_user); assert!(redacted.secret_key.is_empty()); assert!(redacted.session_token.is_empty()); assert!(policy.is_none()); let args = Args { account: &cred.access_key, groups: &groups, action: Action::S3Action(S3Action::ListBucketAction), bucket: "mybucket", conditions: &HashMap::new(), is_owner: false, object: "", claims, deny_only: false, }; let prepared = iam_sys.prepare_auth(&args).await; assert!( matches!(prepared.mode, PreparedIamMode::ServiceAccount { .. }), "created access token must use service-account authorization path" ); assert!( iam_sys.eval_prepared(&prepared, &args).await, "created access token should be allowed through the parent's group policy" ); } /// GHSA-m77q-r63m-pj89 (STS JWTs signed with the shared root secret — /// intentionally UNFIXED). This test characterizes the STS credential /// lifecycle: a session token signed with the active token-signing key /// decodes and authorizes through the STS path. It pins that the SAME key /// both signs and verifies STS tokens. /// /// MUST UPDATE WHEN GHSA-m77q IS FIXED: the fix introduces a dedicated STS /// signing key distinct from the root secret. The narrower pin that captures /// the advisory's exact signature (signing key == root secret) is /// `test_ghsa_m77q_sts_session_token_signed_with_root_secret`; keep both in /// lockstep and flip them red -> green together. /// Advisory: #[tokio::test] async fn test_created_sts_credentials_authorize_with_session_token_claims() { ensure_test_global_credentials(); let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let parent_user = "sts-fallback-test-parent"; let token_secret = crate::root_credentials::token_signing_key() .unwrap_or_else(|| unreachable!("global action credentials should be initialized")); let mut claims = HashMap::new(); claims.insert("parent".to_string(), Value::String(parent_user.to_string())); claims.insert( "exp".to_string(), Value::Number(serde_json::Number::from( (OffsetDateTime::now_utc() + time::Duration::hours(1)).unix_timestamp(), )), ); let mut cred = get_new_credentials_with_metadata(&claims, &token_secret).expect("STS credentials should be created"); cred.parent_user = parent_user.to_string(); iam_sys .set_temp_user(&cred.access_key, &cred, None) .await .expect("STS credentials should be persisted in the temp-user cache"); let stored = iam_sys .get_user(&cred.access_key) .await .expect("created STS credentials should be cached"); assert!(stored.credentials.is_temp()); assert!(!stored.credentials.is_service_account()); assert_eq!(stored.credentials.parent_user, parent_user); let (is_temp, resolved_parent) = iam_sys .is_temp_user(&cred.access_key) .await .expect("created STS credentials should be recognized as temp"); assert!(is_temp); assert_eq!(resolved_parent, parent_user); let listed = iam_sys .list_sts_accounts(parent_user) .await .expect("created STS credentials should be listable by parent"); assert_eq!(listed.len(), 1); assert_eq!(listed[0].access_key, cred.access_key); assert_eq!(listed[0].parent_user, parent_user); assert!(listed[0].secret_key.is_empty()); assert!(listed[0].session_token.is_empty()); let temp_accounts = iam_sys .list_temp_accounts(parent_user) .await .expect("created STS credentials should be listable as temp accounts"); assert_eq!(temp_accounts.len(), 1); assert_eq!(temp_accounts[0].credentials.access_key, cred.access_key); assert_eq!(temp_accounts[0].credentials.parent_user, parent_user); assert!(temp_accounts[0].credentials.secret_key.is_empty()); assert!(temp_accounts[0].credentials.session_token.is_empty()); let (redacted, policy) = iam_sys .get_temporary_account(&cred.access_key) .await .expect("created STS credentials should be readable"); assert_eq!(redacted.access_key, cred.access_key); assert_eq!(redacted.parent_user, parent_user); assert!(redacted.secret_key.is_empty()); assert!(redacted.session_token.is_empty()); assert!(policy.is_none()); let decoded_claims = get_claims_from_token_with_secret(&cred.session_token, &token_secret) .expect("created STS session token should decode with the active signing key"); assert_eq!(decoded_claims.get("parent").and_then(Value::as_str), Some(parent_user)); let groups: Option> = None; let args = Args { account: &cred.access_key, groups: &groups, action: Action::S3Action(S3Action::ListBucketAction), bucket: "mybucket", conditions: &HashMap::new(), is_owner: false, object: "", claims: &decoded_claims, deny_only: false, }; let prepared = iam_sys.prepare_auth(&args).await; assert!( matches!(prepared.mode, PreparedIamMode::Sts { .. }), "created STS credentials must use STS authorization path" ); assert!( iam_sys.eval_prepared(&prepared, &args).await, "created STS credentials should inherit the parent user's group policy" ); } /// GHSA-m77q-r63m-pj89 flow-level pin (STS session tokens are signed with the /// shared root secret — intentionally UNFIXED). Anyone holding the root secret /// can forge STS session tokens because RustFS reuses the root secret as the /// STS token-signing key. This test PINS that vulnerable-by-design behavior end /// to end: (1) the token-signing key IS the root secret, (2) an AssumeRole-style /// session token issued with that key decodes with the ROOT secret and NOT with /// any other secret, and (3) it authorizes through the STS path. /// /// MUST UPDATE WHEN GHSA-m77q IS FIXED: the fix introduces a dedicated STS /// signing key distinct from the root secret. That makes assertion (1) and the /// "root secret decodes the token" assertion fail (red). Whoever fixes m77q must /// replace this pin with a regression asserting the NEW behavior — the root /// secret can no longer decode STS session tokens — i.e. red -> green. /// Advisory: #[tokio::test] async fn test_ghsa_m77q_sts_session_token_signed_with_root_secret() { ensure_test_global_credentials(); let root = crate::root_credentials::credentials().expect("global action (root) credentials should be initialized"); assert!(!root.secret_key.is_empty(), "root secret must be present for the STS signing-key pin"); // (1) m77q signature: the STS token-signing key IS the root secret. A fix // that introduces a dedicated STS key makes this fail -> update red->green. assert_eq!( crate::root_credentials::token_signing_key().as_deref(), Some(root.secret_key.as_str()), "GHSA-m77q pin: STS tokens are signed with the root secret. If this fails, \ m77q may have been fixed (dedicated STS key) — update this test red->green." ); // AssumeRole-style issuance: mint a session token with the active signing key. // Use the mock store's registered STS parent (group `testgroup` -> readwrite) // so authorization resolves through the STS path; the m77q pin is about the // signing key, not the parent identity. let parent_user = "sts-fallback-test-parent"; let signing_key = crate::root_credentials::token_signing_key().expect("STS token-signing key should be present"); let mut claims = HashMap::new(); claims.insert("parent".to_string(), Value::String(parent_user.to_string())); claims.insert( "exp".to_string(), Value::Number(serde_json::Number::from( (OffsetDateTime::now_utc() + time::Duration::hours(1)).unix_timestamp(), )), ); let mut cred = get_new_credentials_with_metadata(&claims, &signing_key).expect("STS credentials should be created"); cred.parent_user = parent_user.to_string(); // (2) The session token decodes with the ROOT secret specifically — the crux // of m77q — and rejects any non-root secret (HMAC verification fails). let decoded = get_claims_from_token_with_secret(&cred.session_token, &root.secret_key) .expect("GHSA-m77q pin: STS session token must decode with the ROOT secret"); assert_eq!(decoded.get("parent").and_then(Value::as_str), Some(parent_user)); assert!( get_claims_from_token_with_secret(&cred.session_token, "not-the-root-secret").is_err(), "GHSA-m77q pin: STS session token must NOT decode with a non-root secret" ); // (3) The root-secret-signed STS credentials authorize end to end. let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); iam_sys .set_temp_user(&cred.access_key, &cred, None) .await .expect("STS credentials should be persisted in the temp-user cache"); let groups: Option> = None; let args = Args { account: &cred.access_key, groups: &groups, action: Action::S3Action(S3Action::ListBucketAction), bucket: "mybucket", conditions: &HashMap::new(), is_owner: false, object: "", claims: &decoded, deny_only: false, }; let prepared = iam_sys.prepare_auth(&args).await; assert!( matches!(prepared.mode, PreparedIamMode::Sts { .. }), "root-secret-signed STS credentials must use the STS authorization path" ); assert!( iam_sys.eval_prepared(&prepared, &args).await, "root-secret-signed STS credentials should authorize through the parent group policy (m77q)" ); } /// Regression test: temp credentials without groups in args still receive group-attached /// policies via the parent user (groups fallback). Without the fallback, policy_db_get /// would get None for groups and the user would have no group policies, so the action /// would be denied. #[tokio::test] async fn test_sts_groups_fallback_temp_creds_receive_parent_group_policies() { let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let parent_user = "sts-fallback-test-parent"; let claims = HashMap::new(); let groups: Option> = None; let args = Args { account: parent_user, groups: &groups, action: Action::S3Action(S3Action::ListBucketAction), bucket: "mybucket", conditions: &HashMap::new(), is_owner: false, object: "", claims: &claims, deny_only: false, }; let prepared = iam_sys.prepare_sts_auth(&args, parent_user).await; let allowed = iam_sys.eval_prepared(&prepared, &args).await; assert!( allowed, "STS temp credentials with no groups in args should still be allowed via parent user's group policy (readwrite)" ); } #[tokio::test] async fn test_sts_claim_policy_resolves_custom_canned_policy() { let store = StsTestMockStore::new(true); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let parent_user = "sts-empty-parent-policy-test"; let sts_access_key = "sts-custom-claim-policy-test-user"; let sts_user = UserIdentity::from(Credentials { access_key: sts_access_key.to_string(), secret_key: "longenoughsecret".to_string(), session_token: "sts-token".to_string(), status: ACCOUNT_ON.to_string(), parent_user: parent_user.to_string(), ..Default::default() }); iam_sys .store .cache .add_or_update_sts_account(sts_access_key, &sts_user, OffsetDateTime::now_utc()); let mut claims = HashMap::new(); claims.insert(POLICYNAME.to_string(), Value::String(CUSTOM_STS_CLAIM_POLICY.to_string())); let groups: Option> = None; let args = Args { account: sts_access_key, groups: &groups, action: Action::S3Action(S3Action::GetObjectAction), bucket: CUSTOM_STS_CLAIM_BUCKET, conditions: &HashMap::new(), is_owner: false, object: "allowed/object.txt", claims: &claims, deny_only: false, }; let prepared = iam_sys.prepare_sts_auth(&args, parent_user).await; assert!(matches!(prepared.mode, PreparedIamMode::Sts { .. })); assert!( iam_sys.eval_prepared(&prepared, &args).await, "STS temp credentials should resolve custom canned policy names carried in JWT policy claims" ); } #[tokio::test] async fn oidc_service_account_uses_verified_policy_and_persisted_boundary() { ensure_test_global_credentials(); let iam_sys = IamSys::new(IamCache::new(StsTestMockStore::new(true)).await.unwrap()); let parent_user = "TwyekekG2eMes0qk9Tgh7KXEitwGi1z2W1f2KccrXGA"; let mut oidc_claims = HashMap::from([ ("iss".to_string(), Value::String("rustfs-oidc".to_string())), ("oidc_provider".to_string(), Value::String("default".to_string())), ("sub".to_string(), Value::String("subject-123".to_string())), ("parent".to_string(), Value::String(parent_user.to_string())), (OIDC_VIRTUAL_PARENT_CLAIM.to_string(), Value::String(parent_user.to_string())), (POLICYNAME.to_string(), Value::String("readwrite".to_string())), ]); oidc_claims.insert( SESSION_POLICY_NAME.to_string(), Value::String(base64_simd::URL_SAFE_NO_PAD.encode_to_string(CUSTOM_STS_CLAIM_POLICY_JSON.as_bytes())), ); oidc_claims.insert("exp".to_string(), Value::Number(123456.into())); let caller = Credentials { access_key: "OIDCTEMPORARYCALLER1".to_string(), secret_key: "oidcTemporaryCallerSecret123".to_string(), session_token: "signed-session-token".to_string(), parent_user: parent_user.to_string(), claims: Some(oidc_claims), ..Default::default() }; let (credential, _, replication_claims) = iam_sys .new_service_account_from_caller( parent_user, None, NewServiceAccountOpts { access_key: "OIDCSERVICEACCOUNT01".to_string(), secret_key: "oidcServiceAccountSecret123".to_string(), ..Default::default() }, &caller, ) .await .expect("OIDC service account should be created"); assert_eq!(replication_claims.get(VERIFIED_FEDERATED_POLICY_CLAIM), Some(&Value::Bool(true))); assert_eq!( replication_claims.get(OIDC_VIRTUAL_PARENT_CLAIM), Some(&Value::String(parent_user.to_string())) ); assert!(replication_claims.contains_key(FEDERATED_POLICY_BOUNDARY_CLAIM)); assert!(!replication_claims.contains_key("exp")); let claims = iam_sys .get_claims_for_svc_acc(&credential.access_key) .await .expect("stored service-account claims should decode"); let groups = None; let conditions = HashMap::new(); for (action, allowed) in [(S3Action::GetObjectAction, true), (S3Action::PutObjectAction, false)] { let args = Args { account: &credential.access_key, groups: &groups, action: Action::S3Action(action), bucket: CUSTOM_STS_CLAIM_BUCKET, conditions: &conditions, is_owner: false, object: "allowed/object.txt", claims: &claims, deny_only: false, }; assert_eq!( iam_sys.is_allowed(&args).await, allowed, "service-account permissions must be the OIDC policy intersected with its session boundary" ); } iam_sys.store.cache.delete_policy_doc("readwrite", OffsetDateTime::now_utc()); let args = Args { account: &credential.access_key, groups: &groups, action: Action::S3Action(S3Action::GetObjectAction), bucket: CUSTOM_STS_CLAIM_BUCKET, conditions: &conditions, is_owner: false, object: "allowed/object.txt", claims: &claims, deny_only: false, }; assert!( !iam_sys.is_allowed(&args).await, "OIDC service accounts must resolve current policy documents instead of retaining a policy snapshot" ); } #[tokio::test] async fn oidc_group_policy_applies_to_sts_and_service_account() { ensure_test_global_credentials(); let iam_sys = IamSys::new(IamCache::new(StsTestMockStore::new(false)).await.unwrap()); let parent_user = "openid=group-policy-parent"; let groups = Some(vec!["testgroup".to_string()]); let claims = HashMap::from([ ("iss".to_string(), Value::String("rustfs-oidc".to_string())), ("oidc_provider".to_string(), Value::String("default".to_string())), ("sub".to_string(), Value::String("subject-123".to_string())), ("parent".to_string(), Value::String(parent_user.to_string())), (OIDC_VIRTUAL_PARENT_CLAIM.to_string(), Value::String(parent_user.to_string())), (POLICYNAME.to_string(), Value::String("readwrite".to_string())), ]); let conditions = HashMap::new(); let sts_args = Args { account: "OIDCGROUPCALLER01", groups: &groups, action: Action::S3Action(S3Action::GetObjectAction), bucket: "bucket", conditions: &conditions, is_owner: false, object: "object.txt", claims: &claims, deny_only: false, }; let prepared = iam_sys.prepare_sts_auth(&sts_args, parent_user).await; assert!( iam_sys.eval_prepared(&prepared, &sts_args).await, "OIDC STS must include policies mapped from verified groups" ); let caller = Credentials { access_key: "OIDCGROUPCALLER01".to_string(), parent_user: parent_user.to_string(), groups: groups.clone(), claims: Some(claims.clone()), ..Default::default() }; let (service_account, _, replication_claims) = iam_sys .new_service_account_from_caller( parent_user, groups.clone(), NewServiceAccountOpts { access_key: "OIDCGROUPSERVICE01".to_string(), secret_key: "oidcGroupServiceSecret123".to_string(), ..Default::default() }, &caller, ) .await .expect("group-authorized OIDC caller should create a service account"); assert_eq!(replication_claims.get(VERIFIED_FEDERATED_POLICY_CLAIM), Some(&Value::Bool(true))); assert_eq!(replication_claims.get(POLICYNAME), Some(&Value::String("readwrite".to_string()))); let service_claims = iam_sys .get_claims_for_svc_acc(&service_account.access_key) .await .expect("service-account claims"); let service_args = Args { account: &service_account.access_key, groups: &groups, claims: &service_claims, ..sts_args }; assert!( iam_sys.is_allowed(&service_args).await, "OIDC service account must resolve the policy names selected at STS issuance" ); } #[tokio::test] async fn verified_oidc_policy_names_do_not_drift_after_group_mapping_changes() { let iam_sys = IamSys::new(IamCache::new(StsTestMockStore::new(false)).await.unwrap()); let deny_delete = Policy::parse_config( br#"{"Version":"2012-10-17","Statement":[{"Effect":"Deny","Action":["s3:DeleteObject"],"Resource":["arn:aws:s3:::bucket/*"]}]}"#, ) .expect("parse deny policy"); let now = OffsetDateTime::now_utc(); iam_sys .store .cache .add_or_update_policy_doc("deny-delete", &PolicyDoc::new(deny_delete), now); iam_sys .store .cache .add_or_update_group_policy("testgroup", &MappedPolicy::new("writeonly"), now); let parent_user = "openid=group-deny-parent"; let claims = HashMap::from([ ("iss".to_string(), Value::String("rustfs-oidc".to_string())), ("oidc_provider".to_string(), Value::String("default".to_string())), ("sub".to_string(), Value::String("subject-123".to_string())), ("parent".to_string(), Value::String(parent_user.to_string())), (OIDC_VIRTUAL_PARENT_CLAIM.to_string(), Value::String(parent_user.to_string())), (POLICYNAME.to_string(), Value::String("readonly,deny-delete".to_string())), ]); let groups = Some(vec!["testgroup".to_string()]); let conditions = HashMap::new(); for (action, allowed) in [ (S3Action::GetObjectAction, true), (S3Action::PutObjectAction, false), (S3Action::DeleteObjectAction, false), ] { let args = Args { account: "OIDCGROUPCALLER02", groups: &groups, action: Action::S3Action(action), bucket: "bucket", conditions: &conditions, is_owner: false, object: "object.txt", claims: &claims, deny_only: false, }; let prepared = iam_sys.prepare_sts_auth(&args, parent_user).await; assert_eq!( iam_sys.eval_prepared(&prepared, &args).await, allowed, "verified OIDC sessions must retain the policy names selected at issuance" ); } } #[tokio::test] async fn active_legacy_oidc_session_can_create_service_account_after_upgrade() { ensure_test_global_credentials(); let iam_sys = IamSys::new(IamCache::new(StsTestMockStore::new(true)).await.unwrap()); let parent_user = "legacy-oidc-user"; let legacy_claims = HashMap::from([ ("iss".to_string(), Value::String("rustfs-oidc".to_string())), ("oidc_provider".to_string(), Value::String("default".to_string())), ("sub".to_string(), Value::String("subject-123".to_string())), ("parent".to_string(), Value::String(parent_user.to_string())), (POLICYNAME.to_string(), Value::String("readwrite".to_string())), ]); let caller = Credentials { access_key: "LEGACYOIDCCALLER01".to_string(), parent_user: parent_user.to_string(), claims: Some(legacy_claims), ..Default::default() }; let (service_account, _, replication_claims) = iam_sys .new_service_account_from_caller( parent_user, None, NewServiceAccountOpts { access_key: "LEGACYOIDCSERVICE01".to_string(), secret_key: "legacyOidcServiceSecret123".to_string(), ..Default::default() }, &caller, ) .await .expect("active pre-upgrade OIDC session should remain usable"); assert!(!replication_claims.contains_key(OIDC_VIRTUAL_PARENT_CLAIM)); assert_eq!(replication_claims.get(POLICYNAME), Some(&Value::String("readwrite".to_string()))); let claims = iam_sys .get_claims_for_svc_acc(&service_account.access_key) .await .expect("legacy-derived service-account claims"); let groups = None; let args = Args { account: &service_account.access_key, groups: &groups, action: Action::S3Action(S3Action::GetObjectAction), bucket: "bucket", conditions: &HashMap::new(), is_owner: false, object: "object.txt", claims: &claims, deny_only: false, }; assert!(iam_sys.is_allowed(&args).await); } #[tokio::test] async fn non_owner_import_cannot_preserve_verified_federated_policy() { ensure_test_global_credentials(); let iam_sys = test_iam_sys().await; let parent_user = "sts-fallback-test-parent"; iam_sys .store .cache .add_or_update_user_policy(parent_user, &MappedPolicy::new("readwrite"), OffsetDateTime::now_utc()); let mut claims = HashMap::from([ ("iss".to_string(), Value::String("rustfs-oidc".to_string())), ("oidc_provider".to_string(), Value::String("default".to_string())), ("sub".to_string(), Value::String("subject-123".to_string())), ("parent".to_string(), Value::String(parent_user.to_string())), (OIDC_VIRTUAL_PARENT_CLAIM.to_string(), Value::String(parent_user.to_string())), (POLICYNAME.to_string(), Value::String("readwrite".to_string())), (VERIFIED_FEDERATED_POLICY_CLAIM.to_string(), Value::Bool(true)), ( FEDERATED_POLICY_BOUNDARY_CLAIM.to_string(), Value::String(base64_simd::URL_SAFE_NO_PAD.encode_to_string(CUSTOM_STS_CLAIM_POLICY_JSON.as_bytes())), ), ]); remove_verified_federated_policy(&mut claims); let (credential, _) = iam_sys .new_service_account( parent_user, None, NewServiceAccountOpts { access_key: "IMPORTEDOIDCSERVICE1".to_string(), secret_key: "importedOidcServiceSecret123".to_string(), claims: Some(claims), ..Default::default() }, ) .await .expect("sanitized service account should be persisted"); let stored_claims = iam_sys .get_claims_for_svc_acc(&credential.access_key) .await .expect("stored service-account claims should decode"); assert!(!stored_claims.contains_key(VERIFIED_FEDERATED_POLICY_CLAIM)); assert!(!stored_claims.contains_key(FEDERATED_POLICY_BOUNDARY_CLAIM)); assert!(!stored_claims.contains_key(OIDC_VIRTUAL_PARENT_CLAIM)); let args = Args { account: &credential.access_key, groups: &None, action: Action::S3Action(S3Action::GetObjectAction), bucket: CUSTOM_STS_CLAIM_BUCKET, conditions: &HashMap::new(), is_owner: false, object: "allowed/object.txt", claims: &stored_claims, deny_only: false, }; assert!(!iam_sys.is_allowed(&args).await); } #[tokio::test] async fn oidc_service_account_rejects_mismatched_virtual_parent() { ensure_test_global_credentials(); let iam_sys = IamSys::new(IamCache::new(StsTestMockStore::new(true)).await.unwrap()); let caller = Credentials { access_key: "OIDCTEMPORARYCALLER2".to_string(), parent_user: "openid=parent-b".to_string(), claims: Some(HashMap::from([ ("iss".to_string(), Value::String("rustfs-oidc".to_string())), ("oidc_provider".to_string(), Value::String("default".to_string())), ("sub".to_string(), Value::String("subject-123".to_string())), ("parent".to_string(), Value::String("openid=parent-b".to_string())), (OIDC_VIRTUAL_PARENT_CLAIM.to_string(), Value::String("openid=parent-a".to_string())), (POLICYNAME.to_string(), Value::String("readwrite".to_string())), ])), ..Default::default() }; let result = iam_sys .new_service_account_from_caller( "openid=parent-b", None, NewServiceAccountOpts { access_key: "OIDCSERVICEACCOUNT02".to_string(), secret_key: "oidcServiceAccountSecret234".to_string(), ..Default::default() }, &caller, ) .await; assert!(matches!(result, Err(IamError::IAMActionNotAllowed))); } #[tokio::test] async fn verified_oidc_sts_missing_policy_fails_closed_for_deny_only_checks() { let iam_sys = IamSys::new(IamCache::new(StsTestMockStore::new(true)).await.unwrap()); let parent_user = "openid=verified-parent"; let claims = HashMap::from([ ("iss".to_string(), Value::String("rustfs-oidc".to_string())), ("oidc_provider".to_string(), Value::String("default".to_string())), ("sub".to_string(), Value::String("subject-123".to_string())), ("parent".to_string(), Value::String(parent_user.to_string())), (OIDC_VIRTUAL_PARENT_CLAIM.to_string(), Value::String(parent_user.to_string())), (POLICYNAME.to_string(), Value::String("readwrite,missing-policy".to_string())), ]); let groups = None; let args = Args { account: "OIDCTEMPORARYCALLER3", groups: &groups, action: Action::AdminAction(AdminAction::ListServiceAccountsAdminAction), bucket: "", conditions: &HashMap::new(), is_owner: false, object: "", claims: &claims, deny_only: true, }; let prepared = iam_sys.prepare_sts_auth(&args, parent_user).await; assert!(matches!(prepared.mode, PreparedIamMode::Deny)); assert!(!iam_sys.eval_prepared(&prepared, &args).await); } #[tokio::test] async fn regular_deny_only_rejects_unresolved_mapped_policy() { let iam_sys = test_iam_sys().await; let user = "regular-unresolved-policy"; let identity = UserIdentity::from(Credentials { access_key: user.to_string(), secret_key: "longenoughsecret".to_string(), status: ACCOUNT_ON.to_string(), ..Default::default() }); iam_sys.store.cache.with_write_lock(|cache| { let now = OffsetDateTime::now_utc(); cache.add_or_update_user(user, &identity, now); cache.add_or_update_user_policy(user, &MappedPolicy::new("readwrite,missing-policy"), now); }); let groups = None; let claims = HashMap::new(); let conditions = HashMap::new(); let args = Args { account: user, groups: &groups, action: Action::StsAction(StsAction::AssumeRoleAction), bucket: "", conditions: &conditions, is_owner: false, object: "", claims: &claims, deny_only: true, }; let prepared = iam_sys.prepare_regular_auth(&args).await; assert!(matches!(prepared.mode, PreparedIamMode::Deny)); assert!( !iam_sys.eval_prepared(&prepared, &args).await, "deny-only evaluation must fail closed when a mapped policy document is unresolved" ); } #[tokio::test] async fn regular_full_evaluation_keeps_resolved_part_of_partial_mapping() { let iam_sys = test_iam_sys().await; let user = "regular-partial-policy"; let identity = UserIdentity::from(Credentials { access_key: user.to_string(), secret_key: "longenoughsecret".to_string(), status: ACCOUNT_ON.to_string(), ..Default::default() }); iam_sys.store.cache.with_write_lock(|cache| { let now = OffsetDateTime::now_utc(); cache.add_or_update_user(user, &identity, now); cache.add_or_update_user_policy(user, &MappedPolicy::new("readwrite,missing-policy"), now); }); let groups = None; let claims = HashMap::new(); let conditions = HashMap::new(); let args = Args { account: user, groups: &groups, action: Action::S3Action(S3Action::GetObjectAction), bucket: "bucket", conditions: &conditions, is_owner: false, object: "object", claims: &claims, deny_only: false, }; let prepared = iam_sys.prepare_regular_auth(&args).await; assert!(matches!(prepared.mode, PreparedIamMode::Regular { .. })); assert!( iam_sys.eval_prepared(&prepared, &args).await, "full evaluation must preserve the historical partial-resolution behavior" ); } #[tokio::test] async fn regular_deny_only_honors_group_derived_explicit_deny() { let iam_sys = test_iam_sys().await; let deny_policy = Policy::parse_config(br#"{"Version":"2012-10-17","Statement":[{"Effect":"Deny","Action":["sts:AssumeRole"]}]}"#) .expect("group deny policy should parse"); let now = OffsetDateTime::now_utc(); iam_sys .store .cache .add_or_update_policy_doc("deny-assume-role", &PolicyDoc::new(deny_policy), now); iam_sys .store .cache .add_or_update_group_policy("testgroup", &MappedPolicy::new("deny-assume-role"), now); let groups = Some(vec!["testgroup".to_string()]); let claims = HashMap::new(); let conditions = HashMap::new(); let args = Args { account: "sts-fallback-test-parent", groups: &groups, action: Action::StsAction(StsAction::AssumeRoleAction), bucket: "", conditions: &conditions, is_owner: false, object: "", claims: &claims, deny_only: true, }; let prepared = iam_sys.prepare_regular_auth(&args).await; assert!(matches!(prepared.mode, PreparedIamMode::Regular { .. })); assert!( !iam_sys.eval_prepared(&prepared, &args).await, "group-derived explicit Deny must remain authoritative" ); } #[tokio::test] async fn regular_deny_only_rejects_unresolved_group_mapping() { let iam_sys = test_iam_sys().await; iam_sys.store.cache.add_or_update_group_policy( "testgroup", &MappedPolicy::new("readwrite,missing-policy"), OffsetDateTime::now_utc(), ); let groups = Some(vec!["testgroup".to_string()]); let claims = HashMap::new(); let conditions = HashMap::new(); let args = Args { account: "sts-fallback-test-parent", groups: &groups, action: Action::StsAction(StsAction::AssumeRoleAction), bucket: "", conditions: &conditions, is_owner: false, object: "", claims: &claims, deny_only: true, }; let prepared = iam_sys.prepare_regular_auth(&args).await; assert!(matches!(prepared.mode, PreparedIamMode::Deny)); assert!( !iam_sys.eval_prepared(&prepared, &args).await, "deny-only evaluation must fail closed for an unresolved group-derived mapping" ); } #[tokio::test] async fn test_sts_claim_policy_ignores_unsafe_and_missing_policy_names() { let store = StsTestMockStore::new(true); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let parent_user = "sts-empty-parent-policy-test"; let sts_access_key = "sts-mixed-claim-policy-test-user"; let sts_user = UserIdentity::from(Credentials { access_key: sts_access_key.to_string(), secret_key: "longenoughsecret".to_string(), session_token: "sts-token".to_string(), status: ACCOUNT_ON.to_string(), parent_user: parent_user.to_string(), ..Default::default() }); iam_sys .store .cache .add_or_update_sts_account(sts_access_key, &sts_user, OffsetDateTime::now_utc()); let mut claims = HashMap::new(); claims.insert( POLICYNAME.to_string(), Value::String(format!("unsafe/policy, missing-sts-claim-policy, {CUSTOM_STS_CLAIM_POLICY}")), ); let groups: Option> = None; let args = Args { account: sts_access_key, groups: &groups, action: Action::S3Action(S3Action::GetObjectAction), bucket: CUSTOM_STS_CLAIM_BUCKET, conditions: &HashMap::new(), is_owner: false, object: "allowed/object.txt", claims: &claims, deny_only: false, }; let prepared = iam_sys.prepare_sts_auth(&args, parent_user).await; assert!(matches!(prepared.mode, PreparedIamMode::Sts { .. })); assert!( iam_sys.eval_prepared(&prepared, &args).await, "STS policy claims should ignore unsafe or unresolved names without dropping a resolvable canned policy" ); } #[tokio::test] async fn test_sts_claim_policy_custom_canned_policy_does_not_grant_other_actions() { let store = StsTestMockStore::new(true); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let parent_user = "sts-empty-parent-policy-test"; let sts_access_key = "sts-custom-claim-policy-deny-test-user"; let sts_user = UserIdentity::from(Credentials { access_key: sts_access_key.to_string(), secret_key: "longenoughsecret".to_string(), session_token: "sts-token".to_string(), status: ACCOUNT_ON.to_string(), parent_user: parent_user.to_string(), ..Default::default() }); iam_sys .store .cache .add_or_update_sts_account(sts_access_key, &sts_user, OffsetDateTime::now_utc()); let mut claims = HashMap::new(); claims.insert(POLICYNAME.to_string(), Value::String(CUSTOM_STS_CLAIM_POLICY.to_string())); let groups: Option> = None; let args = Args { account: sts_access_key, groups: &groups, action: Action::S3Action(S3Action::PutObjectAction), bucket: CUSTOM_STS_CLAIM_BUCKET, conditions: &HashMap::new(), is_owner: false, object: "allowed/object.txt", claims: &claims, deny_only: false, }; let prepared = iam_sys.prepare_sts_auth(&args, parent_user).await; assert!(matches!(prepared.mode, PreparedIamMode::Sts { .. })); assert!( !iam_sys.eval_prepared(&prepared, &args).await, "custom claim policy must not grant S3 actions outside the resolved canned policy" ); } #[tokio::test] async fn test_sts_claim_policy_builtin_policy_remains_compatible() { let store = StsTestMockStore::new(true); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let parent_user = "sts-empty-parent-policy-test"; let mut claims = HashMap::new(); claims.insert(POLICYNAME.to_string(), Value::String("readwrite".to_string())); let groups: Option> = None; let args = Args { account: "sts-builtin-claim-policy-test-user", groups: &groups, action: Action::S3Action(S3Action::ListBucketAction), bucket: "mybucket", conditions: &HashMap::new(), is_owner: false, object: "", claims: &claims, deny_only: false, }; let prepared = iam_sys.prepare_sts_auth(&args, parent_user).await; assert!(matches!(prepared.mode, PreparedIamMode::Sts { .. })); assert!( iam_sys.eval_prepared(&prepared, &args).await, "built-in policy names in STS JWT claims must keep working through the unified policy store path" ); } #[tokio::test] async fn test_sts_claim_policy_missing_policy_denies() { let store = StsTestMockStore::new(true); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let parent_user = "sts-empty-parent-policy-test"; let mut claims = HashMap::new(); claims.insert(POLICYNAME.to_string(), Value::String("missing-sts-claim-policy".to_string())); let groups: Option> = None; let args = Args { account: "sts-missing-claim-policy-test-user", groups: &groups, action: Action::S3Action(S3Action::GetObjectAction), bucket: CUSTOM_STS_CLAIM_BUCKET, conditions: &HashMap::new(), is_owner: false, object: "allowed/object.txt", claims: &claims, deny_only: false, }; let prepared = iam_sys.prepare_sts_auth(&args, parent_user).await; assert!(matches!(prepared.mode, PreparedIamMode::Deny)); assert!( !iam_sys.eval_prepared(&prepared, &args).await, "missing STS claim policy names must deny instead of silently allowing" ); } /// Regression: `deny_only` with empty IAM policies must still evaluate `sessionPolicy-extracted` /// so session policy Deny cannot be bypassed (see PR #2250 review). #[tokio::test] async fn test_sts_deny_only_session_policy_deny_blocks_when_iam_policies_empty() { let store = StsTestMockStore::new(true); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let parent_user = "sts-empty-parent-policy-test"; let session_policy_json = r#"{ "Version": "2012-10-17", "Statement": [ { "Effect": "Deny", "Action": ["admin:CreateUser"], "Resource": ["arn:aws:s3:::*"] } ] }"#; let mut claims = HashMap::new(); claims.insert(SESSION_POLICY_NAME_EXTRACTED.to_string(), Value::String(session_policy_json.to_string())); let groups: Option> = None; let args = Args { account: parent_user, groups: &groups, action: Action::AdminAction(AdminAction::CreateUserAdminAction), bucket: "", conditions: &HashMap::new(), is_owner: false, object: "", claims: &claims, deny_only: true, }; let prepared = iam_sys.prepare_sts_auth(&args, parent_user).await; let allowed = iam_sys.eval_prepared(&prepared, &args).await; assert!( !allowed, "session policy Deny must be evaluated even when IAM policies are empty and deny_only is set" ); } #[tokio::test] async fn test_sts_deny_only_session_policy_allow_when_no_deny_on_action() { let store = StsTestMockStore::new(true); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let parent_user = "sts-empty-parent-policy-test"; let session_policy_json = r#"{ "Version": "2012-10-17", "Statement": [ { "Effect": "Allow", "Action": ["s3:GetObject"], "Resource": ["arn:aws:s3:::bucket/*"] } ] }"#; let mut claims = HashMap::new(); claims.insert(SESSION_POLICY_NAME_EXTRACTED.to_string(), Value::String(session_policy_json.to_string())); let groups: Option> = None; let args = Args { account: parent_user, groups: &groups, action: Action::AdminAction(AdminAction::CreateUserAdminAction), bucket: "", conditions: &HashMap::new(), is_owner: false, object: "", claims: &claims, deny_only: true, }; let prepared = iam_sys.prepare_sts_auth(&args, parent_user).await; let allowed = iam_sys.eval_prepared(&prepared, &args).await; assert!( allowed, "deny_only with no matching Deny in session policy should still allow self-service-style checks" ); } /// Regression test for cross-node IAM notifications: /// `load_user` must populate user cache, and regular-user mapped policy must be written to /// `user_policies` (not `sts_policies`), otherwise list-users and bucket-scoped user listing /// may miss users on follower nodes. #[tokio::test] async fn test_load_user_notification_populates_user_and_policy_caches() { let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); iam_sys.load_user("notify-user", UserType::Reg).await.unwrap(); let users = iam_sys.list_users().await.unwrap(); assert!( users.contains_key("notify-user"), "regular user loaded via notification must appear in list_users cache view" ); let bucket_users = iam_sys.list_bucket_users("notification-regression-bucket").await.unwrap(); assert!( bucket_users.contains_key("notify-user"), "regular user mapped policy must be written to user_policies for bucket user listing" ); } #[tokio::test] async fn test_group_notification_populates_new_membership_entry() { let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); iam_sys.store.cache.with_write_lock(|cache| { cache.replace_user_group_memberships(CacheEntity::default()); }); iam_sys.load_group("notify-group").await.unwrap(); let cache = iam_sys.store.cache.snapshot(); let memberships = &cache.user_group_memberships; assert!( memberships .get("notify-user") .is_some_and(|groups| groups.contains("notify-group")), "group notification must create a reverse membership entry for first-time members" ); } #[tokio::test] async fn test_sts_policy_mapping_notification_updates_sts_policy_cache() { let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); iam_sys .load_policy_mapping("notify-sts-parent", UserType::Sts, false) .await .unwrap(); let cache = iam_sys.store.cache.snapshot(); let sts_policies = &cache.sts_policies; assert!( sts_policies.contains_key("notify-sts-parent"), "STS policy mapping notifications must update sts_policies instead of deleting them" ); } #[tokio::test] async fn test_sts_policy_lookup_loads_missing_mapping_without_regular_user() { let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let policies = iam_sys.sts_policy_db_get("notify-sts-parent", &None).await.unwrap(); assert_eq!(policies, vec!["readwrite"]); assert!(iam_sys.store.cache.snapshot().sts_policies.contains_key("notify-sts-parent")); } #[tokio::test] async fn test_missing_user_notification_cleans_related_cache_state() { let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); const USER: &str = "deleted-notify-user"; const GROUP: &str = "deleted-notify-group"; const SVC_CHILD: &str = "deleted-notify-service-child"; const STS_CHILD: &str = "deleted-notify-sts-child"; const OTHER_USER: &str = "deleted-notify-other-user"; let user = UserIdentity::from(Credentials { access_key: USER.to_string(), secret_key: "longenoughsecret".to_string(), status: ACCOUNT_ON.to_string(), ..Default::default() }); let mut service_claims = HashMap::new(); service_claims.insert(iam_policy_claim_name_sa(), Value::String(INHERITED_POLICY_TYPE.to_string())); let service_child = UserIdentity::from(Credentials { access_key: SVC_CHILD.to_string(), secret_key: "longenoughsecret".to_string(), status: ACCOUNT_ON.to_string(), parent_user: USER.to_string(), claims: Some(service_claims), ..Default::default() }); let sts_child = UserIdentity::from(Credentials { access_key: STS_CHILD.to_string(), secret_key: "longenoughsecret".to_string(), session_token: "session-token".to_string(), status: ACCOUNT_ON.to_string(), parent_user: USER.to_string(), ..Default::default() }); let membership = HashSet::from([GROUP.to_string()]); let group = GroupInfo::new(vec![USER.to_string(), OTHER_USER.to_string()]); let mapped_policy = MappedPolicy::new("readwrite"); iam_sys.store.cache.with_write_lock(|cache| { let now = OffsetDateTime::now_utc(); cache.add_or_update_user(USER, &user, now); cache.add_or_update_user_policy(USER, &mapped_policy, now); cache.add_or_update_group(GROUP, &group, now); cache.add_or_update_user_group_membership(USER, &membership, now); cache.add_or_update_user(SVC_CHILD, &service_child, now); cache.add_or_update_sts_account(STS_CHILD, &sts_child, now); }); iam_sys.load_user(USER, UserType::Reg).await.unwrap(); let cache = iam_sys.store.cache.snapshot(); assert!(!cache.users.contains_key(USER)); assert!(!cache.user_policies.contains_key(USER)); assert!(!cache.user_group_memberships.contains_key(USER)); assert!(!cache.users.contains_key(SVC_CHILD)); assert!(!cache.sts_accounts.contains_key(STS_CHILD)); let group = cache.groups.get(GROUP).expect("group should remain after member removal"); assert!(!group.members.contains(&USER.to_string())); assert!(group.members.contains(&OTHER_USER.to_string())); } #[tokio::test] async fn test_check_key_propagates_cache_miss_load_failure() { let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let result = iam_sys.check_key("load-failure-user").await; assert!(matches!(result, Err(Error::Io(_)))); } #[tokio::test] async fn test_prepare_auth_eval_matches_prepare_sts_auth_for_parent_policy_fallback() { let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let parent_user = "sts-fallback-test-parent"; let claims = HashMap::new(); let groups: Option> = None; let args = Args { account: parent_user, groups: &groups, action: Action::S3Action(S3Action::ListBucketAction), bucket: "mybucket", conditions: &HashMap::new(), is_owner: false, object: "", claims: &claims, deny_only: false, }; let sts_prepared = iam_sys.prepare_sts_auth(&args, parent_user).await; let sts_eval = iam_sys.eval_prepared(&sts_prepared, &args).await; let prepared = iam_sys.prepare_auth(&args).await; let eval = iam_sys.eval_prepared(&prepared, &args).await; assert_eq!(sts_eval, eval, "prepare_auth must match explicit STS preparation for this identity"); } #[tokio::test] async fn test_prepare_auth_detects_existing_object_tag_in_session_policy() { let store = StsTestMockStore::new(true); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let sts_access_key = "sts-session-tag-test-user"; let sts_user = UserIdentity::from(Credentials { access_key: sts_access_key.to_string(), secret_key: "longenoughsecret".to_string(), session_token: "sts-token".to_string(), status: ACCOUNT_ON.to_string(), parent_user: "sts-empty-parent-policy-test".to_string(), ..Default::default() }); iam_sys .store .cache .add_or_update_sts_account(sts_access_key, &sts_user, OffsetDateTime::now_utc()); let mut claims = HashMap::new(); claims.insert( SESSION_POLICY_NAME_EXTRACTED.to_string(), Value::String( r#"{ "Version":"2012-10-17", "Statement":[{"Effect":"Allow","Action":["s3:GetObject"],"Resource":["arn:aws:s3:::bucket/*"],"Condition":{"StringEquals":{"s3:ExistingObjectTag/security":"public"}}}] }"# .to_string(), ), ); let groups: Option> = None; let args = Args { account: sts_access_key, groups: &groups, action: Action::S3Action(S3Action::GetObjectAction), bucket: "bucket", conditions: &HashMap::new(), is_owner: false, object: "obj", claims: &claims, deny_only: true, }; let prepared = iam_sys.prepare_auth(&args).await; assert!( prepared.needs_existing_object_tag, "session policy with ExistingObjectTag must request object tag loading" ); } #[test] fn test_policy_uses_existing_object_tag_matches_condition_keys_only() { let with_value_only = Policy::parse_config( br#"{ "Version":"2012-10-17", "Statement":[{ "Effect":"Allow", "Action":["s3:GetObject"], "Resource":["arn:aws:s3:::bucket/*"], "Condition":{"StringEquals":{"s3:prefix":"ExistingObjectTag/security"}} }] }"#, ) .expect("policy with value-only ExistingObjectTag text should parse"); assert!( !policy_uses_existing_object_tag_conditions(&with_value_only), "ExistingObjectTag text in values should not trigger tag dependency" ); let with_condition_key = Policy::parse_config( br#"{ "Version":"2012-10-17", "Statement":[{ "Effect":"Allow", "Action":["s3:GetObject"], "Resource":["arn:aws:s3:::bucket/*"], "Condition":{"StringEquals":{"s3:ExistingObjectTag/security":"public"}} }] }"#, ) .expect("policy with ExistingObjectTag condition key should parse"); assert!( policy_uses_existing_object_tag_conditions(&with_condition_key), "ExistingObjectTag condition key must trigger tag dependency" ); } #[test] fn test_policy_uses_existing_object_tag_when_only_secondary_action_has_tag_condition() { let split_action_policy = Policy::parse_config( br#"{ "Version":"2012-10-17", "Statement":[ { "Effect":"Allow", "Action":["s3:DeleteObject"], "Resource":["arn:aws:s3:::bucket/*"] }, { "Effect":"Allow", "Action":["s3:DeleteObjectVersion"], "Resource":["arn:aws:s3:::bucket/*"], "Condition":{"StringEquals":{"s3:ExistingObjectTag/security":"public"}} } ] }"#, ) .expect("split-action policy should parse"); assert!( policy_uses_existing_object_tag_conditions(&split_action_policy), "full merged policy must still be detectable as containing ExistingObjectTag keys" ); } #[tokio::test] async fn test_prepare_auth_detects_existing_object_tag_in_encoded_session_policy() { let store = StsTestMockStore::new(true); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let sts_access_key = "sts-session-tag-encoded-test-user"; let sts_user = UserIdentity::from(Credentials { access_key: sts_access_key.to_string(), secret_key: "longenoughsecret".to_string(), session_token: "sts-token".to_string(), status: ACCOUNT_ON.to_string(), parent_user: "sts-empty-parent-policy-test".to_string(), ..Default::default() }); iam_sys .store .cache .add_or_update_sts_account(sts_access_key, &sts_user, OffsetDateTime::now_utc()); let session_policy_json = r#"{ "Version":"2012-10-17", "Statement":[{"Effect":"Allow","Action":["s3:GetObject"],"Resource":["arn:aws:s3:::bucket/*"],"Condition":{"StringEquals":{"s3:ExistingObjectTag/security":"public"}}}] }"#; let mut claims = HashMap::new(); claims.insert( SESSION_POLICY_NAME.to_string(), Value::String(base64_simd::URL_SAFE_NO_PAD.encode_to_string(session_policy_json.as_bytes())), ); let groups: Option> = None; let args = Args { account: sts_access_key, groups: &groups, action: Action::S3Action(S3Action::GetObjectAction), bucket: "bucket", conditions: &HashMap::new(), is_owner: false, object: "obj", claims: &claims, deny_only: true, }; let prepared = iam_sys.prepare_auth(&args).await; assert!( prepared.needs_existing_object_tag, "base64 sessionPolicy with ExistingObjectTag must request object tag loading" ); } #[tokio::test] async fn test_prepare_auth_service_account_inherited_ignores_session_policy_tag_hint() { let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let service_account_access_key = "svc-inherited-tag-hint-test-user"; let parent_user = "sts-fallback-test-parent"; let mut service_account_claims = HashMap::new(); service_account_claims.insert(iam_policy_claim_name_sa(), Value::String(INHERITED_POLICY_TYPE.to_string())); let service_identity = UserIdentity::from(Credentials { access_key: service_account_access_key.to_string(), secret_key: "longenoughsecret".to_string(), status: ACCOUNT_ON.to_string(), parent_user: parent_user.to_string(), claims: Some(service_account_claims), ..Default::default() }); iam_sys .store .cache .add_or_update_user(service_account_access_key, &service_identity, OffsetDateTime::now_utc()); let mut request_claims = HashMap::new(); request_claims.insert("parent".to_string(), Value::String(parent_user.to_string())); request_claims.insert(iam_policy_claim_name_sa(), Value::String(INHERITED_POLICY_TYPE.to_string())); request_claims.insert( SESSION_POLICY_NAME_EXTRACTED.to_string(), Value::String( r#"{ "Version":"2012-10-17", "Statement":[{"Effect":"Allow","Action":["s3:GetObject"],"Resource":["arn:aws:s3:::bucket/*"],"Condition":{"StringEquals":{"s3:ExistingObjectTag/security":"public"}}}] }"# .to_string(), ), ); let groups: Option> = Some(vec!["testgroup".to_string()]); let args = Args { account: service_account_access_key, groups: &groups, action: Action::S3Action(S3Action::GetObjectAction), bucket: "bucket", conditions: &HashMap::new(), is_owner: false, object: "obj", claims: &request_claims, deny_only: false, }; let prepared = iam_sys.prepare_auth(&args).await; assert!( !prepared.needs_existing_object_tag, "inherited service account should not require object tag fetch based on session policy hint" ); } /// Regression test for rustfs#2392: `policy_db_get` must skip non-existent groups /// instead of aborting the entire policy resolution. When a JWT contains groups /// that exist in the IdP but not in IAM, policies from the remaining valid groups /// must still be returned. #[tokio::test] async fn test_policy_db_get_skips_nonexistent_groups() { let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); // "testgroup" exists with "readwrite" policy; "nonexistent-group" does not exist in IAM. let groups = Some(vec!["testgroup".to_string(), "nonexistent-group".to_string()]); let policies = iam_sys .policy_db_get("sts-fallback-test-parent", &groups) .await .expect("policy_db_get should not fail when some groups are missing"); assert!( policies.iter().any(|p| p == "readwrite"), "policies from existing group 'testgroup' should be returned even when other groups are missing; got: {:?}", policies ); } #[tokio::test] async fn test_info_policy_returns_policy_as_json_object() { let store = StsTestMockStore::new(false); let cache_manager = IamCache::new(store).await.unwrap(); let iam_sys = IamSys::new(cache_manager); let policy_info = iam_sys .info_policy("readonly") .await .expect("info_policy should return existing default policy"); assert!( policy_info.policy.is_object(), "policy field should be a JSON object for MinIO-compatible policy readback; got: {}", policy_info.policy ); assert!( policy_info.policy.get("Version").is_some(), "policy object should contain Version field; got: {}", policy_info.policy ); assert!( policy_info.policy.get("Statement").is_some(), "policy object should contain Statement field; got: {}", policy_info.policy ); } #[test] fn test_is_safe_claim_policy_name_allows_k8s_sa_sub() { assert!(IamSys::::is_safe_claim_policy_name( "system:serviceaccount:my-namespace:my-sa" )); } #[test] fn test_is_safe_claim_policy_name_allows_dotted_names() { assert!(IamSys::::is_safe_claim_policy_name("org.team.policy-name")); } #[test] fn test_is_safe_claim_policy_name_allows_simple_names() { assert!(IamSys::::is_safe_claim_policy_name("readwrite")); assert!(IamSys::::is_safe_claim_policy_name("my-custom_policy")); assert!(IamSys::::is_safe_claim_policy_name("Policy123")); } #[test] fn test_is_safe_claim_policy_name_rejects_empty() { assert!(!IamSys::::is_safe_claim_policy_name("")); } #[test] fn test_is_safe_claim_policy_name_rejects_path_traversal() { assert!(!IamSys::::is_safe_claim_policy_name("../etc/passwd")); assert!(!IamSys::::is_safe_claim_policy_name("policy/name")); assert!(!IamSys::::is_safe_claim_policy_name("policy\\name")); assert!(!IamSys::::is_safe_claim_policy_name("..\\etc\\passwd")); } #[test] fn test_is_safe_claim_policy_name_rejects_whitespace() { assert!(!IamSys::::is_safe_claim_policy_name("my policy")); assert!(!IamSys::::is_safe_claim_policy_name("policy\tname")); } #[test] fn test_is_safe_claim_policy_name_rejects_special_chars() { assert!(!IamSys::::is_safe_claim_policy_name("policy;drop")); assert!(!IamSys::::is_safe_claim_policy_name("pol$icy")); assert!(!IamSys::::is_safe_claim_policy_name("name{bad}")); } #[test] fn test_is_safe_claim_policy_name_rejects_no_alphanumeric() { assert!(!IamSys::::is_safe_claim_policy_name(".")); assert!(!IamSys::::is_safe_claim_policy_name("..")); assert!(!IamSys::::is_safe_claim_policy_name(":")); assert!(!IamSys::::is_safe_claim_policy_name("...")); assert!(!IamSys::::is_safe_claim_policy_name(".:.:")); assert!(!IamSys::::is_safe_claim_policy_name("-")); assert!(!IamSys::::is_safe_claim_policy_name("_")); assert!(!IamSys::::is_safe_claim_policy_name("__")); assert!(!IamSys::::is_safe_claim_policy_name("---")); assert!(!IamSys::::is_safe_claim_policy_name("-_-")); assert!(!IamSys::::is_safe_claim_policy_name("-.:_")); } #[tokio::test] async fn verified_federated_policy_rejects_malformed_session_boundary() { let iam_sys = IamSys::new( IamCache::new(StsTestMockStore::new(true)) .await .expect("IAM cache should initialize"), ); let parent_user = "virtual-parent"; let claims = HashMap::from([ ("iss".to_string(), Value::String("rustfs-oidc".to_string())), ("oidc_provider".to_string(), Value::String("default".to_string())), ("sub".to_string(), Value::String("subject-123".to_string())), ("parent".to_string(), Value::String(parent_user.to_string())), (OIDC_VIRTUAL_PARENT_CLAIM.to_string(), Value::String(parent_user.to_string())), (POLICYNAME.to_string(), Value::String("readwrite".to_string())), (SESSION_POLICY_NAME.to_string(), Value::Bool(true)), ]); let result = iam_sys .verified_federated_service_account_claims(&claims, parent_user, &None) .await; assert!(matches!(result, Err(IamError::InvalidArgument))); } }