fix(auth): isolate embedded IAM contexts (#5704)

This commit is contained in:
Zhengchao An
2026-08-04 23:12:35 +08:00
committed by GitHub
parent c26419e357
commit 1695873e55
6 changed files with 161 additions and 22 deletions
+9
View File
@@ -168,6 +168,15 @@ pub async fn build_iam_sys(ecstore: Arc<IamStore>) -> Result<Arc<IamSys<ObjectSt
Ok(Arc::new(IamSys::new(cache_manager)))
}
/// Build an IAM system for an application context and publish the first one
/// as the ambient compatibility default.
#[instrument(skip(ecstore))]
pub async fn init_iam_sys_for_context(ecstore: Arc<IamStore>) -> Result<Arc<IamSys<ObjectStore>>> {
let iam_instance = build_iam_sys(ecstore).await?;
let _ = IAM_SYS.set(iam_instance.clone());
Ok(iam_instance)
}
#[instrument(skip(ecstore))]
pub async fn init_iam_sys(ecstore: Arc<IamStore>) -> Result<Arc<IamSys<ObjectStore>>> {
if let Some(existing) = IAM_SYS.get() {
+23 -2
View File
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::runtime_sources::{AppContext, current_action_credentials, current_ready_iam_handle};
use crate::runtime_sources::{AppContext, ServerContextSlot, current_action_credentials, current_ready_iam_handle};
use http::HeaderMap;
use http::Uri;
use rustfs_credentials::Credentials;
@@ -113,6 +113,7 @@ pub struct IAMAuth {
simple_auth: SimpleAuth,
access_key: String,
secret_key: SecretKey,
server_ctx: Option<std::sync::Arc<ServerContextSlot>>,
}
impl Clone for IAMAuth {
@@ -121,6 +122,7 @@ impl Clone for IAMAuth {
simple_auth: SimpleAuth::from_single(self.access_key.clone(), self.secret_key.clone()),
access_key: self.access_key.clone(),
secret_key: self.secret_key.clone(),
server_ctx: self.server_ctx.clone(),
}
}
}
@@ -134,8 +136,19 @@ impl IAMAuth {
simple_auth,
access_key,
secret_key,
server_ctx: None,
}
}
pub(crate) fn with_server_context(
ak: impl Into<String>,
sk: impl Into<SecretKey>,
server_ctx: std::sync::Arc<ServerContextSlot>,
) -> Self {
let mut auth = Self::new(ak, sk);
auth.server_ctx = Some(server_ctx);
auth
}
}
#[async_trait::async_trait]
@@ -186,7 +199,15 @@ impl S3Auth for IAMAuth {
return Ok(key);
}
if let Ok(iam_store) = current_ready_iam_handle() {
let iam_store = match &self.server_ctx {
Some(server_ctx) => server_ctx
.installed_app_context()
.filter(|context| context.iam().is_ready())
.map(|context| context.iam().handle())
.ok_or(()),
None => current_ready_iam_handle().map_err(|_| ()),
};
if let Ok(iam_store) = iam_store {
// Use check_key instead of get_user to ensure user is loaded from disk if not in cache
// This is important for newly created users that may not be in cache yet.
// check_key will automatically attempt to load the user from disk if not found in cache.
+1 -1
View File
@@ -732,7 +732,7 @@ pub async fn start_http_server(
.transpose()
.map_err(Error::other)?;
b.set_auth(IAMAuth::new(access_key, secret_key));
b.set_auth(IAMAuth::with_server_context(access_key, secret_key, server_ctx.clone()));
b.set_access(store);
b.set_route(storage::metadata_route::with_metadata_route(
admin::make_admin_route(config.console_enable, admin_server_ctx)?,
+18 -7
View File
@@ -16,7 +16,7 @@ use crate::runtime_sources::{AppContext, ServerContextSlot};
use crate::server::{ServiceStateManager, publish_ready_when_runtime_ready};
use crate::storage_api::startup::iam::ECStore;
use rustfs_common::{GlobalReadiness, SystemStage};
use rustfs_iam::init_iam_sys;
use rustfs_iam::init_iam_sys_for_context;
use rustfs_kms::KmsServiceManager;
use std::future::Future;
use std::io::Result;
@@ -86,15 +86,13 @@ where
}
async fn finalize_iam_recovery(
iam: Arc<rustfs_iam::sys::IamSys<rustfs_iam::store::object::ObjectStore>>,
store: Arc<ECStore>,
kms_interface: Arc<KmsServiceManager>,
readiness: Arc<GlobalReadiness>,
state_manager: Option<Arc<ServiceStateManager>>,
server_ctx: Arc<ServerContextSlot>,
) -> Result<()> {
// Resolve the globally published system directly (not context-first —
// the AppContext does not exist yet; creating it is this function's job).
let iam = rustfs_iam::get().map_err(|_| std::io::Error::other("IAM recovered but unavailable"))?;
let endpoint_pools = heal_control_endpoint_pools(store.as_ref())?;
AppContext::ensure_startup_after_iam(store, kms_interface, &server_ctx, iam)?;
let topology_ready =
@@ -147,6 +145,8 @@ fn spawn_iam_recovery_task(
let finalize_readiness = readiness;
let finalize_state_manager = state_manager;
let finalize_server_ctx = server_ctx;
let recovered_iam = Arc::new(std::sync::Mutex::new(None));
let init_recovered_iam = recovered_iam.clone();
tokio::spawn(async move {
run_iam_recovery_loop(
initial_interval,
@@ -154,7 +154,14 @@ fn spawn_iam_recovery_task(
shutdown_token,
move || {
let store = init_store.clone();
Box::pin(async move { attempt_init_iam_sys(store).await.map(|_| ()) })
let recovered_iam = init_recovered_iam.clone();
Box::pin(async move {
let iam = attempt_init_iam_sys(store).await?;
*recovered_iam
.lock()
.map_err(|_| std::io::Error::other("IAM recovery state poisoned"))? = Some(iam);
Ok(())
})
},
move || {
let store = finalize_store.clone();
@@ -162,7 +169,11 @@ fn spawn_iam_recovery_task(
let readiness = finalize_readiness.clone();
let state_manager = finalize_state_manager.clone();
let server_ctx = finalize_server_ctx.clone();
Box::pin(async move { finalize_iam_recovery(store, kms_interface, readiness, state_manager, server_ctx).await })
let iam = recovered_iam.lock().ok().and_then(|mut recovered| recovered.take());
Box::pin(async move {
let iam = iam.ok_or_else(|| std::io::Error::other("IAM recovered but unavailable"))?;
finalize_iam_recovery(iam, store, kms_interface, readiness, state_manager, server_ctx).await
})
},
)
.await;
@@ -359,7 +370,7 @@ async fn attempt_init_iam_sys(
return Err(std::io::Error::other("forced test IAM bootstrap failure"));
}
init_iam_sys(store).await.map_err(std::io::Error::other)
init_iam_sys_for_context(store).await.map_err(std::io::Error::other)
}
/// Attempt IAM bootstrap at startup. If it fails, enter degraded mode and
+9 -1
View File
@@ -677,7 +677,15 @@ pub async fn authorize_request<T>(req: &mut S3Request<T>, action: Action) -> S3R
let version_id = req_info.version_id.clone();
if let Some(cred) = &cred {
let Ok(iam_store) = runtime_sources::current_ready_iam_handle() else {
let iam_store = match req.extensions.get::<std::sync::Arc<ServerContextSlot>>() {
Some(server_ctx) => server_ctx
.installed_app_context()
.filter(|context| context.iam().is_ready())
.map(|context| context.iam().handle())
.ok_or(()),
None => runtime_sources::current_ready_iam_handle().map_err(|_| ()),
};
let Ok(iam_store) = iam_store else {
return Err(S3Error::with_message(
S3ErrorCode::InternalError,
format!("authorize_request {:?}", IamError::IamSysNotInitialized),
+101 -11
View File
@@ -71,22 +71,27 @@ fn hmac(key: &[u8], value: &str) -> Vec<u8> {
fn signed_admin_request(
client: &reqwest::Client,
endpoint: &str,
method: reqwest::Method,
request_path: &str,
access_key: &str,
secret_key: &str,
body: &[u8],
) -> reqwest::RequestBuilder {
let host = endpoint
.strip_prefix("http://")
.or_else(|| endpoint.strip_prefix("https://"))
.expect("embedded endpoint scheme");
let payload_hash = sha256_hex(b"");
let payload_hash = sha256_hex(body);
let now = Utc::now();
let amz_date = now.format("%Y%m%dT%H%M%SZ").to_string();
let date = now.format("%Y%m%d").to_string();
let canonical_headers = format!("host:{host}\nx-amz-content-sha256:{payload_hash}\nx-amz-date:{amz_date}\n");
let signed_headers = "host;x-amz-content-sha256;x-amz-date";
let (path, query) = request_path.split_once('?').unwrap_or((request_path, ""));
let canonical_request = format!("GET\n{path}\n{query}\n{canonical_headers}\n{signed_headers}\n{payload_hash}");
let canonical_request = format!(
"{}\n{path}\n{query}\n{canonical_headers}\n{signed_headers}\n{payload_hash}",
method.as_str()
);
let scope = format!("{date}/us-east-1/s3/aws4_request");
let string_to_sign = format!("AWS4-HMAC-SHA256\n{amz_date}\n{scope}\n{}", sha256_hex(canonical_request.as_bytes()));
let date_key = hmac(format!("AWS4{secret_key}").as_bytes(), &date);
@@ -99,11 +104,12 @@ fn signed_admin_request(
);
client
.get(format!("{endpoint}{request_path}"))
.request(method, format!("{endpoint}{request_path}"))
.header("host", host)
.header("x-amz-content-sha256", payload_hash)
.header("x-amz-date", amz_date)
.header("authorization", authorization)
.body(body.to_vec())
}
// backlog#1052 acceptance: a second embedded server in the same process no
@@ -328,6 +334,81 @@ async fn two_embedded_servers_isolate_auth_and_data_planes_body() {
server_b.shutdown().await;
}
#[cfg(feature = "e2e-test-hooks")]
#[tokio::test]
async fn embedded_servers_isolate_iam_users_and_policies() {
let port_a = find_available_port().expect("find free port for server A");
let server_a = RustFSServerBuilder::new()
.address(format!("127.0.0.1:{port_a}"))
.access_key("iam-root-a")
.secret_key("iam-root-secret-a")
.build()
.await
.expect("start embedded server A");
let port_b = find_available_port().expect("find free port for server B");
let server_b = RustFSServerBuilder::new()
.address(format!("127.0.0.1:{port_b}"))
.access_key("iam-root-b")
.secret_key("iam-root-secret-b")
.build()
.await
.expect("start embedded server B");
let http = reqwest::Client::builder()
.no_proxy()
.build()
.expect("build local admin client");
let user_body = serde_json::json!({"secretKey": "iam-user-secret-a", "status": "enabled"}).to_string();
for (path, body) in [
("/rustfs/admin/v3/add-user?accessKey=iam-user-a", user_body.as_bytes()),
(
"/rustfs/admin/v3/set-user-or-group-policy?policyName=readwrite&userOrGroup=iam-user-a&isGroup=false",
b"".as_slice(),
),
] {
let response = signed_admin_request(
&http,
&server_a.endpoint(),
reqwest::Method::PUT,
path,
server_a.access_key(),
server_a.secret_key(),
body,
)
.send()
.await
.expect("send server A IAM request");
assert!(response.status().is_success(), "server A IAM setup failed: {}", response.status());
}
let client_b = s3_client(&server_b.endpoint(), server_b.access_key(), server_b.secret_key());
client_b
.create_bucket()
.bucket("private-to-b")
.send()
.await
.expect("create server B bucket");
client_b
.put_object()
.bucket("private-to-b")
.key("secret.txt")
.body(ByteStream::from_static(b"server B data"))
.send()
.await
.expect("write server B object");
let cross_instance = s3_client(&server_b.endpoint(), "iam-user-a", "iam-user-secret-a")
.get_object()
.bucket("private-to-b")
.key("secret.txt")
.send()
.await;
assert!(cross_instance.is_err(), "server A IAM user must not authorize against server B");
server_b.shutdown().await;
server_a.shutdown().await;
}
#[cfg(feature = "e2e-test-hooks")]
#[tokio::test]
async fn second_embedded_server_fails_closed_until_its_context_slot_is_installed() {
@@ -392,10 +473,11 @@ async fn second_embedded_server_fails_closed_until_its_context_slot_is_installed
.build()
.expect("build local admin client without proxy");
let inspect_path = "/rustfs/admin/v3/inspect-data?file=marker.txt&volume=startup-window";
let before_install = signed_admin_request(&http, &endpoint_b, inspect_path, b_access_key, b_secret_key)
.send()
.await
.expect("server B HTTP listener must accept the paused request");
let before_install =
signed_admin_request(&http, &endpoint_b, reqwest::Method::GET, inspect_path, b_access_key, b_secret_key, b"")
.send()
.await
.expect("server B HTTP listener must accept the paused request");
let before_install_status = before_install.status();
let before_install_body = before_install.text().await.expect("read paused response body");
assert_eq!(before_install_status, StatusCode::SERVICE_UNAVAILABLE, "{before_install_body}");
@@ -426,10 +508,18 @@ async fn second_embedded_server_fails_closed_until_its_context_slot_is_installed
.await
.expect("server B writes its marker");
let after_install = signed_admin_request(&http, &server_b.endpoint(), inspect_path, b_access_key, b_secret_key)
.send()
.await
.expect("server B admin request after context installation");
let after_install = signed_admin_request(
&http,
&server_b.endpoint(),
reqwest::Method::GET,
inspect_path,
b_access_key,
b_secret_key,
b"",
)
.send()
.await
.expect("server B admin request after context installation");
assert_eq!(after_install.status(), StatusCode::OK);
assert_eq!(after_install.bytes().await.expect("read server B marker"), b"from B".as_slice());