From 2eafd9c6024a0e2f989dae90d549e6261897e4a1 Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Wed, 10 Jun 2026 09:15:51 +0800 Subject: [PATCH] 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 Co-authored-by: houseme --- crates/policy/src/policy/policy.rs | 73 ++++++++++ rustfs/src/admin/auth.rs | 39 ++++- rustfs/src/admin/handlers/table_catalog.rs | 162 ++++++++++++++++++--- 3 files changed, 255 insertions(+), 19 deletions(-) diff --git a/crates/policy/src/policy/policy.rs b/crates/policy/src/policy/policy.rs index 7b3eb231d..c842c84db 100644 --- a/crates/policy/src/policy/policy.rs +++ b/crates/policy/src/policy/policy.rs @@ -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}; diff --git a/rustfs/src/admin/auth.rs b/rustfs/src/admin/auth.rs index b6a8531de..131610785 100644 --- a/rustfs/src/admin/auth.rs +++ b/rustfs/src/admin/auth.rs @@ -32,6 +32,22 @@ struct AuthContext<'a> { remote_addr: Option, } +#[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, remote_addr: Option, 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, + remote_addr: Option, + 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() { diff --git a/rustfs/src/admin/handlers/table_catalog.rs b/rustfs/src/admin/handlers/table_catalog.rs index c019dc9af..32d6f6714 100644 --- a/rustfs/src/admin/handlers/table_catalog.rs +++ b/rustfs/src/admin/handlers/table_catalog.rs @@ -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, } +#[derive(Debug, Serialize)] +struct RestStorageCredential { + prefix: String, + config: BTreeMap, +} + #[derive(Debug, Serialize)] struct RestLoadTableResponse { #[serde(rename = "metadata-location")] metadata_location: String, metadata: serde_json::Value, config: BTreeMap, + #[serde(rename = "storage-credentials")] + storage_credentials: Vec, } #[derive(Debug, Serialize)] @@ -274,7 +286,54 @@ async fn authorize_table_catalog_request(req: &S3Request, action: AdminAct .await } -async fn authorize_table_catalog_warehouse_request(req: &S3Request, warehouse: &str, action: AdminAction) -> S3Result<()> { +#[derive(Debug, Clone)] +struct TableCatalogResource<'a> { + warehouse: &'a str, + namespace: Option, + table: Option, +} + +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 { + 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, + 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, 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::>().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 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, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; - 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, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; - 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::(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, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; - authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::GetTableNamespaceAction).await?; let namespace = namespace_from_params(¶ms)?; + 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, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; - authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::DeleteTableNamespaceAction).await?; let namespace = namespace_from_params(¶ms)?; + 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, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; - authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::GetTableAction).await?; let namespace = namespace_from_params(¶ms)?; + 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, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; - authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::CreateTableAction).await?; let namespace = namespace_from_params(¶ms)?; + let resource = TableCatalogResource::namespace(&warehouse, &namespace); + authorize_table_catalog_resource_request(&req, &resource, AdminAction::CreateTableAction).await?; let request = read_json_body::(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, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; - authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::RegisterTableAction).await?; let namespace = namespace_from_params(¶ms)?; + let resource = TableCatalogResource::namespace(&warehouse, &namespace); + authorize_table_catalog_resource_request(&req, &resource, AdminAction::RegisterTableAction).await?; let request = read_json_body::(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, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; - authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::GetTableMetadataAction).await?; let namespace = namespace_from_params(¶ms)?; let table = table_name_from_params(¶ms)?; + 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, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; - authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::CommitTableAction).await?; let namespace = namespace_from_params(¶ms)?; let table = table_name_from_params(¶ms)?; + let resource = TableCatalogResource::table(&warehouse, &namespace, &table); + authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?; let request = read_json_body::(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, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; - authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::DeleteTableAction).await?; let namespace = namespace_from_params(¶ms)?; let table = table_name_from_params(¶ms)?; + 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, params: Params<'_, '_>) -> S3Result> { let warehouse = warehouse_from_params(¶ms)?; - authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::RunTableMaintenanceAction).await?; let namespace = namespace_from_params(¶ms)?; let table = table_name_from_params(¶ms)?; + let resource = TableCatalogResource::table(&warehouse, &namespace, &table); + authorize_table_catalog_resource_request(&req, &resource, AdminAction::RunTableMaintenanceAction).await?; let request = read_json_body::(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]