From f9e8440a042c2647bf613f152c2682c0905e370d Mon Sep 17 00:00:00 2001 From: GatewayJ <835269233@qq.com> Date: Sun, 19 Jul 2026 14:28:08 +0800 Subject: [PATCH] refactor(iam): introduce federated identity boundary (#5018) --- crates/iam/src/federation/binding.rs | 24 + crates/iam/src/federation/error.rs | 39 ++ crates/iam/src/federation/mod.rs | 31 ++ crates/iam/src/federation/model.rs | 124 +++++ crates/iam/src/federation/oidc/claims.rs | 102 ++++ crates/iam/src/federation/oidc/config.rs | 36 ++ crates/iam/src/federation/oidc/discovery.rs | 23 + crates/iam/src/federation/oidc/flow.rs | 79 +++ crates/iam/src/federation/oidc/http.rs | 57 +++ crates/iam/src/federation/oidc/mod.rs | 21 + crates/iam/src/federation/provider.rs | 37 ++ crates/iam/src/federation/registry.rs | 31 ++ crates/iam/src/federation/transaction.rs | 307 ++++++++++++ crates/iam/src/lib.rs | 1 + rustfs/src/admin/console.rs | 9 +- rustfs/src/admin/handlers/oidc.rs | 153 +++--- rustfs/src/admin/handlers/site_replication.rs | 14 +- rustfs/src/admin/handlers/sts.rs | 448 +++--------------- rustfs/src/admin/runtime_sources.rs | 8 +- .../src/admin/service/federated_identity.rs | 341 +++++++++++++ rustfs/src/admin/service/mod.rs | 2 + rustfs/src/admin/service/session_policy.rs | 85 ++++ rustfs/src/app/context.rs | 83 ++-- rustfs/src/app/context/global.rs | 47 +- rustfs/src/app/context/handles.rs | 121 ++++- rustfs/src/app/context/interfaces.rs | 10 +- rustfs/src/app/context/runtime_sources.rs | 5 - rustfs/src/runtime_sources.rs | 8 +- rustfs/src/startup_auth.rs | 15 +- 29 files changed, 1707 insertions(+), 554 deletions(-) create mode 100644 crates/iam/src/federation/binding.rs create mode 100644 crates/iam/src/federation/error.rs create mode 100644 crates/iam/src/federation/mod.rs create mode 100644 crates/iam/src/federation/model.rs create mode 100644 crates/iam/src/federation/oidc/claims.rs create mode 100644 crates/iam/src/federation/oidc/config.rs create mode 100644 crates/iam/src/federation/oidc/discovery.rs create mode 100644 crates/iam/src/federation/oidc/flow.rs create mode 100644 crates/iam/src/federation/oidc/http.rs create mode 100644 crates/iam/src/federation/oidc/mod.rs create mode 100644 crates/iam/src/federation/provider.rs create mode 100644 crates/iam/src/federation/registry.rs create mode 100644 crates/iam/src/federation/transaction.rs create mode 100644 rustfs/src/admin/service/federated_identity.rs create mode 100644 rustfs/src/admin/service/session_policy.rs diff --git a/crates/iam/src/federation/binding.rs b/crates/iam/src/federation/binding.rs new file mode 100644 index 000000000..69daebe3d --- /dev/null +++ b/crates/iam/src/federation/binding.rs @@ -0,0 +1,24 @@ +// 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 super::{FederatedSessionBindingError, FederatedSessionTransaction}; +use rustfs_credentials::Credentials; + +#[async_trait::async_trait] +pub trait FederatedSessionBinding: Send + Sync { + async fn bind( + &self, + transaction: &FederatedSessionTransaction, + ) -> core::result::Result; +} diff --git a/crates/iam/src/federation/error.rs b/crates/iam/src/federation/error.rs new file mode 100644 index 000000000..bee369ec2 --- /dev/null +++ b/crates/iam/src/federation/error.rs @@ -0,0 +1,39 @@ +// 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. + +pub type Result = core::result::Result; + +#[derive(Debug, thiserror::Error)] +pub enum FederatedSessionBindingError { + #[error("{0}")] + InvalidRequest(String), + #[error("{0}")] + Internal(String), +} + +#[derive(Debug, thiserror::Error)] +pub enum FederationError { + #[error("{0}")] + Authorization(String), + #[error("{0}")] + CodeExchange(String), + #[error("{0}")] + TokenVerification(String), + #[error("{0}")] + Logout(String), + #[error("no policies are available for this OIDC token")] + NoAuthorizationContext, + #[error(transparent)] + Binding(#[from] FederatedSessionBindingError), +} diff --git a/crates/iam/src/federation/mod.rs b/crates/iam/src/federation/mod.rs new file mode 100644 index 000000000..aea1c90f3 --- /dev/null +++ b/crates/iam/src/federation/mod.rs @@ -0,0 +1,31 @@ +// 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. + +mod binding; +mod error; +mod model; +pub mod oidc; +mod provider; +mod registry; +mod transaction; + +pub use binding::FederatedSessionBinding; +pub use error::{FederatedSessionBindingError, FederationError, Result}; +pub use model::{ + FederatedAuthorization, FederatedClaims, FederatedCodeExchange, FederatedLoginSession, FederatedSession, + FederatedSessionTransaction, +}; +pub use provider::FederatedIdentityProvider; +pub use registry::FederatedIdentityRegistry; +pub use transaction::FederatedIdentityService; diff --git a/crates/iam/src/federation/model.rs b/crates/iam/src/federation/model.rs new file mode 100644 index 000000000..1645f8581 --- /dev/null +++ b/crates/iam/src/federation/model.rs @@ -0,0 +1,124 @@ +// 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 rustfs_credentials::Credentials; +use serde_json::Value; +use std::collections::HashMap; + +#[derive(Debug, Clone)] +pub struct FederatedClaims { + pub sub: String, + pub email: String, + pub username: String, + pub groups: Vec, + pub raw: HashMap, +} + +impl FederatedClaims { + pub fn session_identity(&self) -> String { + if !self.username.is_empty() { + self.username.clone() + } else if !self.email.is_empty() { + self.email.clone() + } else if !self.sub.is_empty() { + self.sub.clone() + } else { + "oidc-user-unknown".to_string() + } + } +} + +#[derive(Debug, Clone)] +pub struct FederatedAuthorization { + pub provider_id: String, + pub claims: FederatedClaims, + pub policies: Vec, + pub groups: Vec, + pub roles_claim_key: Option, + pub roles: Vec, +} + +impl FederatedAuthorization { + pub fn has_authorization_context(&self) -> bool { + !self.policies.is_empty() || !self.groups.is_empty() + } +} + +#[derive(Debug)] +pub struct FederatedCodeExchange { + pub authorization: FederatedAuthorization, + pub redirect_after: Option, + pub id_token: String, +} + +#[derive(Debug)] +pub struct FederatedSessionTransaction { + pub authorization: FederatedAuthorization, + pub duration_seconds: usize, + pub session_policy: Option, +} + +#[derive(Debug)] +pub struct FederatedSession { + pub credentials: Credentials, + pub authorization: FederatedAuthorization, +} + +#[derive(Debug)] +pub struct FederatedLoginSession { + pub session: FederatedSession, + pub redirect_after: Option, + pub logout_token: String, +} + +#[cfg(test)] +mod tests { + use super::*; + + fn claims(username: &str, email: &str, sub: &str) -> FederatedClaims { + FederatedClaims { + sub: sub.to_string(), + email: email.to_string(), + username: username.to_string(), + groups: Vec::new(), + raw: HashMap::new(), + } + } + + fn authorization(policies: Vec, groups: Vec) -> FederatedAuthorization { + FederatedAuthorization { + provider_id: "standard_oidc".to_string(), + claims: claims("", "", "subject"), + policies, + groups, + roles_claim_key: None, + roles: Vec::new(), + } + } + + #[test] + fn session_identity_preserves_existing_fallback_order() { + assert_eq!(claims("john", "john@example.com", "sub-1").session_identity(), "john"); + assert_eq!(claims("", "john@example.com", "sub-1").session_identity(), "john@example.com"); + assert_eq!(claims("", "", "sub-1").session_identity(), "sub-1"); + assert_eq!(claims("", "", "").session_identity(), "oidc-user-unknown"); + } + + #[test] + fn authorization_context_accepts_policy_or_group() { + assert!(!authorization(Vec::new(), Vec::new()).has_authorization_context()); + assert!(authorization(vec!["consoleAdmin".to_string()], Vec::new()).has_authorization_context()); + assert!(authorization(Vec::new(), vec!["RustFS.ConsoleAdmin".to_string()]).has_authorization_context()); + } +} diff --git a/crates/iam/src/federation/oidc/claims.rs b/crates/iam/src/federation/oidc/claims.rs new file mode 100644 index 000000000..e54afd62b --- /dev/null +++ b/crates/iam/src/federation/oidc/claims.rs @@ -0,0 +1,102 @@ +// 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 super::config::roles_claim_key; +use crate::{ + federation::{FederatedAuthorization, FederatedClaims}, + oidc::{OidcClaims, OidcSys}, +}; +use rustfs_policy::policy::{ClaimLookup, get_claim_case_insensitive}; + +fn string_list_claim(claims: &OidcClaims, claim_name: &str) -> Vec { + match get_claim_case_insensitive(&claims.raw, claim_name) { + ClaimLookup::Found(serde_json::Value::Array(values)) => values + .iter() + .filter_map(|value| value.as_str().map(ToOwned::to_owned)) + .collect(), + ClaimLookup::Found(serde_json::Value::String(value)) => value + .split(',') + .map(str::trim) + .filter(|value| !value.is_empty()) + .map(ToOwned::to_owned) + .collect(), + ClaimLookup::Missing | ClaimLookup::Ambiguous | ClaimLookup::Found(_) => Vec::new(), + } +} + +pub(super) fn authorization(oidc: &OidcSys, provider_id: String, claims: OidcClaims) -> FederatedAuthorization { + let (policies, groups) = oidc.map_claims_to_policies(&provider_id, &claims); + let roles_claim_key = roles_claim_key(oidc, &provider_id); + let roles = roles_claim_key + .as_deref() + .map(|claim_name| string_list_claim(&claims, claim_name)) + .unwrap_or_default(); + + FederatedAuthorization { + provider_id, + claims: FederatedClaims { + sub: claims.sub, + email: claims.email, + username: claims.username, + groups: claims.groups, + raw: claims.raw, + }, + policies, + groups, + roles_claim_key, + roles, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + use std::collections::HashMap; + + #[test] + fn string_list_claim_matches_existing_array_and_csv_behavior() { + let claims = OidcClaims { + raw: HashMap::from([ + ("array_roles".to_string(), json!(["reader", 7, "writer"])), + ("csv_roles".to_string(), json!("reader, writer, ,auditor")), + ]), + ..Default::default() + }; + + assert_eq!(string_list_claim(&claims, "array_roles"), ["reader", "writer"]); + assert_eq!(string_list_claim(&claims, "csv_roles"), ["reader", "writer", "auditor"]); + } + + #[test] + fn string_list_claim_preserves_exact_and_ambiguous_match_behavior() { + let exact = OidcClaims { + raw: HashMap::from([ + ("Roles".to_string(), json!(["mixed-case"])), + ("roles".to_string(), json!(["exact-match"])), + ]), + ..Default::default() + }; + assert_eq!(string_list_claim(&exact, "roles"), ["exact-match"]); + + let ambiguous = OidcClaims { + raw: HashMap::from([ + ("Roles".to_string(), json!(["mixed-case"])), + ("ROLES".to_string(), json!(["upper-case"])), + ]), + ..Default::default() + }; + assert!(string_list_claim(&ambiguous, "roles").is_empty()); + } +} diff --git a/crates/iam/src/federation/oidc/config.rs b/crates/iam/src/federation/oidc/config.rs new file mode 100644 index 000000000..26287b521 --- /dev/null +++ b/crates/iam/src/federation/oidc/config.rs @@ -0,0 +1,36 @@ +// 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::oidc::{OidcProviderConfig, OidcSys}; + +pub(super) fn provider_config<'a>(oidc: &'a OidcSys, id: &str) -> Option<&'a OidcProviderConfig> { + oidc.get_provider_config(id) +} + +pub(super) fn roles_claim_key(oidc: &OidcSys, provider_id: &str) -> Option { + provider_config(oidc, provider_id) + .map(|config| config.roles_claim.trim().to_string()) + .filter(|claim| !claim.is_empty()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn roles_claim_key_requires_explicit_provider_config() { + let oidc = OidcSys::empty().expect("empty OIDC configuration should be valid"); + assert_eq!(roles_claim_key(&oidc, "default"), None); + } +} diff --git a/crates/iam/src/federation/oidc/discovery.rs b/crates/iam/src/federation/oidc/discovery.rs new file mode 100644 index 000000000..3e883d67b --- /dev/null +++ b/crates/iam/src/federation/oidc/discovery.rs @@ -0,0 +1,23 @@ +// 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::oidc::{OidcProviderSummary, OidcSys}; + +pub(super) fn list_providers(oidc: &OidcSys) -> Vec { + oidc.list_providers() +} + +pub(super) fn list_visible_providers(oidc: &OidcSys) -> Vec { + oidc.list_visible_providers() +} diff --git a/crates/iam/src/federation/oidc/flow.rs b/crates/iam/src/federation/oidc/flow.rs new file mode 100644 index 000000000..04e9badd1 --- /dev/null +++ b/crates/iam/src/federation/oidc/flow.rs @@ -0,0 +1,79 @@ +// 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 super::{claims, config, discovery, http}; +use crate::{ + federation::{FederatedAuthorization, FederatedCodeExchange, FederatedIdentityProvider, FederationError, Result}, + oidc::{OidcProviderConfig, OidcProviderSummary, OidcSys}, +}; +use std::sync::Arc; + +pub struct StandardOidcAdapter { + oidc: Arc, +} + +impl StandardOidcAdapter { + pub fn new(oidc: Arc) -> Self { + Self { oidc } + } +} + +#[async_trait::async_trait] +impl FederatedIdentityProvider for StandardOidcAdapter { + fn has_providers(&self) -> bool { + self.oidc.has_providers() + } + + fn list_providers(&self) -> Vec { + discovery::list_providers(&self.oidc) + } + + fn list_visible_providers(&self) -> Vec { + discovery::list_visible_providers(&self.oidc) + } + + fn provider_config(&self, id: &str) -> Option<&OidcProviderConfig> { + config::provider_config(&self.oidc, id) + } + + async fn authorize_url(&self, provider_id: &str, redirect_uri: &str, redirect_after: Option) -> Result { + http::authorize_url(&self.oidc, provider_id, redirect_uri, redirect_after).await + } + + async fn exchange_code(&self, state: &str, code: &str, redirect_uri: &str) -> Result { + let (oidc_claims, provider_id, session, id_token) = http::exchange_code(&self.oidc, state, code, redirect_uri).await?; + Ok(FederatedCodeExchange { + authorization: claims::authorization(&self.oidc, provider_id, oidc_claims), + redirect_after: session.redirect_after, + id_token, + }) + } + + async fn verify_web_identity_token(&self, jwt: &str) -> Result { + let (oidc_claims, provider_id) = self + .oidc + .verify_web_identity_token(jwt) + .await + .map_err(FederationError::TokenVerification)?; + Ok(claims::authorization(&self.oidc, provider_id, oidc_claims)) + } + + async fn create_logout_token(&self, provider_id: &str, id_token: &str) -> Result { + http::create_logout_token(&self.oidc, provider_id, id_token).await + } + + async fn build_logout_url(&self, logout_token: &str, post_logout_redirect_uri: &str) -> Result> { + http::build_logout_url(&self.oidc, logout_token, post_logout_redirect_uri).await + } +} diff --git a/crates/iam/src/federation/oidc/http.rs b/crates/iam/src/federation/oidc/http.rs new file mode 100644 index 000000000..3b4f0dcd0 --- /dev/null +++ b/crates/iam/src/federation/oidc/http.rs @@ -0,0 +1,57 @@ +// 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::{ + federation::{FederationError, Result}, + oidc::{OidcClaims, OidcSys}, + oidc_state::OidcAuthSession, +}; + +pub(super) async fn authorize_url( + oidc: &OidcSys, + provider_id: &str, + redirect_uri: &str, + redirect_after: Option, +) -> Result { + oidc.authorize_url(provider_id, redirect_uri, redirect_after) + .await + .map_err(FederationError::Authorization) +} + +pub(super) async fn exchange_code( + oidc: &OidcSys, + state: &str, + code: &str, + redirect_uri: &str, +) -> Result<(OidcClaims, String, OidcAuthSession, String)> { + oidc.exchange_code(state, code, redirect_uri) + .await + .map_err(FederationError::CodeExchange) +} + +pub(super) async fn create_logout_token(oidc: &OidcSys, provider_id: &str, id_token: &str) -> Result { + oidc.create_logout_token(provider_id, id_token) + .await + .map_err(FederationError::Logout) +} + +pub(super) async fn build_logout_url( + oidc: &OidcSys, + logout_token: &str, + post_logout_redirect_uri: &str, +) -> Result> { + oidc.build_logout_url(logout_token, post_logout_redirect_uri) + .await + .map_err(FederationError::Logout) +} diff --git a/crates/iam/src/federation/oidc/mod.rs b/crates/iam/src/federation/oidc/mod.rs new file mode 100644 index 000000000..360c38831 --- /dev/null +++ b/crates/iam/src/federation/oidc/mod.rs @@ -0,0 +1,21 @@ +// 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. + +mod claims; +mod config; +mod discovery; +mod flow; +mod http; + +pub use flow::StandardOidcAdapter; diff --git a/crates/iam/src/federation/provider.rs b/crates/iam/src/federation/provider.rs new file mode 100644 index 000000000..3ee79598a --- /dev/null +++ b/crates/iam/src/federation/provider.rs @@ -0,0 +1,37 @@ +// 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 super::{FederatedAuthorization, FederatedCodeExchange, Result}; +use crate::oidc::{OidcProviderConfig, OidcProviderSummary}; + +#[async_trait::async_trait] +pub trait FederatedIdentityProvider: Send + Sync { + fn has_providers(&self) -> bool; + + fn list_providers(&self) -> Vec; + + fn list_visible_providers(&self) -> Vec; + + fn provider_config(&self, id: &str) -> Option<&OidcProviderConfig>; + + async fn authorize_url(&self, provider_id: &str, redirect_uri: &str, redirect_after: Option) -> Result; + + async fn exchange_code(&self, state: &str, code: &str, redirect_uri: &str) -> Result; + + async fn verify_web_identity_token(&self, jwt: &str) -> Result; + + async fn create_logout_token(&self, provider_id: &str, id_token: &str) -> Result; + + async fn build_logout_url(&self, logout_token: &str, post_logout_redirect_uri: &str) -> Result>; +} diff --git a/crates/iam/src/federation/registry.rs b/crates/iam/src/federation/registry.rs new file mode 100644 index 000000000..8c9f9c484 --- /dev/null +++ b/crates/iam/src/federation/registry.rs @@ -0,0 +1,31 @@ +// 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 super::FederatedIdentityProvider; +use std::sync::Arc; + +/// Immutable registry of built-in federation adapters. +pub struct FederatedIdentityRegistry { + standard_oidc: Arc, +} + +impl FederatedIdentityRegistry { + pub fn new(standard_oidc: Arc) -> Self { + Self { standard_oidc } + } + + pub(crate) fn standard_oidc(&self) -> &dyn FederatedIdentityProvider { + self.standard_oidc.as_ref() + } +} diff --git a/crates/iam/src/federation/transaction.rs b/crates/iam/src/federation/transaction.rs new file mode 100644 index 000000000..a0f168df5 --- /dev/null +++ b/crates/iam/src/federation/transaction.rs @@ -0,0 +1,307 @@ +// 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 super::{ + FederatedIdentityRegistry, FederatedLoginSession, FederatedSession, FederatedSessionBinding, FederatedSessionTransaction, + FederationError, Result, +}; +use crate::oidc::{OidcProviderConfig, OidcProviderSummary}; + +pub struct FederatedIdentityService { + registry: FederatedIdentityRegistry, +} + +impl FederatedIdentityService { + pub fn new(registry: FederatedIdentityRegistry) -> Self { + Self { registry } + } + + pub fn has_providers(&self) -> bool { + self.registry.standard_oidc().has_providers() + } + + pub fn list_providers(&self) -> Vec { + self.registry.standard_oidc().list_providers() + } + + pub fn list_visible_providers(&self) -> Vec { + self.registry.standard_oidc().list_visible_providers() + } + + pub fn get_provider_config(&self, id: &str) -> Option<&OidcProviderConfig> { + self.registry.standard_oidc().provider_config(id) + } + + pub async fn authorize_url(&self, provider_id: &str, redirect_uri: &str, redirect_after: Option) -> Result { + self.registry + .standard_oidc() + .authorize_url(provider_id, redirect_uri, redirect_after) + .await + } + + pub async fn complete_authorization_code( + &self, + state: &str, + code: &str, + redirect_uri: &str, + duration_seconds: usize, + binding: &dyn FederatedSessionBinding, + ) -> Result { + let exchange = self.registry.standard_oidc().exchange_code(state, code, redirect_uri).await?; + let provider_id = exchange.authorization.provider_id.clone(); + let transaction = FederatedSessionTransaction { + authorization: exchange.authorization, + duration_seconds, + session_policy: None, + }; + let credentials = binding.bind(&transaction).await?; + let logout_token = self + .registry + .standard_oidc() + .create_logout_token(&provider_id, &exchange.id_token) + .await?; + + Ok(FederatedLoginSession { + session: FederatedSession { + credentials, + authorization: transaction.authorization, + }, + redirect_after: exchange.redirect_after, + logout_token, + }) + } + + pub async fn assume_role_with_web_identity( + &self, + jwt: &str, + duration_seconds: usize, + session_policy: Option, + binding: &dyn FederatedSessionBinding, + ) -> Result { + let authorization = self.registry.standard_oidc().verify_web_identity_token(jwt).await?; + if !authorization.has_authorization_context() { + tracing::warn!( + provider_id = %authorization.provider_id, + username = %authorization.claims.username, + sub = %authorization.claims.sub, + policy_count = authorization.policies.len(), + group_count = authorization.groups.len(), + "AssumeRoleWithWebIdentity has no mapped policies or groups" + ); + return Err(FederationError::NoAuthorizationContext); + } + tracing::debug!( + provider_id = %authorization.provider_id, + username = %authorization.claims.username, + policy_count = authorization.policies.len(), + group_count = authorization.groups.len(), + policies = ?authorization.policies, + groups = ?authorization.groups, + "AssumeRoleWithWebIdentity mapped OIDC policies and groups" + ); + + let transaction = FederatedSessionTransaction { + authorization, + duration_seconds, + session_policy, + }; + let credentials = binding.bind(&transaction).await?; + + Ok(FederatedSession { + credentials, + authorization: transaction.authorization, + }) + } + + pub async fn build_logout_url(&self, logout_token: &str, post_logout_redirect_uri: &str) -> Result> { + self.registry + .standard_oidc() + .build_logout_url(logout_token, post_logout_redirect_uri) + .await + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::federation::{ + FederatedAuthorization, FederatedClaims, FederatedCodeExchange, FederatedIdentityProvider, FederatedSessionBindingError, + }; + use crate::oidc::{OidcProviderConfig, OidcProviderSummary}; + use rustfs_credentials::Credentials; + use std::sync::{ + Arc, Mutex, + atomic::{AtomicUsize, Ordering}, + }; + + struct TestProvider { + with_policy: bool, + events: Arc>>, + } + + impl TestProvider { + fn record(&self, event: &'static str) { + self.events.lock().expect("event log should not be poisoned").push(event); + } + + fn authorization(&self) -> FederatedAuthorization { + FederatedAuthorization { + provider_id: "default".to_string(), + claims: FederatedClaims { + sub: "subject".to_string(), + email: String::new(), + username: "user".to_string(), + groups: Vec::new(), + raw: Default::default(), + }, + policies: if self.with_policy { + vec!["readwrite".to_string()] + } else { + Vec::new() + }, + groups: Vec::new(), + roles_claim_key: None, + roles: Vec::new(), + } + } + } + + #[async_trait::async_trait] + impl FederatedIdentityProvider for TestProvider { + fn has_providers(&self) -> bool { + true + } + + fn list_providers(&self) -> Vec { + Vec::new() + } + + fn list_visible_providers(&self) -> Vec { + Vec::new() + } + + fn provider_config(&self, _id: &str) -> Option<&OidcProviderConfig> { + None + } + + async fn authorize_url( + &self, + _provider_id: &str, + _redirect_uri: &str, + _redirect_after: Option, + ) -> Result { + Ok("https://identity.example/authorize".to_string()) + } + + async fn exchange_code(&self, _state: &str, _code: &str, _redirect_uri: &str) -> Result { + self.record("exchange"); + Ok(FederatedCodeExchange { + authorization: self.authorization(), + redirect_after: Some("/browser".to_string()), + id_token: "id-token".to_string(), + }) + } + + async fn verify_web_identity_token(&self, _jwt: &str) -> Result { + self.record("verify"); + Ok(self.authorization()) + } + + async fn create_logout_token(&self, _provider_id: &str, _id_token: &str) -> Result { + self.record("logout"); + Ok("logout-token".to_string()) + } + + async fn build_logout_url(&self, _logout_token: &str, _post_logout_redirect_uri: &str) -> Result> { + Ok(Some("https://identity.example/logout".to_string())) + } + } + + struct CountingBinding { + calls: AtomicUsize, + events: Arc>>, + } + + #[async_trait::async_trait] + impl FederatedSessionBinding for CountingBinding { + async fn bind( + &self, + transaction: &FederatedSessionTransaction, + ) -> core::result::Result { + self.calls.fetch_add(1, Ordering::Relaxed); + self.events.lock().expect("event log should not be poisoned").push("bind"); + Ok(Credentials { + access_key: transaction.authorization.claims.session_identity(), + ..Default::default() + }) + } + } + + #[tokio::test] + async fn callback_and_web_identity_share_session_binding() { + let events = Arc::new(Mutex::new(Vec::new())); + let provider = Arc::new(TestProvider { + with_policy: true, + events: events.clone(), + }); + let binding = Arc::new(CountingBinding { + calls: AtomicUsize::new(0), + events: events.clone(), + }); + let service = FederatedIdentityService::new(FederatedIdentityRegistry::new(provider)); + + let login = service + .complete_authorization_code("state", "code", "https://console.example/callback", 3600, binding.as_ref()) + .await + .expect("callback flow should complete"); + assert_eq!(login.session.credentials.access_key, "user"); + assert_eq!(login.redirect_after.as_deref(), Some("/browser")); + assert_eq!(login.logout_token, "logout-token"); + assert_eq!( + events.lock().expect("event log should not be poisoned").as_slice(), + ["exchange", "bind", "logout"] + ); + events.lock().expect("event log should not be poisoned").clear(); + + let web_identity = service + .assume_role_with_web_identity("jwt", 3600, None, binding.as_ref()) + .await + .expect("web identity flow should complete"); + assert_eq!(web_identity.credentials.access_key, "user"); + assert_eq!(binding.calls.load(Ordering::Relaxed), 2); + assert_eq!(events.lock().expect("event log should not be poisoned").as_slice(), ["verify", "bind"]); + } + + #[tokio::test] + async fn web_identity_without_policy_or_group_is_not_bound() { + let events = Arc::new(Mutex::new(Vec::new())); + let provider = Arc::new(TestProvider { + with_policy: false, + events: events.clone(), + }); + let binding = Arc::new(CountingBinding { + calls: AtomicUsize::new(0), + events, + }); + let service = FederatedIdentityService::new(FederatedIdentityRegistry::new(provider)); + + let error = service + .assume_role_with_web_identity("jwt", 3600, None, binding.as_ref()) + .await + .expect_err("authorization context is required"); + + assert!(matches!(error, FederationError::NoAuthorizationContext)); + assert_eq!(binding.calls.load(Ordering::Relaxed), 0); + } +} diff --git a/crates/iam/src/lib.rs b/crates/iam/src/lib.rs index 1f49f481a..3f54282cd 100644 --- a/crates/iam/src/lib.rs +++ b/crates/iam/src/lib.rs @@ -28,6 +28,7 @@ const EVENT_OIDC_STATE: &str = "oidc_state"; pub mod cache; pub mod error; +pub mod federation; pub mod keyring; pub mod manager; pub mod oidc; diff --git a/rustfs/src/admin/console.rs b/rustfs/src/admin/console.rs index b86c61431..d4ba186f2 100644 --- a/rustfs/src/admin/console.rs +++ b/rustfs/src/admin/console.rs @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -use crate::admin::runtime_sources::{current_oidc_handle, default_admin_usecase}; +use crate::admin::runtime_sources::{current_federated_identity_service, default_admin_usecase}; use crate::admin::storage_api::access::RequestContext; use crate::license::has_valid_license; use crate::server::has_path_prefix; @@ -139,9 +139,10 @@ impl Config { let http_prefix = rustfs_config::RUSTFS_HTTP_PREFIX; // Collect OIDC provider info if available - let oidc = current_oidc_handle() - .map(|sys| { - sys.list_visible_providers() + let oidc = current_federated_identity_service() + .map(|federation| { + federation + .list_visible_providers() .into_iter() .map(|p| OidcProviderInfo { provider_id: p.provider_id, diff --git a/rustfs/src/admin/handlers/oidc.rs b/rustfs/src/admin/handlers/oidc.rs index 70d6a8599..9e991223f 100644 --- a/rustfs/src/admin/handlers/oidc.rs +++ b/rustfs/src/admin/handlers/oidc.rs @@ -12,12 +12,13 @@ // See the License for the specific language governing permissions and // limitations under the License. -use super::sts::create_oidc_sts_credentials; use crate::admin::auth::validate_admin_request; use crate::admin::router::{AdminOperation, Operation, S3Router}; use crate::admin::runtime_sources::{ - current_app_context, current_object_store_handle_for_context, current_oidc_handle, current_server_config_for_context, + current_app_context, current_federated_identity_service, current_object_store_handle_for_context, + current_server_config_for_context, }; +use crate::admin::service::federated_identity::DefaultFederatedSessionBinding; use crate::admin::storage_api::config::{read_admin_config_without_migrate, save_admin_server_config}; use crate::auth::{check_key_valid, get_session_token}; use crate::server::{ADMIN_PREFIX, CONSOLE_PREFIX, MINIO_ADMIN_PREFIX, RemoteAddr}; @@ -32,6 +33,7 @@ use rustfs_config::oidc::{ }; use rustfs_config::server_config::Config as ServerConfig; use rustfs_config::{DEFAULT_DELIMITER, ENABLE_KEY, ENV_RUSTFS_BROWSER_REDIRECT_URL, EnableState, MAX_ADMIN_REQUEST_BODY_SIZE}; +use rustfs_iam::federation::{FederatedSessionBindingError, FederationError}; use rustfs_policy::policy::action::{Action, AdminAction}; use rustfs_utils::egress::validate_outbound_url; use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; @@ -54,6 +56,24 @@ const CONSOLE_LOGIN_SUFFIX: &str = "/auth/login"; const OIDC_STATE_LB_HINT: &str = "check load balancer session affinity for OIDC authorize/callback requests or configure RUSTFS_BROWSER_REDIRECT_URL"; +fn callback_federation_error(error: FederationError) -> S3Error { + match error { + FederationError::CodeExchange(message) => { + S3Error::with_message(S3ErrorCode::AccessDenied, format!("code exchange failed: {message}")) + } + FederationError::Logout(message) => { + S3Error::with_message(S3ErrorCode::InternalError, format!("logout session creation failed: {message}")) + } + FederationError::Binding(FederatedSessionBindingError::InvalidRequest(message)) => { + S3Error::with_message(S3ErrorCode::InvalidRequest, message) + } + FederationError::Binding(FederatedSessionBindingError::Internal(message)) => { + S3Error::with_message(S3ErrorCode::InternalError, message) + } + other => S3Error::with_message(S3ErrorCode::InternalError, other.to_string()), + } +} + /// Validate that a provider ID contains only safe characters (alphanumeric, underscore, hyphen). fn is_valid_provider_id(id: &str) -> bool { !id.is_empty() && id.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-') @@ -274,9 +294,9 @@ pub struct ListOidcProvidersHandler {} #[async_trait::async_trait] impl Operation for ListOidcProvidersHandler { async fn call(&self, _req: S3Request, _params: Params<'_, '_>) -> S3Result> { - let oidc_sys = current_oidc_handle().ok_or_else(|| s3_error!(InternalError, "OIDC not initialized"))?; + let federation = current_federated_identity_service().ok_or_else(|| s3_error!(InternalError, "OIDC not initialized"))?; - let providers = oidc_sys.list_visible_providers(); + let providers = federation.list_visible_providers(); let json_body = serde_json::to_vec(&providers) .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("serialize error: {e}")))?; @@ -445,7 +465,7 @@ impl Operation for OidcAuthorizeHandler { return Err(s3_error!(InvalidRequest, "invalid provider_id")); } - let oidc_sys = current_oidc_handle().ok_or_else(|| s3_error!(InternalError, "OIDC not initialized"))?; + let federation = current_federated_identity_service().ok_or_else(|| s3_error!(InternalError, "OIDC not initialized"))?; // Derive the callback redirect URI from the request let redirect_uri = derive_callback_uri(&req, provider_id)?; @@ -454,7 +474,7 @@ impl Operation for OidcAuthorizeHandler { let redirect_after = extract_safe_redirect_after(&req.uri)?; let redirect_after_log = redirect_after.clone(); - let auth_url = oidc_sys + let auth_url = federation .authorize_url(provider_id, &redirect_uri, redirect_after) .await .map_err(|e| { @@ -535,35 +555,40 @@ impl Operation for OidcCallbackHandler { let state = extract_query_param(&req.uri, "state").ok_or_else(|| s3_error!(InvalidRequest, "missing 'state' query parameter"))?; - let oidc_sys = current_oidc_handle().ok_or_else(|| s3_error!(InternalError, "OIDC not initialized"))?; + let federation = current_federated_identity_service().ok_or_else(|| s3_error!(InternalError, "OIDC not initialized"))?; let redirect_uri = derive_callback_uri(&req, provider_id)?; - // Exchange authorization code for tokens and extract claims - let (claims, actual_provider_id, session, id_token) = - oidc_sys.exchange_code(&state, &code, &redirect_uri).await.map_err(|e| { - let lb_hint = if is_invalid_oidc_state_error(&e) { - OIDC_STATE_LB_HINT - } else { - "" - }; - error!( - event = EVENT_ADMIN_OIDC_STATE, - component = LOG_COMPONENT_ADMIN, - subsystem = LOG_SUBSYSTEM_OIDC, - result = "code_exchange_failed", - requested_provider_id = %provider_id, - redirect_uri = %redirect_uri, - code = %code, - state = %state, - code_len = code.len(), - state_len = state.len(), - error = %e, - lb_hint = %lb_hint, - "admin oidc state" - ); - S3Error::with_message(S3ErrorCode::AccessDenied, format!("code exchange failed: {e}")) + let login = federation + .complete_authorization_code(&state, &code, &redirect_uri, 3600, &DefaultFederatedSessionBinding) + .await + .map_err(|error| { + if let FederationError::CodeExchange(message) = &error { + let lb_hint = if is_invalid_oidc_state_error(message) { + OIDC_STATE_LB_HINT + } else { + "" + }; + error!( + event = EVENT_ADMIN_OIDC_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_OIDC, + result = "code_exchange_failed", + requested_provider_id = %provider_id, + redirect_uri = %redirect_uri, + code = %code, + state = %state, + code_len = code.len(), + state_len = state.len(), + error = %message, + lb_hint = %lb_hint, + "admin oidc state" + ); + } + callback_federation_error(error) })?; + let authorization = &login.session.authorization; + let actual_provider_id = &authorization.provider_id; debug!( event = EVENT_ADMIN_OIDC_STATE, @@ -574,32 +599,20 @@ impl Operation for OidcCallbackHandler { "admin oidc state" ); - // Map claims to policies and groups - let (policies, groups) = oidc_sys.map_claims_to_policies(&actual_provider_id, &claims); - debug!( event = EVENT_ADMIN_OIDC_STATE, component = LOG_COMPONENT_ADMIN, subsystem = LOG_SUBSYSTEM_OIDC, provider_id = %actual_provider_id, - policy_count = policies.len(), - group_count = groups.len(), - policies = ?policies, - groups = ?groups, + policy_count = authorization.policies.len(), + group_count = authorization.groups.len(), + policies = ?authorization.policies, + groups = ?authorization.groups, state = "claims_mapped", "admin oidc state" ); - // Generate STS credentials using the shared helper. - // Console/OIDC sessions use a fixed 1-hour duration as a security/UX choice. - // Longer-lived credentials (15 min to 12 hours) can be requested via CLI/SDK - // through AssumeRoleWithWebIdentity. - let new_cred = create_oidc_sts_credentials(&claims, &actual_provider_id, &policies, &groups, 3600, None).await?; - - let logout_token = oidc_sys - .create_logout_token(&actual_provider_id, &id_token) - .await - .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("logout session creation failed: {e}")))?; + let new_cred = &login.session.credentials; // Build redirect URL to console with credentials in the fragment let console_redirect = build_console_redirect( @@ -608,8 +621,8 @@ impl Operation for OidcCallbackHandler { &new_cred.secret_key, &new_cred.session_token, new_cred.expiration, - session.redirect_after.as_deref(), - Some(logout_token.as_str()), + login.redirect_after.as_deref(), + Some(login.logout_token.as_str()), )?; let mut resp = S3Response::new((StatusCode::FOUND, Body::empty())); @@ -636,8 +649,8 @@ impl Operation for OidcLogoutHandler { return redirect_response(&fallback_location); }; - let location = match current_oidc_handle() { - Some(oidc_sys) => match oidc_sys.build_logout_url(&logout_token, &fallback_location).await { + let location = match current_federated_identity_service() { + Some(federation) => match federation.build_logout_url(&logout_token, &fallback_location).await { Ok(Some(url)) => url, Ok(None) => fallback_location.clone(), Err(err) => { @@ -664,8 +677,8 @@ impl Operation for OidcLogoutHandler { /// from request headers. For production deployments behind a reverse proxy, configuring /// an explicit redirect_uri is recommended to prevent header manipulation. fn derive_callback_uri(req: &S3Request, provider_id: &str) -> S3Result { - if let Some(oidc_sys) = current_oidc_handle() - && let Some(config) = oidc_sys.get_provider_config(provider_id) + if let Some(federation) = current_federated_identity_service() + && let Some(config) = federation.get_provider_config(provider_id) { return derive_callback_uri_with_provider_config(req, provider_id, Some(config)); } @@ -1326,6 +1339,38 @@ mod tests { assert_eq!(extract_query_param(&uri, "missing"), None); } + #[test] + fn callback_errors_preserve_existing_s3_semantics() { + let cases = [ + ( + FederationError::CodeExchange("invalid state".to_string()), + S3ErrorCode::AccessDenied, + "code exchange failed: invalid state", + ), + ( + FederationError::Logout("session unavailable".to_string()), + S3ErrorCode::InternalError, + "logout session creation failed: session unavailable", + ), + ( + FederationError::Binding(FederatedSessionBindingError::InvalidRequest("invalid policy".to_string())), + S3ErrorCode::InvalidRequest, + "invalid policy", + ), + ( + FederationError::Binding(FederatedSessionBindingError::Internal("failed to store temp user".to_string())), + S3ErrorCode::InternalError, + "failed to store temp user", + ), + ]; + + for (error, expected_code, expected_message) in cases { + let error = callback_federation_error(error); + assert_eq!(error.code(), &expected_code); + assert_eq!(error.message(), Some(expected_message)); + } + } + #[test] fn test_extract_query_param_empty() { let uri: http::Uri = "http://localhost/callback".parse().unwrap(); diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index 9bf982028..03f758941 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -15,10 +15,10 @@ use crate::admin::auth::validate_admin_request; use crate::admin::router::{AdminOperation, Operation, S3Router}; use crate::admin::runtime_sources::{ - current_deployment_id, current_endpoints_handle, current_iam_handle, current_object_store_handle, current_oidc_handle, - current_outbound_tls_generation, current_outbound_tls_state, current_region, current_replication_pool_handle, - current_replication_stats_handle, current_runtime_port, current_server_config, current_token_signing_key, - object_store_from_req, + current_deployment_id, current_endpoints_handle, current_federated_identity_service, current_iam_handle, + current_object_store_handle, current_outbound_tls_generation, current_outbound_tls_state, current_region, + current_replication_pool_handle, current_replication_stats_handle, current_runtime_port, current_server_config, + current_token_signing_key, object_store_from_req, }; use crate::admin::site_replication_identity::{ canonical_endpoint, deployment_id_for_endpoint, normalize_peer_map_by_identity_with, same_identity_endpoint, @@ -2987,13 +2987,13 @@ async fn build_sr_info(state: &SiteReplicationState, local_peer: &PeerInfo) -> S fn local_idp_settings() -> IDPSettings { let mut settings = IDPSettings::default(); - if let Some(oidc) = current_oidc_handle() { - let providers = oidc.list_providers(); + if let Some(federation) = current_federated_identity_service() { + let providers = federation.list_providers(); settings.open_id.enabled = !providers.is_empty(); settings.open_id.region = current_region().map(|region| region.to_string()).unwrap_or_default(); for provider in providers { - let Some(config) = oidc.get_provider_config(&provider.provider_id) else { + let Some(config) = federation.get_provider_config(&provider.provider_id) else { continue; }; let provider_settings = OpenIDProviderSettings { diff --git a/rustfs/src/admin/handlers/sts.rs b/rustfs/src/admin/handlers/sts.rs index e0c94101e..e9e1cb24f 100644 --- a/rustfs/src/admin/handlers/sts.rs +++ b/rustfs/src/admin/handlers/sts.rs @@ -13,14 +13,16 @@ // limitations under the License. use super::is_admin::IsAdminHandler; +use crate::admin::service::federated_identity::DefaultFederatedSessionBinding; +use crate::admin::service::session_policy::populate_session_policy; use crate::admin::storage_api::bucket::utils::serialize; use crate::{ - admin::runtime_sources::{current_action_credentials, current_oidc_handle, current_token_signing_key}, + admin::runtime_sources::{current_action_credentials, current_federated_identity_service, current_token_signing_key}, admin::{ handlers::site_replication::site_replication_iam_change_hook, router::{AdminOperation, Operation, S3Router}, }, - auth::{check_key_valid, extract_string_list_claim, get_session_token}, + auth::{check_key_valid, get_session_token}, server::ADMIN_PREFIX, server::RemoteAddr, }; @@ -29,12 +31,12 @@ use http::header::HeaderValue; use hyper::Method; use matchit::Params; use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE; -use rustfs_iam::{oidc::OidcClaims, sys::SESSION_POLICY_NAME}; +use rustfs_iam::federation::{FederatedSessionBindingError, FederationError}; use rustfs_madmin::{SITE_REPL_API_VERSION, SRIAMItem, SRSTSCredential}; use rustfs_policy::{ auth::get_new_credentials_with_metadata, policy::{ - Args, Policy, + Args, action::{Action, StsAction}, }, }; @@ -46,7 +48,6 @@ use s3s::{ use serde::Deserialize; use serde_json::Value; use serde_urlencoded::from_bytes; -use std::collections::{BTreeMap, HashMap}; use time::{Duration, OffsetDateTime}; use tracing::{debug, error, info, warn}; @@ -75,138 +76,21 @@ fn clamp_assume_role_duration(duration_seconds: usize) -> usize { } } -fn has_identity_authorization_context(policies: &[String], groups: &[String]) -> bool { - !policies.is_empty() || !groups.is_empty() -} - -fn configured_roles_claim_key(provider_id: &str) -> Option { - current_oidc_handle() - .as_ref() - .and_then(|oidc_sys| oidc_sys.get_provider_config(provider_id)) - .map(|cfg| cfg.roles_claim.trim().to_string()) - .filter(|claim| !claim.is_empty()) -} - -fn build_oidc_token_claims( - claims: &OidcClaims, - provider_id: &str, - groups: &[String], - roles_claim_key: Option<&str>, -) -> HashMap { - let mut token_claims: HashMap = HashMap::new(); - token_claims.insert("sub".to_string(), Value::String(claims.sub.clone())); - token_claims.insert("iss".to_string(), Value::String("rustfs-oidc".to_string())); - token_claims.insert("oidc_provider".to_string(), Value::String(provider_id.to_string())); - - if !claims.email.is_empty() { - token_claims.insert("email".to_string(), Value::String(claims.email.clone())); - } - if !claims.username.is_empty() { - token_claims.insert("preferred_username".to_string(), Value::String(claims.username.clone())); - } - if !groups.is_empty() { - token_claims.insert( - "groups".to_string(), - Value::Array(groups.iter().map(|g| Value::String(g.clone())).collect()), - ); - } - if let Some(roles_claim_key) = roles_claim_key { - let roles = extract_string_list_claim(&claims.raw, roles_claim_key); - if !roles.is_empty() { - token_claims.insert("roles".to_string(), Value::Array(roles.into_iter().map(Value::String).collect())); +fn web_identity_federation_error(error: FederationError) -> S3Error { + match error { + FederationError::TokenVerification(message) => { + S3Error::with_message(S3ErrorCode::AccessDenied, format!("token verification failed: {message}")) } - } - token_claims -} - -fn resolve_oidc_session_identity(claims: &OidcClaims) -> String { - if !claims.username.is_empty() { - claims.username.clone() - } else if !claims.email.is_empty() { - claims.email.clone() - } else if !claims.sub.is_empty() { - claims.sub.clone() - } else { - "oidc-user-unknown".to_string() - } -} - -async fn log_oidc_policy_diagnostics( - iam_store: &rustfs_iam::sys::IamSys, - provider_id: &str, - parent_user: &str, - policies: &[String], - groups: &[String], -) { - if policies.is_empty() { - let policy_documents = BTreeMap::::new(); - let missing_policies = Vec::::new(); - let combined_policy = Value::Null; - debug!( - provider_id = %provider_id, - parent_user = %parent_user, - policy_count = 0, - group_count = groups.len(), - policies = ?policies, - groups = ?groups, - policy_documents = ?policy_documents, - missing_policies = ?missing_policies, - combined_policy = ?combined_policy, - "OIDC STS policy diagnostics" - ); - return; - } - - match iam_store.list_policy_docs("").await { - Ok(policy_docs) => { - let mut policy_documents = BTreeMap::new(); - let mut missing_policies = Vec::new(); - for policy_name in policies { - match policy_docs.get(policy_name) { - Some(policy_doc) => { - let policy_doc_json = serde_json::to_value(policy_doc).unwrap_or_else(|err| { - serde_json::json!({ - "serialization_error": err.to_string(), - }) - }); - policy_documents.insert(policy_name.clone(), policy_doc_json); - } - None => missing_policies.push(policy_name.clone()), - } - } - - let combined_policy = iam_store.get_combined_policy(policies).await; - let combined_policy_json = serde_json::to_value(&combined_policy).unwrap_or_else(|err| { - serde_json::json!({ - "serialization_error": err.to_string(), - }) - }); - - debug!( - provider_id = %provider_id, - parent_user = %parent_user, - policy_count = policies.len(), - group_count = groups.len(), - policies = ?policies, - groups = ?groups, - missing_policies = ?missing_policies, - policy_documents = ?policy_documents, - combined_policy = ?combined_policy_json, - "OIDC STS policy diagnostics" - ); + FederationError::NoAuthorizationContext => { + s3_error!(InvalidArgument, "no policies are available for this OIDC token") } - Err(err) => { - warn!( - provider_id = %provider_id, - parent_user = %parent_user, - policy_count = policies.len(), - group_count = groups.len(), - policies = ?policies, - groups = ?groups, - error = %err, - "OIDC STS policy diagnostics failed" - ); + FederationError::Binding(FederatedSessionBindingError::InvalidRequest(message)) => { + S3Error::with_message(S3ErrorCode::InvalidRequest, message) } + FederationError::Binding(FederatedSessionBindingError::Internal(message)) => { + S3Error::with_message(S3ErrorCode::InternalError, message) + } + other => S3Error::with_message(S3ErrorCode::InternalError, other.to_string()), } } @@ -411,42 +295,6 @@ async fn handle_assume_role_with_web_identity(body: AssumeRoleRequest) -> S3Resu return Err(s3_error!(InvalidArgument, "not support version")); } - // Verify the JWT and extract claims - let oidc_sys = current_oidc_handle().ok_or_else(|| s3_error!(InternalError, "OIDC not initialized"))?; - - let (claims, provider_id) = oidc_sys - .verify_web_identity_token(&body.web_identity_token) - .await - .map_err(|e| { - warn!("AssumeRoleWithWebIdentity JWT verification failed: {}", e); - S3Error::with_message(S3ErrorCode::AccessDenied, format!("token verification failed: {e}")) - })?; - - // Map claims to policies and groups - let (policies, groups) = oidc_sys.map_claims_to_policies(&provider_id, &claims); - - if !has_identity_authorization_context(&policies, &groups) { - warn!( - provider_id = %provider_id, - username = %claims.username, - sub = %claims.sub, - policy_count = policies.len(), - group_count = groups.len(), - "AssumeRoleWithWebIdentity has no mapped policies or groups" - ); - return Err(s3_error!(InvalidArgument, "no policies are available for this OIDC token")); - } - - debug!( - provider_id = %provider_id, - username = %claims.username, - policy_count = policies.len(), - group_count = groups.len(), - policies = ?policies, - groups = ?groups, - "AssumeRoleWithWebIdentity mapped OIDC policies and groups" - ); - let mut duration = if body.duration_seconds > 0 { body.duration_seconds } else { @@ -455,18 +303,24 @@ async fn handle_assume_role_with_web_identity(body: AssumeRoleRequest) -> S3Resu // Enforce reasonable bounds for STS credentials duration (similar to AWS STS) duration = duration.clamp(900, 43200); - // Generate STS credentials using the shared helper - let new_cred = create_oidc_sts_credentials( - &claims, - &provider_id, - &policies, - &groups, - duration, - if body.policy.is_empty() { None } else { Some(&body.policy) }, - ) - .await?; - - let subject = resolve_oidc_session_identity(&claims); + let federation = current_federated_identity_service().ok_or_else(|| s3_error!(InternalError, "OIDC not initialized"))?; + let session = federation + .assume_role_with_web_identity( + &body.web_identity_token, + duration, + if body.policy.is_empty() { None } else { Some(body.policy) }, + &DefaultFederatedSessionBinding, + ) + .await + .map_err(|error| { + if let FederationError::TokenVerification(message) = &error { + warn!("AssumeRoleWithWebIdentity JWT verification failed: {}", message); + } + web_identity_federation_error(error) + })?; + let authorization = &session.authorization; + let new_cred = &session.credentials; + let subject = authorization.claims.session_identity(); // Build XML response (AssumeRoleWithWebIdentityResponse) let expiration = new_cred @@ -502,121 +356,6 @@ async fn handle_assume_role_with_web_identity(body: AssumeRoleRequest) -> S3Resu Ok(resp) } -/// Shared helper to generate STS credentials from OIDC claims. -/// Used by both the OIDC callback handler and AssumeRoleWithWebIdentity. -pub async fn create_oidc_sts_credentials( - claims: &OidcClaims, - provider_id: &str, - policies: &[String], - groups: &[String], - duration_seconds: usize, - session_policy: Option<&str>, -) -> S3Result { - let roles_claim_key = configured_roles_claim_key(provider_id); - let mut token_claims = build_oidc_token_claims(claims, provider_id, groups, roles_claim_key.as_deref()); - - // Set expiration - let exp = OffsetDateTime::now_utc().saturating_add(Duration::seconds(duration_seconds as i64)); - token_claims.insert("exp".to_string(), Value::Number(serde_json::Number::from(exp.unix_timestamp()))); - - // Set the parent user: prefer username, then email, then sub - let parent_user = resolve_oidc_session_identity(claims); - debug!( - provider_id = %provider_id, - parent_user = %parent_user, - email = %claims.email, - username = %claims.username, - sub = %claims.sub, - policy_count = policies.len(), - group_count = groups.len(), - policies = ?policies, - groups = ?groups, - roles_claim_key = ?roles_claim_key, - has_session_policy = session_policy.is_some(), - "OIDC STS credential claims prepared" - ); - token_claims.insert("parent".to_string(), Value::String(parent_user.clone())); - - // Set policies as a comma-separated string - if !policies.is_empty() { - token_claims.insert("policy".to_string(), Value::String(policies.join(","))); - } - - // Optionally apply session policy - if let Some(policy_str) = session_policy { - populate_session_policy(&mut token_claims, policy_str)?; - } - - // Generate STS temp credentials - let secret = current_token_signing_key().ok_or_else(|| s3_error!(InternalError, "token signing key not initialized"))?; - - let mut new_cred = get_new_credentials_with_metadata(&token_claims, &secret) - .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("credential generation failed: {e}")))?; - - new_cred.parent_user = parent_user; - new_cred.groups = Some(groups.to_vec()); - - // Store temp user in IAM - let iam_store = - crate::admin::runtime_sources::current_ready_iam_handle().map_err(|_| s3_error!(InternalError, "IAM not initialized"))?; - if tracing::enabled!(tracing::Level::DEBUG) { - log_oidc_policy_diagnostics(&iam_store, provider_id, &new_cred.parent_user, policies, groups).await; - } - - let updated_at = iam_store - .set_temp_user(&new_cred.access_key, &new_cred, None) - .await - .map_err(|_| s3_error!(InternalError, "failed to store temp user"))?; - - if let Err(err) = site_replication_iam_change_hook(SRIAMItem { - r#type: "sts-credential".to_string(), - sts_credential: Some(SRSTSCredential { - access_key: new_cred.access_key.clone(), - secret_key: new_cred.secret_key.clone(), - session_token: new_cred.session_token.clone(), - parent_user: new_cred.parent_user.clone(), - parent_policy_mapping: policies.join(","), - api_version: Some(SITE_REPL_API_VERSION.to_string()), - }), - updated_at: Some(updated_at), - api_version: Some(SITE_REPL_API_VERSION.to_string()), - ..Default::default() - }) - .await - { - warn!("site replication OIDC STS hook failed, err: {err}"); - } - - Ok(new_cred) -} - -pub fn populate_session_policy(claims: &mut HashMap, policy: &str) -> S3Result<()> { - if !policy.is_empty() { - let session_policy = Policy::parse_config(policy.as_bytes()) - .map_err(|e| { - let error_msg = format!("Failed to parse session policy: {}. Please check that the policy is valid JSON format with standard brackets [] for arrays.", e); - S3Error::with_message(S3ErrorCode::InvalidRequest, error_msg) - })?; - if session_policy.version.is_empty() { - return Err(s3_error!(InvalidRequest, "invalid policy")); - } - - let policy_buf = serde_json::to_vec(&session_policy) - .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("marshal policy err {e}")))?; - - if policy_buf.len() > 2048 { - return Err(s3_error!(InvalidRequest, "policy too large")); - } - - claims.insert( - SESSION_POLICY_NAME.to_string(), - Value::String(base64_simd::URL_SAFE_NO_PAD.encode_to_string(&policy_buf)), - ); - } - - Ok(()) -} - /// Escape special XML characters in a string. fn xml_escape(s: &str) -> String { // Fast path: if there are no escapable characters, just clone the string. @@ -684,93 +423,34 @@ mod tests { } #[test] - fn test_has_identity_authorization_context() { - let empty: Vec = vec![]; - let groups = vec!["RustFS.ConsoleAdmin".to_string()]; - let policies = vec!["consoleAdmin".to_string()]; + fn web_identity_errors_preserve_existing_s3_semantics() { + let cases = [ + ( + FederationError::TokenVerification("invalid token".to_string()), + S3ErrorCode::AccessDenied, + "token verification failed: invalid token", + ), + ( + FederationError::NoAuthorizationContext, + S3ErrorCode::InvalidArgument, + "no policies are available for this OIDC token", + ), + ( + FederationError::Binding(FederatedSessionBindingError::InvalidRequest("invalid policy".to_string())), + S3ErrorCode::InvalidRequest, + "invalid policy", + ), + ( + FederationError::Binding(FederatedSessionBindingError::Internal("failed to store temp user".to_string())), + S3ErrorCode::InternalError, + "failed to store temp user", + ), + ]; - assert!(!has_identity_authorization_context(&empty, &empty)); - assert!(has_identity_authorization_context(&policies, &empty)); - assert!(has_identity_authorization_context(&empty, &groups)); - } - - #[test] - fn test_extract_string_list_claim_supports_array_and_csv() { - let mut claims = HashMap::new(); - claims.insert("roles".to_string(), serde_json::json!(["admin", "reader"])); - claims.insert("groups".to_string(), serde_json::json!("devs, ops")); - - assert_eq!(extract_string_list_claim(&claims, "roles"), vec!["admin", "reader"]); - assert_eq!(extract_string_list_claim(&claims, "groups"), vec!["devs", "ops"]); - } - - #[test] - fn test_extract_string_list_claim_prefers_exact_match() { - let mut claims = HashMap::new(); - claims.insert("Roles".to_string(), serde_json::json!(["mixed-case"])); - claims.insert("roles".to_string(), serde_json::json!(["exact-match"])); - - assert_eq!(extract_string_list_claim(&claims, "roles"), vec!["exact-match"]); - } - - #[test] - fn test_extract_string_list_claim_ambiguous_case_insensitive_match_returns_empty() { - let mut claims = HashMap::new(); - claims.insert("Roles".to_string(), serde_json::json!(["mixed-case"])); - claims.insert("ROLES".to_string(), serde_json::json!(["upper-case"])); - - assert!(extract_string_list_claim(&claims, "roles").is_empty()); - } - - #[test] - fn test_build_oidc_token_claims_includes_normalized_roles() { - let mut raw = HashMap::new(); - raw.insert("Roles".to_string(), serde_json::json!("admin, reader")); - let claims = OidcClaims { - sub: "user-sub".to_string(), - raw, - ..Default::default() - }; - let token_claims = build_oidc_token_claims(&claims, "default", &["devs".to_string()], Some("roles")); - - assert_eq!(token_claims.get("roles"), Some(&serde_json::json!(["admin", "reader"]))); - } - - #[test] - fn test_configured_roles_claim_key_requires_explicit_config() { - assert_eq!(configured_roles_claim_key("default"), None); - } - - #[test] - fn test_resolve_oidc_session_identity_prefers_username_over_email() { - let claims = OidcClaims { - username: "john".to_string(), - email: "john@example.com".to_string(), - sub: "sub-1".to_string(), - ..Default::default() - }; - - assert_eq!(resolve_oidc_session_identity(&claims), "john"); - } - - #[test] - fn test_resolve_oidc_session_identity_falls_back_to_email_then_sub() { - let claims_with_email = OidcClaims { - email: "john@example.com".to_string(), - sub: "sub-1".to_string(), - ..Default::default() - }; - assert_eq!(resolve_oidc_session_identity(&claims_with_email), "john@example.com"); - - let claims_with_sub = OidcClaims { - sub: "sub-1".to_string(), - ..Default::default() - }; - assert_eq!(resolve_oidc_session_identity(&claims_with_sub), "sub-1"); - } - - #[test] - fn test_resolve_oidc_session_identity_uses_unknown_when_all_empty() { - assert_eq!(resolve_oidc_session_identity(&OidcClaims::default()), "oidc-user-unknown"); + for (error, expected_code, expected_message) in cases { + let error = web_identity_federation_error(error); + assert_eq!(error.code(), &expected_code); + assert_eq!(error.message(), Some(expected_message)); + } } } diff --git a/rustfs/src/admin/runtime_sources.rs b/rustfs/src/admin/runtime_sources.rs index 8c824b451..50e6d9f28 100644 --- a/rustfs/src/admin/runtime_sources.rs +++ b/rustfs/src/admin/runtime_sources.rs @@ -23,10 +23,10 @@ use crate::app::object_usecase::DefaultObjectUsecase; use crate::runtime_sources as root_runtime_sources; pub(crate) use crate::runtime_sources::{ AppContext, ServerContextSlot, current_action_credentials, current_boot_time, current_bucket_metadata_handle, - current_bucket_monitor_handle, current_deployment_id, current_endpoints_handle, current_iam_handle, - current_kms_runtime_service_manager, current_notification_system_for_context, current_object_data_cache_handle_for_context, - current_object_store_handle_for_context, current_oidc_handle, current_ready_iam_handle, current_region, - current_replication_pool_handle, current_replication_stats_handle, current_server_config_for_context, + current_bucket_monitor_handle, current_deployment_id, current_endpoints_handle, current_federated_identity_service, + current_iam_handle, current_kms_runtime_service_manager, current_notification_system_for_context, + current_object_data_cache_handle_for_context, current_object_store_handle_for_context, current_ready_iam_handle, + current_region, current_replication_pool_handle, current_replication_stats_handle, current_server_config_for_context, current_token_signing_key, }; use rustfs_config::server_config::Config; diff --git a/rustfs/src/admin/service/federated_identity.rs b/rustfs/src/admin/service/federated_identity.rs new file mode 100644 index 000000000..1f46144ba --- /dev/null +++ b/rustfs/src/admin/service/federated_identity.rs @@ -0,0 +1,341 @@ +// 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 super::session_policy::populate_session_policy; +use crate::admin::{ + handlers::site_replication::site_replication_iam_change_hook, + runtime_sources::{current_ready_iam_handle, current_token_signing_key}, +}; +use rustfs_iam::{ + federation::{FederatedSessionBinding, FederatedSessionBindingError, FederatedSessionTransaction}, + store::object::ObjectStore, + sys::IamSys, +}; +use rustfs_madmin::{SITE_REPL_API_VERSION, SRIAMItem, SRSTSCredential}; +use rustfs_policy::auth::get_new_credentials_with_metadata; +use s3s::{S3Error, S3ErrorCode}; +use serde_json::Value; +use std::collections::{BTreeMap, HashMap}; +use time::{Duration, OffsetDateTime}; +use tracing::{debug, warn}; + +pub(crate) struct DefaultFederatedSessionBinding; + +fn build_oidc_token_claims(transaction: &FederatedSessionTransaction) -> HashMap { + let authorization = &transaction.authorization; + let claims = &authorization.claims; + let mut token_claims = HashMap::new(); + token_claims.insert("sub".to_string(), Value::String(claims.sub.clone())); + token_claims.insert("iss".to_string(), Value::String("rustfs-oidc".to_string())); + token_claims.insert("oidc_provider".to_string(), Value::String(authorization.provider_id.clone())); + + if !claims.email.is_empty() { + token_claims.insert("email".to_string(), Value::String(claims.email.clone())); + } + if !claims.username.is_empty() { + token_claims.insert("preferred_username".to_string(), Value::String(claims.username.clone())); + } + if !authorization.groups.is_empty() { + token_claims.insert( + "groups".to_string(), + Value::Array( + authorization + .groups + .iter() + .map(|group| Value::String(group.clone())) + .collect(), + ), + ); + } + if !authorization.roles.is_empty() { + token_claims.insert( + "roles".to_string(), + Value::Array(authorization.roles.iter().cloned().map(Value::String).collect()), + ); + } + + token_claims +} + +async fn log_oidc_policy_diagnostics( + iam_store: &IamSys, + provider_id: &str, + parent_user: &str, + policies: &[String], + groups: &[String], +) { + if policies.is_empty() { + let policy_documents = BTreeMap::::new(); + let missing_policies = Vec::::new(); + let combined_policy = Value::Null; + debug!( + provider_id = %provider_id, + parent_user = %parent_user, + policy_count = 0, + group_count = groups.len(), + policies = ?policies, + groups = ?groups, + policy_documents = ?policy_documents, + missing_policies = ?missing_policies, + combined_policy = ?combined_policy, + "OIDC STS policy diagnostics" + ); + return; + } + + match iam_store.list_policy_docs("").await { + Ok(policy_docs) => { + let mut policy_documents = BTreeMap::new(); + let mut missing_policies = Vec::new(); + for policy_name in policies { + match policy_docs.get(policy_name) { + Some(policy_doc) => { + let policy_doc_json = serde_json::to_value(policy_doc).unwrap_or_else(|err| { + serde_json::json!({ + "serialization_error": err.to_string(), + }) + }); + policy_documents.insert(policy_name.clone(), policy_doc_json); + } + None => missing_policies.push(policy_name.clone()), + } + } + + let combined_policy = iam_store.get_combined_policy(policies).await; + let combined_policy_json = serde_json::to_value(&combined_policy).unwrap_or_else(|err| { + serde_json::json!({ + "serialization_error": err.to_string(), + }) + }); + + debug!( + provider_id = %provider_id, + parent_user = %parent_user, + policy_count = policies.len(), + group_count = groups.len(), + policies = ?policies, + groups = ?groups, + policy_documents = ?policy_documents, + missing_policies = ?missing_policies, + combined_policy = ?combined_policy_json, + "OIDC STS policy diagnostics" + ); + } + Err(err) => { + warn!( + provider_id = %provider_id, + parent_user = %parent_user, + policy_count = policies.len(), + group_count = groups.len(), + policies = ?policies, + groups = ?groups, + error = %err, + "OIDC STS policy diagnostics failed" + ); + } + } +} + +fn binding_error_from_s3(error: S3Error) -> FederatedSessionBindingError { + let message = error.message().unwrap_or("federated session binding failed").to_string(); + if error.code() == &S3ErrorCode::InvalidRequest { + FederatedSessionBindingError::InvalidRequest(message) + } else { + FederatedSessionBindingError::Internal(message) + } +} + +fn issue_credentials( + transaction: &FederatedSessionTransaction, + secret: Option<&str>, +) -> Result { + let authorization = &transaction.authorization; + let claims = &authorization.claims; + let mut token_claims = build_oidc_token_claims(transaction); + let duration = i64::try_from(transaction.duration_seconds) + .map_err(|_| FederatedSessionBindingError::InvalidRequest("invalid duration".to_string()))?; + let exp = OffsetDateTime::now_utc().saturating_add(Duration::seconds(duration)); + token_claims.insert("exp".to_string(), Value::Number(serde_json::Number::from(exp.unix_timestamp()))); + + let parent_user = claims.session_identity(); + debug!( + provider_id = %authorization.provider_id, + parent_user = %parent_user, + email = %claims.email, + username = %claims.username, + sub = %claims.sub, + policy_count = authorization.policies.len(), + group_count = authorization.groups.len(), + policies = ?authorization.policies, + groups = ?authorization.groups, + roles_claim_key = ?authorization.roles_claim_key, + has_session_policy = transaction.session_policy.is_some(), + "OIDC STS credential claims prepared" + ); + token_claims.insert("parent".to_string(), Value::String(parent_user.clone())); + + if !authorization.policies.is_empty() { + token_claims.insert("policy".to_string(), Value::String(authorization.policies.join(","))); + } + if let Some(policy) = transaction.session_policy.as_deref() { + populate_session_policy(&mut token_claims, policy).map_err(binding_error_from_s3)?; + } + + let secret = secret.ok_or_else(|| FederatedSessionBindingError::Internal("token signing key not initialized".to_string()))?; + let mut credentials = get_new_credentials_with_metadata(&token_claims, secret) + .map_err(|error| FederatedSessionBindingError::Internal(format!("credential generation failed: {error}")))?; + credentials.parent_user = parent_user; + credentials.groups = Some(authorization.groups.clone()); + Ok(credentials) +} + +fn site_replication_item( + credentials: &rustfs_credentials::Credentials, + transaction: &FederatedSessionTransaction, + updated_at: OffsetDateTime, +) -> SRIAMItem { + SRIAMItem { + r#type: "sts-credential".to_string(), + sts_credential: Some(SRSTSCredential { + access_key: credentials.access_key.clone(), + secret_key: credentials.secret_key.clone(), + session_token: credentials.session_token.clone(), + parent_user: credentials.parent_user.clone(), + parent_policy_mapping: transaction.authorization.policies.join(","), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + }), + updated_at: Some(updated_at), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + ..Default::default() + } +} + +#[async_trait::async_trait] +impl FederatedSessionBinding for DefaultFederatedSessionBinding { + async fn bind( + &self, + transaction: &FederatedSessionTransaction, + ) -> Result { + let authorization = &transaction.authorization; + let secret = current_token_signing_key(); + let credentials = issue_credentials(transaction, secret.as_deref())?; + + let iam_store = + current_ready_iam_handle().map_err(|_| FederatedSessionBindingError::Internal("IAM not initialized".to_string()))?; + if tracing::enabled!(tracing::Level::DEBUG) { + log_oidc_policy_diagnostics( + &iam_store, + &authorization.provider_id, + &credentials.parent_user, + &authorization.policies, + &authorization.groups, + ) + .await; + } + + let updated_at = iam_store + .set_temp_user(&credentials.access_key, &credentials, None) + .await + .map_err(|_| FederatedSessionBindingError::Internal("failed to store temp user".to_string()))?; + + if let Err(err) = site_replication_iam_change_hook(site_replication_item(&credentials, transaction, updated_at)).await { + warn!("site replication OIDC STS hook failed, err: {err}"); + } + + Ok(credentials) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use rustfs_iam::federation::{FederatedAuthorization, FederatedClaims}; + + fn transaction() -> FederatedSessionTransaction { + FederatedSessionTransaction { + authorization: FederatedAuthorization { + provider_id: "default".to_string(), + claims: FederatedClaims { + sub: "subject".to_string(), + email: "user@example.com".to_string(), + username: "user".to_string(), + groups: vec!["source-group".to_string()], + raw: HashMap::new(), + }, + policies: vec!["readwrite".to_string()], + groups: vec!["devs".to_string()], + roles_claim_key: Some("roles".to_string()), + roles: vec!["admin".to_string(), "reader".to_string()], + }, + duration_seconds: 3600, + session_policy: None, + } + } + + #[test] + fn token_claims_preserve_existing_oidc_shape() { + let transaction = transaction(); + + let claims = build_oidc_token_claims(&transaction); + assert_eq!(claims.get("sub"), Some(&serde_json::json!("subject"))); + assert_eq!(claims.get("iss"), Some(&serde_json::json!("rustfs-oidc"))); + assert_eq!(claims.get("oidc_provider"), Some(&serde_json::json!("default"))); + assert_eq!(claims.get("email"), Some(&serde_json::json!("user@example.com"))); + assert_eq!(claims.get("preferred_username"), Some(&serde_json::json!("user"))); + assert_eq!(claims.get("groups"), Some(&serde_json::json!(["devs"]))); + assert_eq!(claims.get("roles"), Some(&serde_json::json!(["admin", "reader"]))); + } + + #[test] + fn issued_credentials_and_replication_item_preserve_existing_shape() { + let transaction = transaction(); + let secret = "federated-session-test-signing-secret"; + + let credentials = issue_credentials(&transaction, Some(secret)).expect("credential issuance should succeed"); + assert_eq!(credentials.parent_user, "user"); + assert_eq!(credentials.groups, Some(vec!["devs".to_string()])); + + let claims = rustfs_iam::sys::get_claims_from_token_with_secret(&credentials.session_token, secret) + .expect("issued session token should verify"); + assert_eq!(claims.get("iss"), Some(&serde_json::json!("rustfs-oidc"))); + assert_eq!(claims.get("oidc_provider"), Some(&serde_json::json!("default"))); + assert_eq!(claims.get("parent"), Some(&serde_json::json!("user"))); + assert_eq!(claims.get("policy"), Some(&serde_json::json!("readwrite"))); + assert_eq!(claims.get("groups"), Some(&serde_json::json!(["devs"]))); + assert_eq!(claims.get("roles"), Some(&serde_json::json!(["admin", "reader"]))); + assert!(!claims.contains_key("oidc_issuer")); + + let updated_at = OffsetDateTime::UNIX_EPOCH; + let item = site_replication_item(&credentials, &transaction, updated_at); + assert_eq!(item.r#type, "sts-credential"); + assert_eq!(item.updated_at, Some(updated_at)); + assert_eq!(item.api_version.as_deref(), Some(SITE_REPL_API_VERSION)); + let replicated = item.sts_credential.expect("replication item should contain STS credentials"); + assert_eq!(replicated.access_key, credentials.access_key); + assert_eq!(replicated.secret_key, credentials.secret_key); + assert_eq!(replicated.session_token, credentials.session_token); + assert_eq!(replicated.parent_user, "user"); + assert_eq!(replicated.parent_policy_mapping, "readwrite"); + assert_eq!(replicated.api_version.as_deref(), Some(SITE_REPL_API_VERSION)); + } + + #[test] + fn invalid_session_policy_precedes_missing_signing_key() { + let mut transaction = transaction(); + transaction.session_policy = Some("not-json".to_string()); + + let error = issue_credentials(&transaction, None).expect_err("invalid policy should fail first"); + assert!(matches!(error, FederatedSessionBindingError::InvalidRequest(_))); + } +} diff --git a/rustfs/src/admin/service/mod.rs b/rustfs/src/admin/service/mod.rs index 5af8b3f7e..569802ed0 100644 --- a/rustfs/src/admin/service/mod.rs +++ b/rustfs/src/admin/service/mod.rs @@ -13,4 +13,6 @@ // limitations under the License. pub mod config; +pub(crate) mod federated_identity; +pub(crate) mod session_policy; pub mod site_replication; diff --git a/rustfs/src/admin/service/session_policy.rs b/rustfs/src/admin/service/session_policy.rs new file mode 100644 index 000000000..67c894982 --- /dev/null +++ b/rustfs/src/admin/service/session_policy.rs @@ -0,0 +1,85 @@ +// 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 rustfs_iam::sys::SESSION_POLICY_NAME; +use rustfs_policy::policy::Policy; +use s3s::{S3Error, S3ErrorCode, S3Result, s3_error}; +use serde_json::Value; +use std::collections::HashMap; + +pub(crate) fn populate_session_policy(claims: &mut HashMap, policy: &str) -> S3Result<()> { + if !policy.is_empty() { + let session_policy = Policy::parse_config(policy.as_bytes()).map_err(|e| { + let error_msg = format!( + "Failed to parse session policy: {}. Please check that the policy is valid JSON format with standard brackets [] for arrays.", + e + ); + S3Error::with_message(S3ErrorCode::InvalidRequest, error_msg) + })?; + if session_policy.version.is_empty() { + return Err(s3_error!(InvalidRequest, "invalid policy")); + } + + let policy_buf = serde_json::to_vec(&session_policy) + .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("marshal policy err {e}")))?; + + if policy_buf.len() > 2048 { + return Err(s3_error!(InvalidRequest, "policy too large")); + } + + claims.insert( + SESSION_POLICY_NAME.to_string(), + Value::String(base64_simd::URL_SAFE_NO_PAD.encode_to_string(&policy_buf)), + ); + } + + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn empty_policy_does_not_change_claims() { + let mut claims = HashMap::new(); + populate_session_policy(&mut claims, "").expect("empty policy should be accepted"); + assert!(claims.is_empty()); + } + + #[test] + fn valid_policy_is_encoded_in_the_session_claim() { + let mut claims = HashMap::new(); + let policy = r#"{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Action":["s3:GetObject"],"Resource":["arn:aws:s3:::bucket/*"]}]}"#; + + populate_session_policy(&mut claims, policy).expect("valid policy should be accepted"); + + let encoded = claims + .get(SESSION_POLICY_NAME) + .and_then(Value::as_str) + .expect("session policy claim should be present"); + let decoded = base64_simd::URL_SAFE_NO_PAD + .decode_to_vec(encoded.as_bytes()) + .expect("session policy claim should be base64url encoded"); + let decoded_policy: Value = serde_json::from_slice(&decoded).expect("session policy claim should contain JSON"); + assert_eq!(decoded_policy["Version"], "2012-10-17"); + } + + #[test] + fn invalid_policy_is_rejected_as_invalid_request() { + let mut claims = HashMap::new(); + let error = populate_session_policy(&mut claims, "not-json").expect_err("invalid policy should fail"); + assert_eq!(error.code(), &S3ErrorCode::InvalidRequest); + } +} diff --git a/rustfs/src/app/context.rs b/rustfs/src/app/context.rs index dd0ca2977..65cbf2fbb 100644 --- a/rustfs/src/app/context.rs +++ b/rustfs/src/app/context.rs @@ -38,7 +38,7 @@ use crate::app::object_data_cache::ObjectDataCacheAdapter; use crate::config::RustFSBufferConfig; use rustfs_config::server_config::Config; use rustfs_credentials::Credentials; -use rustfs_iam::{error::Error as IamError, oidc::OidcSys, store::object::ObjectStore, sys::IamSys}; +use rustfs_iam::{error::Error as IamError, federation::FederatedIdentityService, store::object::ObjectStore, sys::IamSys}; use rustfs_io_metrics::{PerformanceMetrics, internode_metrics::InternodeMetrics}; use rustfs_kms::{KmsServiceManager, ObjectEncryptionService}; use rustfs_lock::LockClient; @@ -88,14 +88,14 @@ pub fn resolve_iam_handle() -> Option>> { resolve_iam_handle_with(get_global_app_context()) } -/// Resolve OIDC system handle using AppContext-first precedence. -pub fn resolve_oidc_handle() -> Option> { - resolve_oidc_handle_with(get_global_app_context()) +/// Resolve the federated identity service using AppContext-first precedence. +pub fn resolve_federated_identity_service() -> Option> { + resolve_federated_identity_service_with(get_global_app_context()) } -/// Publish the initialized OIDC system handle into the global AppContext. -pub fn publish_oidc_handle(oidc: Arc) -> bool { - publish_oidc_handle_with(get_global_app_context(), oidc) +/// Publish the initialized federated identity service into the global AppContext. +pub fn publish_federated_identity_service(service: Arc) -> bool { + handles::publish_default_federated_identity_service(service) } /// Resolve a ready IAM system handle using AppContext-first precedence. @@ -329,12 +329,13 @@ fn resolve_iam_handle_with(context: Option>) -> Option>) -> Option> { - context.and_then(|context| context.oidc().handle()) +fn resolve_federated_identity_service_with(context: Option>) -> Option> { + context.and_then(|context| context.federated_identity().handle()) } -fn publish_oidc_handle_with(context: Option>, oidc: Arc) -> bool { - context.is_some_and(|context| context.publish_oidc_handle(oidc)) +#[cfg(test)] +fn publish_federated_identity_service_with(context: Option>, service: Arc) -> bool { + context.is_some_and(|context| context.publish_federated_identity_service(service)) } fn resolve_ready_iam_handle_with(context: Arc) -> rustfs_iam::error::Result>> { @@ -509,14 +510,19 @@ mod tests { }; use crate::app::context::interfaces::{ ActionCredentialInterface, BootTimeInterface, BucketMetadataInterface, BufferConfigInterface, DeploymentIdInterface, - EndpointsInterface, IamInterface, InternodeMetricsInterface, KmsInterface, KmsRuntimeInterface, LocalNodeNameInterface, - LockClientInterface, LockClientsInterface, OidcInterface, OutboundTlsRuntimeInterface, PerformanceMetricsInterface, - RegionInterface, ReplicationStatsInterface, RuntimePortInterface, S3SelectDbInterface, ScannerMetricsInterface, - ServerConfigInterface, StorageClassInterface, TierConfigInterface, TransitionStateInterface, + EndpointsInterface, FederatedIdentityInterface, IamInterface, InternodeMetricsInterface, KmsInterface, + KmsRuntimeInterface, LocalNodeNameInterface, LockClientInterface, LockClientsInterface, OutboundTlsRuntimeInterface, + PerformanceMetricsInterface, RegionInterface, ReplicationStatsInterface, RuntimePortInterface, S3SelectDbInterface, + ScannerMetricsInterface, ServerConfigInterface, StorageClassInterface, TierConfigInterface, TransitionStateInterface, }; use crate::config::{RustFSBufferConfig, WorkloadProfile}; use async_trait::async_trait; - use rustfs_iam::{oidc::OidcSys, store::object::ObjectStore, sys::IamSys}; + use rustfs_iam::{ + federation::{FederatedIdentityRegistry, FederatedIdentityService, oidc::StandardOidcAdapter}, + oidc::OidcSys, + store::object::ObjectStore, + sys::IamSys, + }; use rustfs_io_metrics::{PerformanceMetrics, internode_metrics::InternodeMetrics}; use rustfs_lock::{LocalClient, LockClient}; use rustfs_s3select_api::{ @@ -552,32 +558,37 @@ mod tests { } } - struct TestOidcInterface { - oidc: StdRwLock>>, + struct TestFederatedIdentityInterface { + service: StdRwLock>>, } - impl TestOidcInterface { - fn new(oidc: Option>) -> Self { + impl TestFederatedIdentityInterface { + fn new(service: Option>) -> Self { Self { - oidc: StdRwLock::new(oidc), + service: StdRwLock::new(service), } } } - impl OidcInterface for TestOidcInterface { - fn handle(&self) -> Option> { - self.oidc.read().ok().and_then(|oidc| oidc.as_ref().cloned()) + impl FederatedIdentityInterface for TestFederatedIdentityInterface { + fn handle(&self) -> Option> { + self.service.read().ok().and_then(|service| service.as_ref().cloned()) } - fn publish_handle(&self, oidc: Arc) -> bool { - let Ok(mut published_oidc) = self.oidc.write() else { + fn publish_handle(&self, service: Arc) -> bool { + let Ok(mut published_service) = self.service.write() else { return false; }; - *published_oidc = Some(oidc); + *published_service = Some(service); true } } + fn test_federated_identity_service(oidc: OidcSys) -> Arc { + let adapter = Arc::new(StandardOidcAdapter::new(Arc::new(oidc))); + Arc::new(FederatedIdentityService::new(FederatedIdentityRegistry::new(adapter))) + } + struct TestKmsInterface { kms: Arc, } @@ -1001,8 +1012,8 @@ mod tests { Ok(sys) => sys, Err(err) => unreachable!("test OIDC fallback sys should initialize: {err}"), }; - let context_oidc = Arc::new(context_oidc_sys); - let fallback_oidc = Arc::new(fallback_oidc_sys); + let context_oidc = test_federated_identity_service(context_oidc_sys); + let fallback_oidc = test_federated_identity_service(fallback_oidc_sys); let context_token_signing_key = "context-token-signing-key".to_string(); let context = Arc::new(AppContext::with_test_interfaces( @@ -1012,7 +1023,7 @@ mod tests { ready: true, token_signing_key: Some(context_token_signing_key.clone()), }), - oidc: Arc::new(TestOidcInterface::new(Some(context_oidc.clone()))), + federated_identity: Arc::new(TestFederatedIdentityInterface::new(None)), kms: Arc::new(TestKmsInterface { kms: context_kms.clone(), }), @@ -1112,10 +1123,12 @@ mod tests { context_outbound_tls_state.generation ); assert!(resolve_iam_ready_with(Some(context.clone())).expect("context IAM ready")); - let resolved_oidc = resolve_oidc_handle_with(Some(context.clone())); + assert!(resolve_federated_identity_service_with(Some(context.clone())).is_none()); + assert!(publish_federated_identity_service_with(Some(context.clone()), context_oidc.clone())); + let resolved_oidc = resolve_federated_identity_service_with(Some(context.clone())); assert!(resolved_oidc.as_ref().is_some_and(|oidc| Arc::ptr_eq(oidc, &context_oidc))); - assert!(publish_oidc_handle_with(Some(context.clone()), fallback_oidc.clone())); - let resolved_oidc = resolve_oidc_handle_with(Some(context.clone())); + assert!(publish_federated_identity_service_with(Some(context.clone()), fallback_oidc.clone())); + let resolved_oidc = resolve_federated_identity_service_with(Some(context.clone())); assert!(resolved_oidc.as_ref().is_some_and(|oidc| Arc::ptr_eq(oidc, &fallback_oidc))); assert_eq!( resolve_token_signing_key_with(Some(context.clone())).as_deref(), @@ -1236,10 +1249,10 @@ mod tests { assert!(resolve_outbound_tls_state_with(None).await.is_none()); assert!(resolve_iam_ready_with(None).is_none()); assert!(resolve_iam_handle_with(None).is_none()); - assert!(resolve_oidc_handle_with(None).is_none()); + assert!(resolve_federated_identity_service_with(None).is_none()); assert!(resolve_token_signing_key_with(None).is_none()); assert!(resolve_notify_interface_for_context(None).is_none()); - assert!(!publish_oidc_handle_with(None, context_oidc)); + assert!(!publish_federated_identity_service_with(None, context_oidc)); assert!(resolve_bucket_metadata_handle_with(None).is_none()); assert!(resolve_bucket_monitor_handle_with(None).is_none()); assert!(resolve_replication_pool_handle_with(None).is_none()); diff --git a/rustfs/src/app/context/global.rs b/rustfs/src/app/context/global.rs index 9275a00f8..309de0b77 100644 --- a/rustfs/src/app/context/global.rs +++ b/rustfs/src/app/context/global.rs @@ -17,24 +17,25 @@ use super::super::storage_api::context::runtime::set_object_store_resolver; use super::handles::{ IamHandle, KmsHandle, default_action_credential_interface, default_boot_time_interface, default_bucket_metadata_interface, default_bucket_monitor_interface, default_buffer_config_interface, default_deployment_id_interface, - default_endpoints_interface, default_expiry_state_interface, default_internode_metrics_interface, - default_kms_runtime_interface, default_local_node_name_interface, default_lock_client_interface, - default_lock_clients_interface, default_notification_system_interface, default_notify_interface, default_oidc_interface, - default_outbound_tls_runtime_interface, default_performance_metrics_interface, default_region_interface, - default_replication_pool_interface, default_replication_stats_interface, default_runtime_port_interface, - default_s3select_db_interface, default_scanner_metrics_interface, default_server_config_interface, - default_storage_class_interface, default_tier_config_interface, default_transition_state_interface, oidc_interface, + default_endpoints_interface, default_expiry_state_interface, default_federated_identity_interface, + default_internode_metrics_interface, default_kms_runtime_interface, default_local_node_name_interface, + default_lock_client_interface, default_lock_clients_interface, default_notification_system_interface, + default_notify_interface, default_outbound_tls_runtime_interface, default_performance_metrics_interface, + default_region_interface, default_replication_pool_interface, default_replication_stats_interface, + default_runtime_port_interface, default_s3select_db_interface, default_scanner_metrics_interface, + default_server_config_interface, default_storage_class_interface, default_tier_config_interface, + default_transition_state_interface, }; use super::interfaces::{ ActionCredentialInterface, BootTimeInterface, BucketMetadataInterface, BucketMonitorInterface, BufferConfigInterface, - DeploymentIdInterface, EndpointsInterface, ExpiryStateInterface, IamInterface, InternodeMetricsInterface, KmsInterface, - KmsRuntimeInterface, LocalNodeNameInterface, LockClientInterface, LockClientsInterface, NotificationSystemInterface, - NotifyInterface, OidcInterface, OutboundTlsRuntimeInterface, PerformanceMetricsInterface, RegionInterface, - ReplicationPoolInterface, ReplicationStatsInterface, RuntimePortInterface, S3SelectDbInterface, ScannerMetricsInterface, - ServerConfigInterface, StorageClassInterface, TierConfigInterface, TransitionStateInterface, + DeploymentIdInterface, EndpointsInterface, ExpiryStateInterface, FederatedIdentityInterface, IamInterface, + InternodeMetricsInterface, KmsInterface, KmsRuntimeInterface, LocalNodeNameInterface, LockClientInterface, + LockClientsInterface, NotificationSystemInterface, NotifyInterface, OutboundTlsRuntimeInterface, PerformanceMetricsInterface, + RegionInterface, ReplicationPoolInterface, ReplicationStatsInterface, RuntimePortInterface, S3SelectDbInterface, + ScannerMetricsInterface, ServerConfigInterface, StorageClassInterface, TierConfigInterface, TransitionStateInterface, }; use crate::app::object_data_cache::ObjectDataCacheAdapter; -use rustfs_iam::{oidc::OidcSys, store::object::ObjectStore, sys::IamSys}; +use rustfs_iam::{federation::FederatedIdentityService, store::object::ObjectStore, sys::IamSys}; use rustfs_kms::KmsServiceManager; use std::sync::{Arc, OnceLock}; @@ -43,7 +44,7 @@ use std::sync::{Arc, OnceLock}; pub struct AppContext { object_store: Arc, iam: Arc, - oidc: Arc, + federated_identity: Arc, #[allow(dead_code)] kms: Arc, kms_runtime: Arc, @@ -91,7 +92,7 @@ impl AppContext { Self { object_store, iam, - oidc: default_oidc_interface(), + federated_identity: default_federated_identity_interface(), kms, kms_runtime: default_kms_runtime_interface(), outbound_tls_runtime: default_outbound_tls_runtime_interface(), @@ -129,9 +130,7 @@ impl AppContext { iam: Arc>, kms: Arc, ) -> Self { - let mut context = Self::new(object_store, Arc::new(IamHandle::new(iam)), Arc::new(KmsHandle::new(kms))); - context.oidc = oidc_interface(super::runtime_sources::oidc_handle()); - context + Self::new(object_store, Arc::new(IamHandle::new(iam)), Arc::new(KmsHandle::new(kms))) } pub fn object_store(&self) -> Arc { @@ -142,12 +141,12 @@ impl AppContext { self.iam.clone() } - pub fn oidc(&self) -> Arc { - self.oidc.clone() + pub fn federated_identity(&self) -> Arc { + self.federated_identity.clone() } - pub fn publish_oidc_handle(&self, oidc: Arc) -> bool { - self.oidc.publish_handle(oidc) + pub fn publish_federated_identity_service(&self, service: Arc) -> bool { + self.federated_identity.publish_handle(service) } #[allow(dead_code)] @@ -276,7 +275,7 @@ impl AppContext { #[cfg(test)] pub(super) struct AppContextTestInterfaces { pub(super) iam: Arc, - pub(super) oidc: Arc, + pub(super) federated_identity: Arc, pub(super) kms: Arc, pub(super) kms_runtime: Arc, pub(super) outbound_tls_runtime: Arc, @@ -313,7 +312,7 @@ impl AppContext { Self { object_store, iam: interfaces.iam, - oidc: interfaces.oidc, + federated_identity: interfaces.federated_identity, kms: interfaces.kms, kms_runtime: interfaces.kms_runtime, outbound_tls_runtime: interfaces.outbound_tls_runtime, diff --git a/rustfs/src/app/context/handles.rs b/rustfs/src/app/context/handles.rs index b2b864eec..7eb4e3f29 100644 --- a/rustfs/src/app/context/handles.rs +++ b/rustfs/src/app/context/handles.rs @@ -20,18 +20,18 @@ use super::super::storage_api::context::runtime::{ }; use super::interfaces::{ ActionCredentialInterface, BootTimeInterface, BucketMetadataInterface, BucketMonitorInterface, BufferConfigInterface, - DeploymentIdInterface, EndpointsInterface, ExpiryStateInterface, IamInterface, InternodeMetricsInterface, KmsInterface, - KmsRuntimeInterface, LocalNodeNameInterface, LockClientInterface, LockClientsInterface, NotificationSystemInterface, - NotifyInterface, OidcInterface, OutboundTlsRuntimeInterface, PerformanceMetricsInterface, RegionInterface, - ReplicationPoolInterface, ReplicationStatsInterface, RuntimePortInterface, S3SelectDbInterface, ScannerMetricsInterface, - ServerConfigInterface, StorageClassInterface, TierConfigInterface, TransitionStateInterface, + DeploymentIdInterface, EndpointsInterface, ExpiryStateInterface, FederatedIdentityInterface, IamInterface, + InternodeMetricsInterface, KmsInterface, KmsRuntimeInterface, LocalNodeNameInterface, LockClientInterface, + LockClientsInterface, NotificationSystemInterface, NotifyInterface, OutboundTlsRuntimeInterface, PerformanceMetricsInterface, + RegionInterface, ReplicationPoolInterface, ReplicationStatsInterface, RuntimePortInterface, S3SelectDbInterface, + ScannerMetricsInterface, ServerConfigInterface, StorageClassInterface, TierConfigInterface, TransitionStateInterface, }; use super::runtime_sources; use crate::config::RustFSBufferConfig; use async_trait::async_trait; use rustfs_config::server_config::Config; use rustfs_credentials::Credentials; -use rustfs_iam::{oidc::OidcSys, store::object::ObjectStore, sys::IamSys}; +use rustfs_iam::{federation::FederatedIdentityService, store::object::ObjectStore, sys::IamSys}; use rustfs_io_metrics::{PerformanceMetrics, internode_metrics::InternodeMetrics}; use rustfs_kms::KmsServiceManager; use rustfs_lock::LockClient; @@ -42,7 +42,7 @@ use rustfs_tls_runtime::{GlobalPublishedOutboundTlsState, TlsGeneration}; use s3s::dto::SelectObjectContentInput; use std::{ collections::HashMap, - sync::{Arc, RwLock as StdRwLock}, + sync::{Arc, OnceLock, RwLock as StdRwLock}, time::SystemTime, }; use tokio::sync::RwLock; @@ -76,35 +76,35 @@ impl IamInterface for IamHandle { } } -/// Default OIDC interface adapter. -pub struct OidcHandle { - oidc: StdRwLock>>, +/// Default federated identity service interface adapter. +pub struct FederatedIdentityHandle { + service: StdRwLock>>, } -impl OidcHandle { - pub fn new(oidc: Option>) -> Self { +impl FederatedIdentityHandle { + pub fn new(service: Option>) -> Self { Self { - oidc: StdRwLock::new(oidc), + service: StdRwLock::new(service), } } } -impl Default for OidcHandle { +impl Default for FederatedIdentityHandle { fn default() -> Self { Self::new(None) } } -impl OidcInterface for OidcHandle { - fn handle(&self) -> Option> { - self.oidc.read().ok().and_then(|oidc| oidc.as_ref().cloned()) +impl FederatedIdentityInterface for FederatedIdentityHandle { + fn handle(&self) -> Option> { + self.service.read().ok().and_then(|service| service.as_ref().cloned()) } - fn publish_handle(&self, oidc: Arc) -> bool { - let Ok(mut published_oidc) = self.oidc.write() else { + fn publish_handle(&self, service: Arc) -> bool { + let Ok(mut published_service) = self.service.write() else { return false; }; - *published_oidc = Some(oidc); + *published_service = Some(service); true } } @@ -556,12 +556,24 @@ pub fn default_action_credential_interface() -> Arc Arc { - Arc::new(OidcHandle::default()) +static DEFAULT_FEDERATED_IDENTITY_HANDLE: OnceLock> = OnceLock::new(); + +fn default_federated_identity_handle() -> Arc { + DEFAULT_FEDERATED_IDENTITY_HANDLE + .get_or_init(|| Arc::new(FederatedIdentityHandle::default())) + .clone() } -pub fn oidc_interface(oidc: Option>) -> Arc { - Arc::new(OidcHandle::new(oidc)) +pub fn default_federated_identity_interface() -> Arc { + default_federated_identity_handle() +} + +pub fn publish_default_federated_identity_service(service: Arc) -> bool { + default_federated_identity_handle().publish_handle(service) +} + +pub fn federated_identity_interface(service: Option>) -> Arc { + Arc::new(FederatedIdentityHandle::new(service)) } pub fn default_region_interface() -> Arc { @@ -594,10 +606,69 @@ pub fn default_buffer_config_interface() -> Arc { #[cfg(test)] mod tests { - use super::{ServerConfigHandle, runtime_sources}; + use super::{ + ServerConfigHandle, default_federated_identity_interface, federated_identity_interface, + publish_default_federated_identity_service, runtime_sources, + }; use crate::app::context::interfaces::ServerConfigInterface; use rustfs_config::server_config::Config; + use rustfs_iam::{ + federation::{FederatedIdentityRegistry, FederatedIdentityService, oidc::StandardOidcAdapter}, + oidc::OidcSys, + }; use std::collections::HashMap; + use std::sync::Arc; + + fn test_federated_identity_service() -> Arc { + let oidc = OidcSys::empty().expect("empty OIDC configuration should be valid"); + let adapter = Arc::new(StandardOidcAdapter::new(Arc::new(oidc))); + Arc::new(FederatedIdentityService::new(FederatedIdentityRegistry::new(adapter))) + } + + #[test] + fn federated_identity_handle_preserves_early_publish_and_allows_replacement() { + let first = test_federated_identity_service(); + let second = test_federated_identity_service(); + let interface = federated_identity_interface(Some(first.clone())); + + assert!( + interface + .handle() + .as_ref() + .is_some_and(|resolved| Arc::ptr_eq(resolved, &first)) + ); + + assert!(interface.publish_handle(second.clone())); + assert!( + interface + .handle() + .as_ref() + .is_some_and(|resolved| Arc::ptr_eq(resolved, &second)) + ); + } + + #[test] + fn default_federated_identity_handle_shares_early_publish_and_replacement() { + let first = test_federated_identity_service(); + let second = test_federated_identity_service(); + + assert!(publish_default_federated_identity_service(first.clone())); + let interface = default_federated_identity_interface(); + assert!( + interface + .handle() + .as_ref() + .is_some_and(|resolved| Arc::ptr_eq(resolved, &first)) + ); + + assert!(publish_default_federated_identity_service(second.clone())); + assert!( + interface + .handle() + .as_ref() + .is_some_and(|resolved| Arc::ptr_eq(resolved, &second)) + ); + } // backlog#1052 S3: each context's server-config handle keeps its own copy, // so two handles that both published stay isolated even though the shared diff --git a/rustfs/src/app/context/interfaces.rs b/rustfs/src/app/context/interfaces.rs index 728dd1b11..17e029b76 100644 --- a/rustfs/src/app/context/interfaces.rs +++ b/rustfs/src/app/context/interfaces.rs @@ -22,7 +22,7 @@ use crate::config::RustFSBufferConfig; use async_trait::async_trait; use rustfs_config::server_config::Config; use rustfs_credentials::Credentials; -use rustfs_iam::{oidc::OidcSys, store::object::ObjectStore, sys::IamSys}; +use rustfs_iam::{federation::FederatedIdentityService, store::object::ObjectStore, sys::IamSys}; use rustfs_io_metrics::{PerformanceMetrics, internode_metrics::InternodeMetrics}; use rustfs_kms::KmsServiceManager; use rustfs_lock::LockClient; @@ -44,10 +44,10 @@ pub trait IamInterface: Send + Sync { } } -/// OIDC interface for admin and runtime consumers. -pub trait OidcInterface: Send + Sync { - fn handle(&self) -> Option>; - fn publish_handle(&self, _oidc: Arc) -> bool { +/// Federated identity service interface for admin and runtime consumers. +pub trait FederatedIdentityInterface: Send + Sync { + fn handle(&self) -> Option>; + fn publish_handle(&self, _service: Arc) -> bool { false } } diff --git a/rustfs/src/app/context/runtime_sources.rs b/rustfs/src/app/context/runtime_sources.rs index 8d7a8a7b0..f93b83a69 100644 --- a/rustfs/src/app/context/runtime_sources.rs +++ b/rustfs/src/app/context/runtime_sources.rs @@ -25,7 +25,6 @@ use super::super::storage_api::context::runtime::{ use crate::config::{RustFSBufferConfig, get_global_buffer_config}; use rustfs_config::server_config::{Config, get_global_server_config, set_global_server_config}; use rustfs_credentials::{Credentials, get_global_action_cred}; -use rustfs_iam::oidc::OidcSys; use rustfs_io_metrics::{ PerformanceMetrics, global_metrics::get_global_metrics, @@ -64,10 +63,6 @@ pub async fn outbound_tls_state() -> GlobalPublishedOutboundTlsState { load_global_outbound_tls_state().await } -pub fn oidc_handle() -> Option> { - rustfs_iam::get_oidc() -} - pub fn token_signing_key() -> Option { rustfs_iam::manager::get_token_signing_key() } diff --git a/rustfs/src/runtime_sources.rs b/rustfs/src/runtime_sources.rs index 9b5491cf6..6e3a6cd8f 100644 --- a/rustfs/src/runtime_sources.rs +++ b/rustfs/src/runtime_sources.rs @@ -22,13 +22,14 @@ pub(crate) use context::{ default_s3select_db_interface as fallback_s3select_db_interface, default_scanner_metrics_interface as fallback_scanner_metrics_interface, default_server_config_interface as fallback_server_config_interface, - default_storage_class_interface as fallback_storage_class_interface, publish_oidc_handle, publish_server_config, - publish_storage_class_config, resolve_action_credentials as current_action_credentials, + default_storage_class_interface as fallback_storage_class_interface, publish_federated_identity_service, + publish_server_config, publish_storage_class_config, resolve_action_credentials as current_action_credentials, resolve_boot_time as current_boot_time, resolve_bucket_metadata_handle as current_bucket_metadata_handle, resolve_bucket_monitor_handle as current_bucket_monitor_handle, resolve_buffer_config as current_buffer_config, resolve_daily_tier_stats as current_daily_tier_stats, resolve_deployment_id as current_deployment_id, resolve_encryption_service as current_encryption_service, resolve_endpoints_handle as current_endpoints_handle, - resolve_expiry_state_handle as current_expiry_state_handle, resolve_iam_handle as current_iam_handle, + resolve_expiry_state_handle as current_expiry_state_handle, + resolve_federated_identity_service as current_federated_identity_service, resolve_iam_handle as current_iam_handle, resolve_iam_ready as current_iam_ready, resolve_internode_metrics as current_internode_metrics, resolve_kms_runtime_service_manager as current_kms_runtime_service_manager, resolve_local_node_name as current_local_node_name, resolve_lock_client as current_lock_client, @@ -39,7 +40,6 @@ pub(crate) use context::{ resolve_object_data_cache_handle_for_context as current_object_data_cache_handle_for_context, resolve_object_store_handle as current_object_store_handle, resolve_object_store_handle_for_context as current_object_store_handle_for_context, - resolve_oidc_handle as current_oidc_handle, resolve_or_init_kms_runtime_service_manager as current_or_init_kms_runtime_service_manager, resolve_outbound_tls_generation as current_outbound_tls_generation, resolve_outbound_tls_state as current_outbound_tls_state, resolve_performance_metrics as current_performance_metrics, resolve_ready_iam_handle as current_ready_iam_handle, diff --git a/rustfs/src/startup_auth.rs b/rustfs/src/startup_auth.rs index e8a8202d1..02029e538 100644 --- a/rustfs/src/startup_auth.rs +++ b/rustfs/src/startup_auth.rs @@ -12,8 +12,14 @@ // See the License for the specific language governing permissions and // limitations under the License. -use rustfs_iam::{get_oidc, init_oidc_sys}; -use std::io::{Error, Result}; +use rustfs_iam::{ + federation::{FederatedIdentityRegistry, FederatedIdentityService, oidc::StandardOidcAdapter}, + get_oidc, init_oidc_sys, +}; +use std::{ + io::{Error, Result}, + sync::Arc, +}; use tracing::{error, info, warn}; const LOG_COMPONENT_MAIN: &str = "main"; @@ -47,7 +53,10 @@ pub(crate) async fn init_auth_integrations() -> Result<()> { match init_oidc_sys().await { Ok(()) => { if let Some(oidc) = get_oidc() { - crate::runtime_sources::publish_oidc_handle(oidc); + let adapter = Arc::new(StandardOidcAdapter::new(oidc)); + let registry = FederatedIdentityRegistry::new(adapter); + let service = Arc::new(FederatedIdentityService::new(registry)); + crate::runtime_sources::publish_federated_identity_service(service); } } Err(e) => {