mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 08:18:18 +00:00
fix: improve IAM and quota authorization (#1781)
Signed-off-by: yxrxy <yxrxytrigger@gmail.com> Co-authored-by: loverustfs <hello@rustfs.com>
This commit is contained in:
@@ -133,8 +133,7 @@ where
|
||||
if let Err(e) = self.clone().load().await {
|
||||
if attempt == MAX_RETRIES - 1 {
|
||||
self.state.store(IamState::Error as u8, Ordering::SeqCst);
|
||||
error!("IAM fail to load initial data after {} attempts: {:?}", MAX_RETRIES, e);
|
||||
return Err(e);
|
||||
warn!("IAM failed to load initial data after {} attempts: {:?}", MAX_RETRIES, e);
|
||||
} else {
|
||||
warn!("IAM load failed, retrying... attempt {}", attempt + 1);
|
||||
tokio::time::sleep(Duration::from_secs(1)).await;
|
||||
@@ -162,7 +161,7 @@ where
|
||||
_ = ticker.tick() => {
|
||||
info!("iam load ticker");
|
||||
if let Err(err) =s.clone().load().await{
|
||||
error!("iam load err {:?}", err);
|
||||
warn!("iam load err {:?}", err);
|
||||
}
|
||||
},
|
||||
i = receiver.recv() => {
|
||||
@@ -173,7 +172,7 @@ where
|
||||
if last <= t {
|
||||
info!("iam load receiver load");
|
||||
if let Err(err) =s.clone().load().await{
|
||||
error!("iam load err {:?}", err);
|
||||
warn!("iam load err {:?}", err);
|
||||
}
|
||||
ticker.reset();
|
||||
}
|
||||
|
||||
+61
-14
@@ -24,6 +24,15 @@ use s3s::s3_error;
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
|
||||
#[derive(Clone)]
|
||||
struct AuthContext<'a> {
|
||||
headers: &'a HeaderMap,
|
||||
cred: &'a Credentials,
|
||||
is_owner: bool,
|
||||
deny_only: bool,
|
||||
remote_addr: Option<std::net::SocketAddr>,
|
||||
}
|
||||
|
||||
pub async fn validate_admin_request(
|
||||
headers: &HeaderMap,
|
||||
cred: &Credentials,
|
||||
@@ -35,8 +44,16 @@ pub async fn validate_admin_request(
|
||||
let Ok(iam_store) = rustfs_iam::get() else {
|
||||
return Err(s3_error!(InternalError, "iam not init"));
|
||||
};
|
||||
let ctx = AuthContext {
|
||||
headers,
|
||||
cred,
|
||||
is_owner,
|
||||
deny_only,
|
||||
remote_addr,
|
||||
};
|
||||
|
||||
for action in actions {
|
||||
match check_admin_request_auth(iam_store.clone(), headers, cred, is_owner, deny_only, action, remote_addr).await {
|
||||
match check_admin_request_auth(iam_store.clone(), &ctx, action, "", "").await {
|
||||
Ok(_) => return Ok(()),
|
||||
Err(_) => {
|
||||
continue;
|
||||
@@ -49,26 +66,24 @@ pub async fn validate_admin_request(
|
||||
|
||||
async fn check_admin_request_auth(
|
||||
iam_store: Arc<IamSys<ObjectStore>>,
|
||||
headers: &HeaderMap,
|
||||
cred: &Credentials,
|
||||
is_owner: bool,
|
||||
deny_only: bool,
|
||||
ctx: &AuthContext<'_>,
|
||||
action: Action,
|
||||
remote_addr: Option<std::net::SocketAddr>,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
) -> S3Result<()> {
|
||||
let conditions = get_condition_values(headers, cred, None, None, remote_addr);
|
||||
let conditions = get_condition_values(ctx.headers, ctx.cred, None, None, ctx.remote_addr);
|
||||
|
||||
if !iam_store
|
||||
.is_allowed(&Args {
|
||||
account: &cred.access_key,
|
||||
groups: &cred.groups,
|
||||
account: &ctx.cred.access_key,
|
||||
groups: &ctx.cred.groups,
|
||||
action,
|
||||
conditions: &conditions,
|
||||
is_owner,
|
||||
claims: cred.claims.as_ref().unwrap_or(&HashMap::new()),
|
||||
deny_only,
|
||||
bucket: "",
|
||||
object: "",
|
||||
is_owner: ctx.is_owner,
|
||||
claims: ctx.cred.claims.as_ref().unwrap_or(&HashMap::new()),
|
||||
deny_only: ctx.deny_only,
|
||||
bucket,
|
||||
object,
|
||||
})
|
||||
.await
|
||||
{
|
||||
@@ -77,3 +92,35 @@ async fn check_admin_request_auth(
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn validate_admin_request_with_bucket(
|
||||
headers: &HeaderMap,
|
||||
cred: &Credentials,
|
||||
is_owner: bool,
|
||||
deny_only: bool,
|
||||
actions: Vec<Action>,
|
||||
remote_addr: Option<std::net::SocketAddr>,
|
||||
bucket: &str,
|
||||
) -> S3Result<()> {
|
||||
let Ok(iam_store) = rustfs_iam::get() else {
|
||||
return Err(s3_error!(InternalError, "iam not init"));
|
||||
};
|
||||
let ctx = AuthContext {
|
||||
headers,
|
||||
cred,
|
||||
is_owner,
|
||||
deny_only,
|
||||
remote_addr,
|
||||
};
|
||||
|
||||
for action in actions {
|
||||
match check_admin_request_auth(iam_store.clone(), &ctx, action, bucket, "").await {
|
||||
Ok(_) => return Ok(()),
|
||||
Err(_) => {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Err(s3_error!(AccessDenied, "Access Denied"))
|
||||
}
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
|
||||
//! Quota admin handlers for HTTP API
|
||||
|
||||
use crate::admin::auth::validate_admin_request;
|
||||
use crate::admin::auth::{validate_admin_request, validate_admin_request_with_bucket};
|
||||
use crate::admin::router::{AdminOperation, Operation, S3Router};
|
||||
use crate::auth::{check_key_valid, get_session_token};
|
||||
use crate::server::ADMIN_PREFIX;
|
||||
@@ -215,21 +215,22 @@ impl Operation for GetBucketQuotaHandler {
|
||||
let (cred, owner) =
|
||||
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?;
|
||||
|
||||
validate_admin_request(
|
||||
let bucket = params.get("bucket").unwrap_or("").to_string();
|
||||
if bucket.is_empty() {
|
||||
return Err(s3_error!(InvalidRequest, "bucket name is required"));
|
||||
}
|
||||
|
||||
validate_admin_request_with_bucket(
|
||||
&req.headers,
|
||||
&cred,
|
||||
owner,
|
||||
false,
|
||||
vec![Action::S3Action(S3Action::GetBucketQuotaAction)],
|
||||
None,
|
||||
&bucket,
|
||||
)
|
||||
.await?;
|
||||
|
||||
let bucket = params.get("bucket").unwrap_or("").to_string();
|
||||
if bucket.is_empty() {
|
||||
return Err(s3_error!(InvalidRequest, "bucket name is required"));
|
||||
}
|
||||
|
||||
let metadata_sys_lock = rustfs_ecstore::bucket::metadata_sys::GLOBAL_BucketMetadataSys
|
||||
.get()
|
||||
.ok_or_else(|| s3_error!(InternalError, "Bucket metadata system not initialized"))?;
|
||||
@@ -343,21 +344,22 @@ impl Operation for GetBucketQuotaStatsHandler {
|
||||
let (cred, owner) =
|
||||
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?;
|
||||
|
||||
validate_admin_request(
|
||||
let bucket = params.get("bucket").unwrap_or("").to_string();
|
||||
if bucket.is_empty() {
|
||||
return Err(s3_error!(InvalidRequest, "bucket name is required"));
|
||||
}
|
||||
|
||||
validate_admin_request_with_bucket(
|
||||
&req.headers,
|
||||
&cred,
|
||||
owner,
|
||||
false,
|
||||
vec![Action::S3Action(S3Action::GetBucketQuotaAction)],
|
||||
None,
|
||||
&bucket,
|
||||
)
|
||||
.await?;
|
||||
|
||||
let bucket = params.get("bucket").unwrap_or("").to_string();
|
||||
if bucket.is_empty() {
|
||||
return Err(s3_error!(InvalidRequest, "bucket name is required"));
|
||||
}
|
||||
|
||||
let metadata_sys_lock = rustfs_ecstore::bucket::metadata_sys::GLOBAL_BucketMetadataSys
|
||||
.get()
|
||||
.ok_or_else(|| s3_error!(InternalError, "Bucket metadata system not initialized"))?;
|
||||
@@ -410,21 +412,22 @@ impl Operation for CheckBucketQuotaHandler {
|
||||
let (cred, owner) =
|
||||
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &cred.access_key).await?;
|
||||
|
||||
validate_admin_request(
|
||||
let bucket = params.get("bucket").unwrap_or("").to_string();
|
||||
if bucket.is_empty() {
|
||||
return Err(s3_error!(InvalidRequest, "bucket name is required"));
|
||||
}
|
||||
|
||||
validate_admin_request_with_bucket(
|
||||
&req.headers,
|
||||
&cred,
|
||||
owner,
|
||||
false,
|
||||
vec![Action::S3Action(S3Action::GetBucketQuotaAction)],
|
||||
None,
|
||||
&bucket,
|
||||
)
|
||||
.await?;
|
||||
|
||||
let bucket = params.get("bucket").unwrap_or("").to_string();
|
||||
if bucket.is_empty() {
|
||||
return Err(s3_error!(InvalidRequest, "bucket name is required"));
|
||||
}
|
||||
|
||||
let body = req
|
||||
.input
|
||||
.store_all_limited(rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE)
|
||||
|
||||
Reference in New Issue
Block a user