feat(table-catalog): add credential response boundary (#3305)

* feat(table-catalog): add credential response boundary

* fix(table-catalog): enforce scoped catalog resources

* fix(table-catalog): reduce admin auth scope arguments

---------

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
Henry Guo
2026-06-10 09:15:51 +08:00
committed by GitHub
parent a73c90c811
commit 2eafd9c602
3 changed files with 255 additions and 19 deletions
+73
View File
@@ -1587,6 +1587,79 @@ mod test {
Ok(())
}
#[tokio::test]
async fn test_table_admin_action_with_resource_is_limited_to_table_object_scope() -> Result<()> {
use crate::policy::action::{Action, AdminAction};
let data = r#"
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": ["admin:GetTableMetadata"],
"Resource": ["arn:aws:s3:::warehouse-a/namespaces/analytics/tables/events"]
}
]
}
"#;
let policy = Policy::parse_config(data.as_bytes())?;
let conditions = HashMap::new();
let claims = HashMap::new();
let groups = None;
let matching_args = Args {
account: "testuser",
groups: &groups,
action: Action::AdminAction(AdminAction::GetTableMetadataAction),
bucket: "warehouse-a",
conditions: &conditions,
is_owner: false,
object: "namespaces/analytics/tables/events",
claims: &claims,
deny_only: false,
};
assert!(
policy.is_allowed(&matching_args).await,
"table admin action should allow the explicitly granted table resource"
);
let namespace_only_args = Args {
account: "testuser",
groups: &groups,
action: Action::AdminAction(AdminAction::GetTableMetadataAction),
bucket: "warehouse-a",
conditions: &conditions,
is_owner: false,
object: "namespaces/analytics",
claims: &claims,
deny_only: false,
};
assert!(
!policy.is_allowed(&namespace_only_args).await,
"table admin action must not match a namespace-only resource when a table resource is required"
);
let other_table_args = Args {
account: "testuser",
groups: &groups,
action: Action::AdminAction(AdminAction::GetTableMetadataAction),
bucket: "warehouse-a",
conditions: &conditions,
is_owner: false,
object: "namespaces/analytics/tables/orders",
claims: &claims,
deny_only: false,
};
assert!(
!policy.is_allowed(&other_table_args).await,
"table admin action must not match a different table resource"
);
Ok(())
}
#[tokio::test]
async fn test_table_admin_action_with_not_resource_excludes_bucket() -> Result<()> {
use crate::policy::action::{Action, AdminAction};
+38 -1
View File
@@ -32,6 +32,22 @@ struct AuthContext<'a> {
remote_addr: Option<std::net::SocketAddr>,
}
#[derive(Clone, Copy, Debug)]
pub struct AdminResourceScope<'a> {
bucket: &'a str,
object: &'a str,
}
impl<'a> AdminResourceScope<'a> {
pub fn bucket(bucket: &'a str) -> Self {
Self { bucket, object: "" }
}
pub fn bucket_object(bucket: &'a str, object: &'a str) -> Self {
Self { bucket, object }
}
}
pub async fn validate_admin_request(
headers: &HeaderMap,
cred: &Credentials,
@@ -108,6 +124,27 @@ pub async fn validate_admin_request_with_bucket(
actions: Vec<Action>,
remote_addr: Option<std::net::SocketAddr>,
bucket: &str,
) -> S3Result<()> {
validate_admin_request_with_bucket_object(
headers,
cred,
is_owner,
deny_only,
actions,
remote_addr,
AdminResourceScope::bucket(bucket),
)
.await
}
pub async fn validate_admin_request_with_bucket_object(
headers: &HeaderMap,
cred: &Credentials,
is_owner: bool,
deny_only: bool,
actions: Vec<Action>,
remote_addr: Option<std::net::SocketAddr>,
resource: AdminResourceScope<'_>,
) -> S3Result<()> {
let Ok(iam_store) = rustfs_iam::get() else {
return Err(s3_error!(InternalError, "iam not init"));
@@ -121,7 +158,7 @@ pub async fn validate_admin_request_with_bucket(
};
for action in &actions {
if check_admin_request_auth(iam_store.clone(), &ctx, *action, bucket, "")
if check_admin_request_auth(iam_store.clone(), &ctx, *action, resource.bucket, resource.object)
.await
.is_ok()
{
+144 -18
View File
@@ -13,7 +13,7 @@
// limitations under the License.
use crate::admin::{
auth::{validate_admin_request, validate_admin_request_with_bucket},
auth::{AdminResourceScope, validate_admin_request, validate_admin_request_with_bucket_object},
router::{AdminOperation, Operation, S3Router},
};
use crate::auth::{check_key_valid, get_session_token};
@@ -33,6 +33,10 @@ use uuid::Uuid;
const JSON_CONTENT_TYPE: &str = "application/json";
const WAREHOUSE_PROPERTY: &str = "warehouse";
const CREDENTIAL_VENDING_CONFIG_KEY: &str = "rustfs.credential-vending";
const CREDENTIAL_VENDING_UNSUPPORTED: &str = "unsupported";
const TABLE_CATALOG_NAMESPACE_RESOURCE_ROOT: &str = "namespaces";
const TABLE_CATALOG_TABLE_RESOURCE_ROOT: &str = "tables";
const TABLE_CATALOG_ENDPOINTS: &[&str] = &[
"GET /{warehouse}/namespaces",
"POST /{warehouse}/namespaces",
@@ -155,12 +159,20 @@ struct RestListTablesResponse {
identifiers: Vec<RestTableIdentifier>,
}
#[derive(Debug, Serialize)]
struct RestStorageCredential {
prefix: String,
config: BTreeMap<String, String>,
}
#[derive(Debug, Serialize)]
struct RestLoadTableResponse {
#[serde(rename = "metadata-location")]
metadata_location: String,
metadata: serde_json::Value,
config: BTreeMap<String, String>,
#[serde(rename = "storage-credentials")]
storage_credentials: Vec<RestStorageCredential>,
}
#[derive(Debug, Serialize)]
@@ -274,7 +286,54 @@ async fn authorize_table_catalog_request(req: &S3Request<Body>, action: AdminAct
.await
}
async fn authorize_table_catalog_warehouse_request(req: &S3Request<Body>, warehouse: &str, action: AdminAction) -> S3Result<()> {
#[derive(Debug, Clone)]
struct TableCatalogResource<'a> {
warehouse: &'a str,
namespace: Option<String>,
table: Option<String>,
}
impl<'a> TableCatalogResource<'a> {
fn warehouse(warehouse: &'a str) -> Self {
Self {
warehouse,
namespace: None,
table: None,
}
}
fn namespace(warehouse: &'a str, namespace: &crate::table_catalog::Namespace) -> Self {
Self {
warehouse,
namespace: Some(namespace.storage_id()),
table: None,
}
}
fn table(warehouse: &'a str, namespace: &crate::table_catalog::Namespace, table: &str) -> Self {
Self {
warehouse,
namespace: Some(namespace.storage_id()),
table: Some(table.to_string()),
}
}
fn object_path(&self) -> Option<String> {
match (&self.namespace, &self.table) {
(Some(namespace), Some(table)) => Some(format!(
"{TABLE_CATALOG_NAMESPACE_RESOURCE_ROOT}/{namespace}/{TABLE_CATALOG_TABLE_RESOURCE_ROOT}/{table}"
)),
(Some(namespace), None) => Some(format!("{TABLE_CATALOG_NAMESPACE_RESOURCE_ROOT}/{namespace}")),
(None, _) => None,
}
}
}
async fn authorize_table_catalog_resource_request(
req: &S3Request<Body>,
resource: &TableCatalogResource<'_>,
action: AdminAction,
) -> S3Result<()> {
let Some(input_cred) = &req.credentials else {
return Err(s3_error!(InvalidRequest, "authentication required"));
};
@@ -282,14 +341,15 @@ async fn authorize_table_catalog_warehouse_request(req: &S3Request<Body>, wareho
let (cred, owner) =
check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?;
validate_admin_request_with_bucket(
let object_path = resource.object_path();
validate_admin_request_with_bucket_object(
&req.headers,
&cred,
owner,
false,
vec![Action::AdminAction(action)],
req.extensions.get::<Option<RemoteAddr>>().and_then(|opt| opt.map(|a| a.0)),
warehouse,
AdminResourceScope::bucket_object(resource.warehouse, object_path.as_deref().unwrap_or("")),
)
.await
}
@@ -427,10 +487,15 @@ fn list_tables_response_from_entries(entries: Vec<crate::table_catalog::TableEnt
}
fn load_table_response_from_entry(entry: crate::table_catalog::TableEntry, metadata: serde_json::Value) -> RestLoadTableResponse {
let mut config = BTreeMap::new();
config.insert("warehouse-location".to_string(), entry.warehouse_location);
config.insert(CREDENTIAL_VENDING_CONFIG_KEY.to_string(), CREDENTIAL_VENDING_UNSUPPORTED.to_string());
RestLoadTableResponse {
metadata_location: entry.metadata_location,
metadata,
config: BTreeMap::from([("warehouse-location".to_string(), entry.warehouse_location)]),
config,
storage_credentials: Vec::new(),
}
}
@@ -1513,7 +1578,8 @@ pub struct RestListNamespacesHandler {}
impl Operation for RestListNamespacesHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::GetTableNamespaceAction).await?;
let resource = TableCatalogResource::warehouse(&warehouse);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableNamespaceAction).await?;
let store = table_catalog_store()?;
let response = list_namespaces_response(&store, &warehouse).await?;
build_json_response(StatusCode::OK, &response)
@@ -1526,7 +1592,8 @@ pub struct RestCreateNamespaceHandler {}
impl Operation for RestCreateNamespaceHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::SetTableNamespaceAction).await?;
let resource = TableCatalogResource::warehouse(&warehouse);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::SetTableNamespaceAction).await?;
let request = read_json_body::<CreateNamespaceRequest>(req.input).await?;
let store = table_catalog_store()?;
let table_bucket_enabled = table_bucket_enabled_from_metadata(&warehouse).await?;
@@ -1541,8 +1608,9 @@ pub struct RestGetNamespaceHandler {}
impl Operation for RestGetNamespaceHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::GetTableNamespaceAction).await?;
let namespace = namespace_from_params(&params)?;
let resource = TableCatalogResource::namespace(&warehouse, &namespace);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableNamespaceAction).await?;
let store = table_catalog_store()?;
let response = get_namespace_response(&store, &warehouse, &namespace).await?;
build_json_response(StatusCode::OK, &response)
@@ -1555,8 +1623,9 @@ pub struct RestDropNamespaceHandler {}
impl Operation for RestDropNamespaceHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::DeleteTableNamespaceAction).await?;
let namespace = namespace_from_params(&params)?;
let resource = TableCatalogResource::namespace(&warehouse, &namespace);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::DeleteTableNamespaceAction).await?;
let store = table_catalog_store()?;
drop_namespace_in_store(&store, &warehouse, &namespace.public_name()).await?;
Ok(S3Response::new((StatusCode::NO_CONTENT, Body::default())))
@@ -1569,8 +1638,9 @@ pub struct RestListTablesHandler {}
impl Operation for RestListTablesHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::GetTableAction).await?;
let namespace = namespace_from_params(&params)?;
let resource = TableCatalogResource::namespace(&warehouse, &namespace);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableAction).await?;
let store = table_catalog_store()?;
let response = list_tables_response(&store, &warehouse, &namespace).await?;
build_json_response(StatusCode::OK, &response)
@@ -1583,8 +1653,9 @@ pub struct RestCreateTableHandler {}
impl Operation for RestCreateTableHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::CreateTableAction).await?;
let namespace = namespace_from_params(&params)?;
let resource = TableCatalogResource::namespace(&warehouse, &namespace);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::CreateTableAction).await?;
let request = read_json_body::<CreateTableRequest>(req.input).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
@@ -1601,8 +1672,9 @@ pub struct RestRegisterTableHandler {}
impl Operation for RestRegisterTableHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::RegisterTableAction).await?;
let namespace = namespace_from_params(&params)?;
let resource = TableCatalogResource::namespace(&warehouse, &namespace);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::RegisterTableAction).await?;
let request = read_json_body::<RegisterTableRequest>(req.input).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
@@ -1619,9 +1691,10 @@ pub struct RestLoadTableHandler {}
impl Operation for RestLoadTableHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::GetTableMetadataAction).await?;
let namespace = namespace_from_params(&params)?;
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableMetadataAction).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
let response = load_table_response(&store, &metadata_backend, &warehouse, &namespace, &table).await?;
@@ -1635,9 +1708,10 @@ pub struct RestCommitTableHandler {}
impl Operation for RestCommitTableHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::CommitTableAction).await?;
let namespace = namespace_from_params(&params)?;
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?;
let request = read_json_body::<RestCommitTableRequest>(req.input).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
@@ -1652,9 +1726,10 @@ pub struct RestDropTableHandler {}
impl Operation for RestDropTableHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::DeleteTableAction).await?;
let namespace = namespace_from_params(&params)?;
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::DeleteTableAction).await?;
let store = table_catalog_store()?;
drop_table_in_store(&store, &warehouse, &namespace, &table).await?;
Ok(S3Response::new((StatusCode::NO_CONTENT, Body::default())))
@@ -1667,9 +1742,10 @@ pub struct RestTableMetadataMaintenanceHandler {}
impl Operation for RestTableMetadataMaintenanceHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::RunTableMaintenanceAction).await?;
let namespace = namespace_from_params(&params)?;
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::RunTableMaintenanceAction).await?;
let request = read_json_body::<TableMetadataMaintenanceRequest>(req.input).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend);
@@ -1722,6 +1798,10 @@ mod tests {
operation_block(src, "GetCatalogConfigHandler")
.contains("authorize_table_catalog_request(&req, AdminAction::GetTableCatalogAction).await?;")
);
assert!(
src.contains("validate_admin_request_with_bucket_object("),
"catalog resource auth should pass namespace/table scope into IAM object matching"
);
for (handler, action) in [
("RestListNamespacesHandler", "AdminAction::GetTableNamespaceAction"),
@@ -1738,14 +1818,55 @@ mod tests {
] {
let block = operation_block(src, handler);
assert!(
block.contains(&format!("authorize_table_catalog_warehouse_request(&req, &warehouse, {action}).await?;")),
"{handler} should require {action} with warehouse-scoped auth"
block.contains(&format!("authorize_table_catalog_resource_request(&req, &resource, {action}).await?;")),
"{handler} should require {action} with catalog resource auth"
);
assert!(
!block.contains("authorize_table_catalog_request(&req,"),
"{handler} must not use unscoped table catalog authorization"
);
assert!(
!block.contains("authorize_table_catalog_warehouse_request(&req, &warehouse,"),
"{handler} should not bypass catalog resource auth"
);
}
for (handler, action) in [
("RestLoadTableHandler", "AdminAction::GetTableMetadataAction"),
("RestCommitTableHandler", "AdminAction::CommitTableAction"),
("RestDropTableHandler", "AdminAction::DeleteTableAction"),
("RestTableMetadataMaintenanceHandler", "AdminAction::RunTableMaintenanceAction"),
] {
let block = operation_block(src, handler);
assert!(
block.contains("TableCatalogResource::table(&warehouse, &namespace, &table)"),
"{handler} should build a table-aware catalog resource"
);
assert!(
block.contains(&format!("authorize_table_catalog_resource_request(&req, &resource, {action}).await?;")),
"{handler} should authorize against the table-aware catalog resource"
);
}
}
#[test]
fn table_catalog_resource_builds_policy_object_scope() {
let namespace = crate::table_catalog::Namespace::parse("analytics.daily_events").expect("namespace should parse");
let table = crate::table_catalog::IdentifierSegment::parse("events").expect("table should parse");
assert_eq!(TableCatalogResource::warehouse("warehouse-a").object_path(), None);
assert_eq!(
TableCatalogResource::namespace("warehouse-a", &namespace)
.object_path()
.as_deref(),
Some("namespaces/analytics/daily_events")
);
assert_eq!(
TableCatalogResource::table("warehouse-a", &namespace, table.as_str())
.object_path()
.as_deref(),
Some("namespaces/analytics/daily_events/tables/events")
);
}
fn operation_block<'a>(src: &'a str, handler: &str) -> &'a str {
@@ -2325,6 +2446,11 @@ mod tests {
);
assert_eq!(response.metadata, metadata);
assert!(response.storage_credentials.is_empty());
assert_eq!(response.config.get("rustfs.credential-vending"), Some(&"unsupported".to_string()));
assert!(!response.config.contains_key("s3.access-key-id"));
assert!(!response.config.contains_key("s3.secret-access-key"));
assert!(!response.config.contains_key("s3.session-token"));
}
#[test]