fix(s3tables): harden catalog guard checks (#3801)

This commit is contained in:
GatewayJ
2026-06-24 12:51:40 +08:00
committed by GitHub
parent 31b5db9a75
commit 89dfb1357c
2 changed files with 203 additions and 11 deletions
+123 -7
View File
@@ -1272,6 +1272,13 @@ async fn table_bucket_enabled_from_metadata(bucket: &str) -> S3Result<bool> {
Ok(metadata.table_bucket_enabled())
}
async fn ensure_table_bucket_enabled(bucket: &str) -> S3Result<()> {
if table_bucket_enabled_from_metadata(bucket).await? {
return Ok(());
}
Err(s3_error!(InvalidRequest, "bucket {bucket} is not table-enabled"))
}
fn table_bucket_entry_from_metadata_marker(bucket: &str) -> crate::table_catalog::TableBucketEntry {
crate::table_catalog::TableBucketEntry {
version: crate::table_catalog::TABLE_CATALOG_ENTRY_VERSION,
@@ -1345,8 +1352,8 @@ async fn enable_table_bucket_response<S>(store: &S, bucket: &str) -> S3Result<Ta
where
S: crate::table_catalog::TableCatalogStore + ?Sized,
{
ensure_table_bucket_entry(store, bucket, true).await?;
enable_table_bucket_marker(bucket).await?;
ensure_table_bucket_entry(store, bucket, true).await?;
table_bucket_response(store, bucket, true).await
}
@@ -1699,10 +1706,7 @@ fn table_commit_request_from_rest_request(
}
fn validate_table_location_in_bucket(bucket: &str, location: &str) -> S3Result<()> {
if !location.starts_with(&format!("s3://{bucket}/")) {
return Err(s3_error!(InvalidRequest, "table location must be inside the table bucket"));
}
Ok(())
crate::table_catalog::validate_table_warehouse_location(bucket, location).map_err(catalog_store_error)
}
fn metadata_table_uuid(metadata: &serde_json::Value) -> S3Result<&str> {
@@ -4632,6 +4636,7 @@ impl Operation for RestListNamespacesHandler {
let warehouse = warehouse_from_params(&params)?;
let resource = TableCatalogResource::warehouse(&warehouse);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableNamespaceAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
let response = list_namespaces_response(&store, &warehouse).await?;
build_json_response(StatusCode::OK, &response)
@@ -4663,6 +4668,7 @@ impl Operation for RestGetNamespaceHandler {
let namespace = namespace_from_params(&params)?;
let resource = TableCatalogResource::namespace(&warehouse, &namespace);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableNamespaceAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
let response = get_namespace_response(&store, &warehouse, &namespace).await?;
build_json_response(StatusCode::OK, &response)
@@ -4678,6 +4684,7 @@ impl Operation for RestDropNamespaceHandler {
let namespace = namespace_from_params(&params)?;
let resource = TableCatalogResource::namespace(&warehouse, &namespace);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::DeleteTableNamespaceAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
drop_namespace_in_store(&store, &warehouse, &namespace.public_name()).await?;
Ok(empty_response(StatusCode::NO_CONTENT))
@@ -4693,6 +4700,7 @@ impl Operation for RestNamespaceExistsHandler {
let namespace = namespace_from_params(&params)?;
let resource = TableCatalogResource::namespace(&warehouse, &namespace);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableNamespaceAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
Ok(empty_response(namespace_exists_status(&store, &warehouse, &namespace).await?))
}
@@ -4707,6 +4715,7 @@ impl Operation for RestListTablesHandler {
let namespace = namespace_from_params(&params)?;
let resource = TableCatalogResource::namespace(&warehouse, &namespace);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
let response = list_tables_response(&store, &warehouse, &namespace).await?;
build_json_response(StatusCode::OK, &response)
@@ -4760,6 +4769,7 @@ impl Operation for RestListViewsHandler {
let namespace = namespace_from_params(&params)?;
let resource = TableCatalogResource::namespace(&warehouse, &namespace);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableMetadataAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
let response = list_views_response(&store, &warehouse, &namespace).await?;
build_json_response(StatusCode::OK, &response)
@@ -4795,6 +4805,7 @@ impl Operation for RestLoadTableHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableMetadataAction).await?;
ensure_table_bucket_enabled(&warehouse).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?;
@@ -4812,6 +4823,7 @@ impl Operation for RestTableExistsHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
Ok(empty_response(table_exists_status(&store, &warehouse, &namespace, &table).await?))
}
@@ -4827,6 +4839,7 @@ impl Operation for RestLoadCredentialsHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableCredentialsAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let principal = table_catalog_request_principal(&req).await?;
let store = table_catalog_store()?;
let issuer = IamTableCredentialIssuer::from_env();
@@ -4845,6 +4858,7 @@ impl Operation for RestCommitTableHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?;
ensure_table_bucket_enabled(&warehouse).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());
@@ -4863,6 +4877,7 @@ impl Operation for RestDropTableHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::DeleteTableAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
drop_table_in_store(&store, &warehouse, &namespace, &table).await?;
Ok(empty_response(StatusCode::NO_CONTENT))
@@ -4879,6 +4894,7 @@ impl Operation for RestLoadViewHandler {
let view = view_name_from_params(&params)?;
let resource = TableCatalogResource::view(&warehouse, &namespace, &view);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableMetadataAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
let response = load_view_response(&store, &metadata_backend, &warehouse, &namespace, &view).await?;
@@ -4896,6 +4912,7 @@ impl Operation for RestViewExistsHandler {
let view = view_name_from_params(&params)?;
let resource = TableCatalogResource::view(&warehouse, &namespace, &view);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
Ok(empty_response(view_exists_status(&store, &warehouse, &namespace, &view).await?))
}
@@ -4911,6 +4928,7 @@ impl Operation for RestReplaceViewHandler {
let view = view_name_from_params(&params)?;
let resource = TableCatalogResource::view(&warehouse, &namespace, &view);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let request = read_json_body::<RestCommitViewRequest>(req.input).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
@@ -4929,6 +4947,7 @@ impl Operation for RestDropViewHandler {
let view = view_name_from_params(&params)?;
let resource = TableCatalogResource::view(&warehouse, &namespace, &view);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::DeleteTableAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
drop_view_in_store(&store, &warehouse, &namespace, &view).await?;
Ok(empty_response(StatusCode::NO_CONTENT))
@@ -4945,6 +4964,7 @@ impl Operation for ListTableRefsHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableMetadataAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
let response = table_refs_response(&store, &metadata_backend, &warehouse, &namespace, &table).await?;
@@ -4963,6 +4983,7 @@ impl Operation for PutTableRefHandler {
let ref_name = ref_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let request = read_json_body::<PutTableRefRequest>(req.input).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
@@ -4983,6 +5004,7 @@ impl Operation for DeleteTableRefHandler {
let ref_name = ref_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let request = read_json_body_or_default::<DeleteTableRefRequest>(req.input).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
@@ -5002,6 +5024,7 @@ impl Operation for GetTableMetadataLocationHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableMetadataLocationAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
let response = get_table_metadata_location_response(&store, &warehouse, &namespace, &table).await?;
build_json_response(StatusCode::OK, &response)
@@ -5018,6 +5041,7 @@ impl Operation for UpdateTableMetadataLocationHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::SetTableMetadataLocationAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let request = read_json_body::<UpdateTableMetadataLocationRequest>(req.input).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
@@ -5037,6 +5061,7 @@ impl Operation for RestTableMetadataMaintenanceHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::RunTableMaintenanceAction).await?;
ensure_table_bucket_enabled(&warehouse).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.clone());
@@ -5056,6 +5081,7 @@ impl Operation for GetTableMaintenanceConfigHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableLifecycleAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
let response = store
.get_table_maintenance_config(&warehouse, &namespace.public_name(), &table)
@@ -5075,6 +5101,7 @@ impl Operation for PutTableMaintenanceConfigHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::SetTableLifecycleAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let request = read_json_body::<crate::table_catalog::TableMaintenanceConfig>(req.input).await?;
let store = table_catalog_store()?;
let response = store
@@ -5096,6 +5123,7 @@ impl Operation for GetTableMaintenanceJobHandler {
let job = job_id_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableLifecycleAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
let Some(response) = store
.get_table_metadata_maintenance_report(&warehouse, &namespace.public_name(), &table, &job)
@@ -5118,6 +5146,7 @@ impl Operation for RunTableMaintenanceWorkerHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::RunTableMaintenanceAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let request = read_json_body::<TableMaintenanceWorkerRunRequest>(req.input).await?;
let store = table_catalog_store()?;
let response = store
@@ -5144,6 +5173,7 @@ impl Operation for HeartbeatTableMaintenanceJobHandler {
let job = job_id_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::RunTableMaintenanceAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let request = read_json_body::<TableMaintenanceHeartbeatRequest>(req.input).await?;
let store = table_catalog_store()?;
let response = store
@@ -5171,6 +5201,7 @@ impl Operation for ExportTableCatalogHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableMetadataAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
let started = Instant::now();
let result = store
@@ -5214,6 +5245,7 @@ impl Operation for ExternalCatalogBridgeHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableMetadataAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
let response = external_catalog_bridge_response(&store, &warehouse, &namespace, &table).await?;
build_json_response(StatusCode::OK, &response)
@@ -5230,6 +5262,7 @@ impl Operation for PutExternalCatalogBridgeHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::RegisterTableAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let request = read_json_body::<ExternalCatalogBridgeRequest>(req.input).await?;
let store = table_catalog_store()?;
let response = put_external_catalog_bridge_response(&store, &warehouse, &namespace, &table, request).await?;
@@ -5247,6 +5280,7 @@ impl Operation for SyncExternalCatalogBridgeHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::SetTableMetadataLocationAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
if store
@@ -5283,6 +5317,7 @@ impl Operation for GetTableCatalogDiagnosticsHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::GetTableMetadataAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
let config = store
.get_table_maintenance_config(&warehouse, &namespace.public_name(), &table)
@@ -5316,6 +5351,7 @@ impl Operation for RecoverTableCatalogHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let store = table_catalog_store()?;
let started = Instant::now();
let result = store
@@ -5338,6 +5374,7 @@ impl Operation for RollbackTableCatalogHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::CommitTableAction).await?;
ensure_table_bucket_enabled(&warehouse).await?;
let request = read_json_body::<RollbackTableRequest>(req.input).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
@@ -5583,6 +5620,76 @@ mod tests {
}
}
#[test]
fn table_catalog_handlers_require_enabled_table_bucket_marker_before_catalog_state() {
let src = include_str!("table_catalog.rs");
for handler in [
"RestListNamespacesHandler",
"RestCreateNamespaceHandler",
"RestGetNamespaceHandler",
"RestNamespaceExistsHandler",
"RestDropNamespaceHandler",
"RestListTablesHandler",
"RestCreateTableHandler",
"RestRegisterTableHandler",
"RestListViewsHandler",
"RestCreateViewHandler",
"RestLoadTableHandler",
"RestTableExistsHandler",
"RestLoadCredentialsHandler",
"RestCommitTableHandler",
"RestDropTableHandler",
"RestLoadViewHandler",
"RestViewExistsHandler",
"RestReplaceViewHandler",
"RestDropViewHandler",
"ListTableRefsHandler",
"PutTableRefHandler",
"DeleteTableRefHandler",
"GetTableMetadataLocationHandler",
"UpdateTableMetadataLocationHandler",
"RestTableMetadataMaintenanceHandler",
"GetTableMaintenanceConfigHandler",
"PutTableMaintenanceConfigHandler",
"GetTableMaintenanceJobHandler",
"RunTableMaintenanceWorkerHandler",
"HeartbeatTableMaintenanceJobHandler",
"ExportTableCatalogHandler",
"ImportTableCatalogHandler",
"ExternalCatalogBridgeHandler",
"PutExternalCatalogBridgeHandler",
"SyncExternalCatalogBridgeHandler",
"GetTableCatalogDiagnosticsHandler",
"RecoverTableCatalogHandler",
"RollbackTableCatalogHandler",
] {
let block = operation_block(src, handler);
assert!(
block.contains("ensure_table_bucket_enabled(&warehouse).await?;")
|| block.contains("table_bucket_enabled_from_metadata(&warehouse).await?;"),
"{handler} should require the table bucket metadata marker before catalog state access"
);
}
}
#[test]
fn enable_table_bucket_response_writes_metadata_marker_before_catalog_entry() {
let src = include_str!("table_catalog.rs");
let block = function_block(src, "async fn enable_table_bucket_response");
let marker_write = block
.find("enable_table_bucket_marker(bucket).await?;")
.expect("enable should write the metadata marker");
let catalog_entry_write = block
.find("ensure_table_bucket_entry(store, bucket, true).await?;")
.expect("enable should write the catalog entry");
assert!(
marker_write < catalog_entry_write,
"enable should write the metadata marker before the catalog entry"
);
}
#[test]
fn table_catalog_resource_builds_policy_object_scope() {
let namespace = crate::table_catalog::Namespace::parse("analytics.daily_events").expect("namespace should parse");
@@ -5613,6 +5720,15 @@ mod tests {
&block[..end]
}
fn function_block<'a>(src: &'a str, signature: &str) -> &'a str {
let block = src.split_once(signature).expect("function should exist").1;
let end = block
.find("\nfn ")
.or_else(|| block.find("\nasync fn "))
.unwrap_or(block.len());
&block[..end]
}
#[test]
fn rest_catalog_mvp_routes_use_implemented_handlers() {
fn assert_operation<T: Operation>() {}
@@ -9697,11 +9813,11 @@ mod tests {
}
#[tokio::test]
async fn enable_table_bucket_response_fails_before_marker_when_catalog_entry_fails() {
async fn ensure_table_bucket_entry_propagates_catalog_entry_failure() {
let store = TestTableCatalogStore::default();
*store.fail_put_table_bucket.lock().await = true;
assert!(enable_table_bucket_response(&store, "warehouse").await.is_err());
assert!(ensure_table_bucket_entry(&store, "warehouse", true).await.is_err());
assert!(
store
.get_table_bucket("warehouse")
+80 -4
View File
@@ -1194,6 +1194,10 @@ fn table_warehouse_object_prefix_from_location(table_bucket: &str, warehouse_loc
normalize_warehouse_object_prefix(object_prefix)
}
pub(crate) fn validate_table_warehouse_location(table_bucket: &str, warehouse_location: &str) -> TableCatalogStoreResult<()> {
table_warehouse_object_prefix_from_location(table_bucket, warehouse_location).map(|_| ())
}
pub(crate) fn table_warehouse_object_prefix(entry: &TableEntry) -> TableCatalogStoreResult<String> {
table_warehouse_object_prefix_from_location(&entry.table_bucket, &entry.warehouse_location)
}
@@ -1814,6 +1818,7 @@ where
self.require_table_bucket(&entry.table_bucket).await?;
let namespace = parse_namespace_for_store(&entry.namespace)?;
let table = parse_table_for_store(&entry.table)?;
validate_table_warehouse_location(&entry.table_bucket, &entry.warehouse_location)?;
if self.get_namespace(&entry.table_bucket, &entry.namespace).await?.is_none() {
return Err(TableCatalogStoreError::NotFound(format!(
"namespace {}/{}",
@@ -1830,6 +1835,7 @@ where
self.require_table_bucket(&entry.table_bucket).await?;
let namespace = parse_namespace_for_store(&entry.namespace)?;
let view = parse_table_for_store(&entry.view)?;
validate_table_warehouse_location(&entry.table_bucket, &entry.warehouse_location)?;
if self.get_namespace(&entry.table_bucket, &entry.namespace).await?.is_none() {
return Err(TableCatalogStoreError::NotFound(format!(
"namespace {}/{}",
@@ -7002,9 +7008,14 @@ pub fn validate_object_mutation(table_bucket_enabled: bool, object_key: &str) ->
}
pub(crate) async fn validate_bucket_object_mutation(bucket: &str, object_key: &str) -> Result<(), TableObjectMutationError> {
if !is_reserved_table_object_key(object_key) {
return Ok(());
}
let table_bucket_enabled = get_bucket_metadata(bucket)
.await
.is_ok_and(|metadata| metadata.table_bucket_enabled());
.map(|metadata| metadata.table_bucket_enabled())
.unwrap_or(true);
validate_object_mutation(table_bucket_enabled, object_key)
}
@@ -7080,9 +7091,15 @@ mod tests {
}
#[tokio::test]
async fn bucket_object_mutation_guard_allows_when_bucket_metadata_is_unavailable() {
assert!(
async fn bucket_object_mutation_guard_fails_closed_for_reserved_prefix_when_bucket_metadata_is_unavailable() {
assert_eq!(
validate_bucket_object_mutation("missing-bucket", ".rustfs-table/current.json")
.await
.unwrap_err(),
TableObjectMutationError::ReservedCatalogObject
);
assert!(
validate_bucket_object_mutation("missing-bucket", "ordinary/current.json")
.await
.is_ok()
);
@@ -8055,6 +8072,56 @@ mod tests {
assert_eq!(resource.catalog_resource_object(), "namespaces/sales/tables/orders_child");
}
#[tokio::test]
async fn object_table_catalog_store_rejects_invalid_table_warehouse_location() {
let backend = TestCatalogObjectBackend::default();
let store = ObjectTableCatalogStore::new(backend);
let bucket = "analytics";
let namespace = Namespace::parse("sales").unwrap();
let table = IdentifierSegment::parse("orders").unwrap();
let current = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json");
store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap();
store
.create_namespace(test_namespace_entry(bucket, &namespace))
.await
.unwrap();
let mut entry = test_table_entry(bucket, &namespace, &table, current);
entry.warehouse_location = format!("s3://{bucket}/tables/../table-id");
let error = store.create_table(entry).await.unwrap_err();
assert!(matches!(
error,
TableCatalogStoreError::Invalid(message) if message.contains("invalid path segment")
));
}
#[tokio::test]
async fn object_table_catalog_store_rejects_invalid_view_warehouse_location() {
let backend = TestCatalogObjectBackend::default();
let store = ObjectTableCatalogStore::new(backend);
let bucket = "analytics";
let namespace = Namespace::parse("sales").unwrap();
let view = IdentifierSegment::parse("recent_orders").unwrap();
let current = default_view_metadata_file_path(&namespace, &view, "00001.metadata.json");
store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap();
store
.create_namespace(test_namespace_entry(bucket, &namespace))
.await
.unwrap();
let mut entry = test_view_entry(bucket, &namespace, &view, current);
entry.warehouse_location = format!("s3://{bucket}/views/../view-id");
let error = store.create_view(entry).await.unwrap_err();
assert!(matches!(
error,
TableCatalogStoreError::Invalid(message) if message.contains("invalid path segment")
));
}
#[tokio::test]
async fn table_data_plane_resource_skips_invalid_warehouse_locations() {
let backend = TestCatalogObjectBackend::default();
@@ -8074,7 +8141,16 @@ mod tests {
let mut invalid_entry = test_table_entry(bucket, &namespace, &invalid_table, current.clone());
invalid_entry.table_id = "bad-table-id".to_string();
invalid_entry.warehouse_location = format!("s3://{bucket}/");
store.create_table(invalid_entry).await.unwrap();
let invalid_path = store.paths.table_entry_path(bucket, &namespace, &invalid_table);
store
.write_entry(
store.catalog_bucket(),
&invalid_path,
&invalid_entry,
TableCatalogPutPrecondition::IfAbsent,
)
.await
.unwrap();
store
.create_table(test_table_entry(bucket, &namespace, &valid_table, current))
.await