feat(table-catalog): add REST exists endpoints (#3395)

This commit is contained in:
GatewayJ
2026-06-13 13:54:39 +08:00
committed by GitHub
parent 74604201bb
commit bda0b1f3dd
3 changed files with 244 additions and 3 deletions
+195 -2
View File
@@ -43,11 +43,13 @@ const TABLE_CATALOG_ENDPOINTS: &[&str] = &[
"GET /v1/{prefix}/namespaces",
"POST /v1/{prefix}/namespaces",
"GET /v1/{prefix}/namespaces/{namespace}",
"HEAD /v1/{prefix}/namespaces/{namespace}",
"DELETE /v1/{prefix}/namespaces/{namespace}",
"GET /v1/{prefix}/namespaces/{namespace}/tables",
"POST /v1/{prefix}/namespaces/{namespace}/tables",
"POST /v1/{prefix}/namespaces/{namespace}/register",
"GET /v1/{prefix}/namespaces/{namespace}/tables/{table}",
"HEAD /v1/{prefix}/namespaces/{namespace}/tables/{table}",
"POST /v1/{prefix}/namespaces/{namespace}/tables/{table}",
"DELETE /v1/{prefix}/namespaces/{namespace}/tables/{table}",
"PUT /buckets/{warehouse}",
@@ -55,11 +57,13 @@ const TABLE_CATALOG_ENDPOINTS: &[&str] = &[
"GET /{warehouse}/namespaces",
"POST /{warehouse}/namespaces",
"GET /{warehouse}/namespaces/{namespace}",
"HEAD /{warehouse}/namespaces/{namespace}",
"DELETE /{warehouse}/namespaces/{namespace}",
"GET /{warehouse}/namespaces/{namespace}/tables",
"POST /{warehouse}/namespaces/{namespace}/tables",
"POST /{warehouse}/namespaces/{namespace}/register",
"GET /{warehouse}/namespaces/{namespace}/tables/{table}",
"HEAD /{warehouse}/namespaces/{namespace}/tables/{table}",
"POST /{warehouse}/namespaces/{namespace}/tables/{table}",
"DELETE /{warehouse}/namespaces/{namespace}/tables/{table}",
"POST /{warehouse}/namespaces/{namespace}/tables/{table}/maintenance/metadata",
@@ -80,11 +84,13 @@ static GET_TABLE_BUCKET_HANDLER: GetTableBucketHandler = GetTableBucketHandler {
static LIST_NAMESPACES_HANDLER: RestListNamespacesHandler = RestListNamespacesHandler {};
static CREATE_NAMESPACE_HANDLER: RestCreateNamespaceHandler = RestCreateNamespaceHandler {};
static GET_NAMESPACE_HANDLER: RestGetNamespaceHandler = RestGetNamespaceHandler {};
static NAMESPACE_EXISTS_HANDLER: RestNamespaceExistsHandler = RestNamespaceExistsHandler {};
static DROP_NAMESPACE_HANDLER: RestDropNamespaceHandler = RestDropNamespaceHandler {};
static LIST_TABLES_HANDLER: RestListTablesHandler = RestListTablesHandler {};
static CREATE_TABLE_HANDLER: RestCreateTableHandler = RestCreateTableHandler {};
static REGISTER_TABLE_HANDLER: RestRegisterTableHandler = RestRegisterTableHandler {};
static LOAD_TABLE_HANDLER: RestLoadTableHandler = RestLoadTableHandler {};
static TABLE_EXISTS_HANDLER: RestTableExistsHandler = RestTableExistsHandler {};
static COMMIT_TABLE_HANDLER: RestCommitTableHandler = RestCommitTableHandler {};
static DROP_TABLE_HANDLER: RestDropTableHandler = RestDropTableHandler {};
static GET_TABLE_METADATA_LOCATION_HANDLER: GetTableMetadataLocationHandler = GetTableMetadataLocationHandler {};
@@ -326,6 +332,11 @@ fn register_table_catalog_prefix_routes(r: &mut S3Router<AdminOperation>, prefix
format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}").as_str(),
AdminOperation(&GET_NAMESPACE_HANDLER),
)?;
r.insert(
Method::HEAD,
format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}").as_str(),
AdminOperation(&NAMESPACE_EXISTS_HANDLER),
)?;
r.insert(
Method::DELETE,
format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}").as_str(),
@@ -351,6 +362,11 @@ fn register_table_catalog_prefix_routes(r: &mut S3Router<AdminOperation>, prefix
format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}/tables/{{table}}").as_str(),
AdminOperation(&LOAD_TABLE_HANDLER),
)?;
r.insert(
Method::HEAD,
format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}/tables/{{table}}").as_str(),
AdminOperation(&TABLE_EXISTS_HANDLER),
)?;
r.insert(
Method::POST,
format!("{prefix}/{{warehouse}}/namespaces/{{namespace}}/tables/{{table}}").as_str(),
@@ -434,6 +450,18 @@ fn build_json_response<T: Serialize>(status: StatusCode, body: &T) -> S3Result<S
Ok(S3Response::with_headers((status, Body::from(data)), headers))
}
fn empty_response(status: StatusCode) -> S3Response<(StatusCode, Body)> {
S3Response::new((status, Body::default()))
}
fn exists_status(exists: bool) -> StatusCode {
if exists {
StatusCode::NO_CONTENT
} else {
StatusCode::NOT_FOUND
}
}
async fn authorize_table_catalog_request(req: &S3Request<Body>, action: AdminAction) -> S3Result<()> {
let Some(input_cred) = &req.credentials else {
return Err(s3_error!(InvalidRequest, "authentication required"));
@@ -1651,6 +1679,18 @@ where
namespace_response_from_entry(entry)
}
async fn namespace_exists_status<S>(store: &S, bucket: &str, namespace: &crate::table_catalog::Namespace) -> S3Result<StatusCode>
where
S: crate::table_catalog::TableCatalogStore + ?Sized,
{
let exists = store
.get_namespace(bucket, &namespace.public_name())
.await
.map_err(catalog_store_error)?
.is_some();
Ok(exists_status(exists))
}
async fn drop_namespace_in_store<S>(store: &S, bucket: &str, namespace: &str) -> S3Result<()>
where
S: crate::table_catalog::TableCatalogStore + ?Sized,
@@ -1762,6 +1802,23 @@ where
Ok(load_table_response_from_entry(entry, metadata))
}
async fn table_exists_status<S>(
store: &S,
bucket: &str,
namespace: &crate::table_catalog::Namespace,
table: &str,
) -> S3Result<StatusCode>
where
S: crate::table_catalog::TableCatalogStore + ?Sized,
{
let exists = store
.load_table(bucket, &namespace.public_name(), table)
.await
.map_err(catalog_store_error)?
.is_some();
Ok(exists_status(exists))
}
async fn get_table_metadata_location_response<S>(
store: &S,
bucket: &str,
@@ -2125,7 +2182,21 @@ impl Operation for RestDropNamespaceHandler {
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())))
Ok(empty_response(StatusCode::NO_CONTENT))
}
}
pub struct RestNamespaceExistsHandler {}
#[async_trait::async_trait]
impl Operation for RestNamespaceExistsHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
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()?;
Ok(empty_response(namespace_exists_status(&store, &warehouse, &namespace).await?))
}
}
@@ -2199,6 +2270,21 @@ impl Operation for RestLoadTableHandler {
}
}
pub struct RestTableExistsHandler {}
#[async_trait::async_trait]
impl Operation for RestTableExistsHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let warehouse = warehouse_from_params(&params)?;
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::GetTableAction).await?;
let store = table_catalog_store()?;
Ok(empty_response(table_exists_status(&store, &warehouse, &namespace, &table).await?))
}
}
pub struct RestCommitTableHandler {}
#[async_trait::async_trait]
@@ -2229,7 +2315,7 @@ impl Operation for RestDropTableHandler {
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())))
Ok(empty_response(StatusCode::NO_CONTENT))
}
}
@@ -2447,13 +2533,20 @@ mod tests {
);
assert!(response.overrides.is_empty());
assert!(response.endpoints.contains(&"GET /v1/{prefix}/namespaces"));
assert!(response.endpoints.contains(&"HEAD /v1/{prefix}/namespaces/{namespace}"));
assert!(
response
.endpoints
.contains(&"GET /v1/{prefix}/namespaces/{namespace}/tables/{table}")
);
assert!(
response
.endpoints
.contains(&"HEAD /v1/{prefix}/namespaces/{namespace}/tables/{table}")
);
assert!(response.endpoints.contains(&"GET /{warehouse}/namespaces"));
assert!(response.endpoints.contains(&"POST /{warehouse}/namespaces"));
assert!(response.endpoints.contains(&"HEAD /{warehouse}/namespaces/{namespace}"));
assert!(
response
.endpoints
@@ -2469,6 +2562,11 @@ mod tests {
.endpoints
.contains(&"GET /{warehouse}/namespaces/{namespace}/tables/{table}")
);
assert!(
response
.endpoints
.contains(&"HEAD /{warehouse}/namespaces/{namespace}/tables/{table}")
);
assert!(
response
.endpoints
@@ -2495,11 +2593,13 @@ mod tests {
("RestListNamespacesHandler", "AdminAction::GetTableNamespaceAction"),
("RestCreateNamespaceHandler", "AdminAction::SetTableNamespaceAction"),
("RestGetNamespaceHandler", "AdminAction::GetTableNamespaceAction"),
("RestNamespaceExistsHandler", "AdminAction::GetTableNamespaceAction"),
("RestDropNamespaceHandler", "AdminAction::DeleteTableNamespaceAction"),
("RestListTablesHandler", "AdminAction::GetTableAction"),
("RestCreateTableHandler", "AdminAction::CreateTableAction"),
("RestRegisterTableHandler", "AdminAction::RegisterTableAction"),
("RestLoadTableHandler", "AdminAction::GetTableMetadataAction"),
("RestTableExistsHandler", "AdminAction::GetTableAction"),
("RestCommitTableHandler", "AdminAction::CommitTableAction"),
("RestDropTableHandler", "AdminAction::DeleteTableAction"),
("GetTableMetadataLocationHandler", "AdminAction::GetTableMetadataLocationAction"),
@@ -2530,6 +2630,7 @@ mod tests {
for (handler, action) in [
("RestLoadTableHandler", "AdminAction::GetTableMetadataAction"),
("RestTableExistsHandler", "AdminAction::GetTableAction"),
("RestCommitTableHandler", "AdminAction::CommitTableAction"),
("RestDropTableHandler", "AdminAction::DeleteTableAction"),
("GetTableMetadataLocationHandler", "AdminAction::GetTableMetadataLocationAction"),
@@ -2594,11 +2695,13 @@ mod tests {
let _: &RestListNamespacesHandler = &LIST_NAMESPACES_HANDLER;
let _: &RestCreateNamespaceHandler = &CREATE_NAMESPACE_HANDLER;
let _: &RestGetNamespaceHandler = &GET_NAMESPACE_HANDLER;
let _: &RestNamespaceExistsHandler = &NAMESPACE_EXISTS_HANDLER;
let _: &RestDropNamespaceHandler = &DROP_NAMESPACE_HANDLER;
let _: &RestListTablesHandler = &LIST_TABLES_HANDLER;
let _: &RestCreateTableHandler = &CREATE_TABLE_HANDLER;
let _: &RestRegisterTableHandler = &REGISTER_TABLE_HANDLER;
let _: &RestLoadTableHandler = &LOAD_TABLE_HANDLER;
let _: &RestTableExistsHandler = &TABLE_EXISTS_HANDLER;
let _: &RestCommitTableHandler = &COMMIT_TABLE_HANDLER;
let _: &RestDropTableHandler = &DROP_TABLE_HANDLER;
let _: &GetTableMetadataLocationHandler = &GET_TABLE_METADATA_LOCATION_HANDLER;
@@ -2617,11 +2720,13 @@ mod tests {
assert_operation::<RestListNamespacesHandler>();
assert_operation::<RestCreateNamespaceHandler>();
assert_operation::<RestGetNamespaceHandler>();
assert_operation::<RestNamespaceExistsHandler>();
assert_operation::<RestDropNamespaceHandler>();
assert_operation::<RestListTablesHandler>();
assert_operation::<RestCreateTableHandler>();
assert_operation::<RestRegisterTableHandler>();
assert_operation::<RestLoadTableHandler>();
assert_operation::<RestTableExistsHandler>();
assert_operation::<RestCommitTableHandler>();
assert_operation::<RestDropTableHandler>();
assert_operation::<GetTableMetadataLocationHandler>();
@@ -2821,6 +2926,94 @@ mod tests {
assert_eq!(response.identifiers[0].name, "events");
}
#[tokio::test]
async fn namespace_exists_status_uses_head_rest_semantics() {
let store = TestTableCatalogStore::default();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
assert_eq!(
namespace_exists_status(&store, "warehouse", &namespace)
.await
.expect("missing namespace check should succeed"),
StatusCode::NOT_FOUND
);
create_namespace_response(
&store,
"warehouse",
CreateNamespaceRequest {
namespace: vec!["analytics".to_string()],
properties: BTreeMap::new(),
},
true,
)
.await
.expect("namespace should be created");
assert_eq!(
namespace_exists_status(&store, "warehouse", &namespace)
.await
.expect("existing namespace check should succeed"),
StatusCode::NO_CONTENT
);
}
#[tokio::test]
async fn table_exists_status_uses_head_rest_semantics() {
let store = TestTableCatalogStore::default();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
create_namespace_response(
&store,
"warehouse",
CreateNamespaceRequest {
namespace: vec!["analytics".to_string()],
properties: BTreeMap::new(),
},
true,
)
.await
.expect("namespace should be created");
assert_eq!(
table_exists_status(&store, "warehouse", &namespace, "events")
.await
.expect("missing table check should succeed"),
StatusCode::NOT_FOUND
);
store
.create_table(crate::table_catalog::TableEntry {
version: crate::table_catalog::TABLE_CATALOG_ENTRY_VERSION,
table_bucket: "warehouse".to_string(),
namespace: namespace.public_name(),
table: "events".to_string(),
table_id: "table-id".to_string(),
table_uuid: "table-uuid".to_string(),
format: "ICEBERG".to_string(),
format_version: 2,
warehouse_location: "s3://warehouse/tables/table-id".to_string(),
metadata_location:
".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00001.metadata.json"
.to_string(),
version_token: "token-v1".to_string(),
generation: 1,
state: crate::table_catalog::TableCatalogEntryState::Active,
properties: BTreeMap::new(),
created_at: None,
updated_at: None,
})
.await
.expect("table should be created");
assert_eq!(
table_exists_status(&store, "warehouse", &namespace, "events")
.await
.expect("existing table check should succeed"),
StatusCode::NO_CONTENT
);
}
#[test]
fn register_table_request_builds_initial_table_entry() {
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
+37 -1
View File
@@ -633,6 +633,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
GET_TABLE_NAMESPACE,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Head,
"/iceberg/v1/{warehouse}/namespaces/{namespace}",
GET_TABLE_NAMESPACE,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Delete,
"/iceberg/v1/{warehouse}/namespaces/{namespace}",
@@ -663,6 +669,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
GET_TABLE_METADATA,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Head,
"/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}",
GET_TABLE,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Post,
"/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}",
@@ -766,6 +778,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
GET_TABLE_NAMESPACE,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Head,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}",
GET_TABLE_NAMESPACE,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Delete,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}",
@@ -796,6 +814,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
GET_TABLE_METADATA,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Head,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}",
GET_TABLE,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Post,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}",
@@ -1041,11 +1065,13 @@ mod tests {
let table_specs = ADMIN_ROUTE_POLICY_SPECS
.iter()
.filter(|spec| spec.path().starts_with("/iceberg/v1") || spec.path().starts_with("/_iceberg/v1"));
assert_eq!(table_specs.count(), 46);
assert_eq!(table_specs.count(), 50);
assert_action(HttpMethod::Put, "/iceberg/v1/buckets/{warehouse}", SET_TABLE_BUCKET);
assert_action(HttpMethod::Get, "/_iceberg/v1/buckets/{warehouse}", GET_TABLE_BUCKET);
assert_action(HttpMethod::Get, "/iceberg/v1/{warehouse}/namespaces", GET_TABLE_NAMESPACE);
assert_action(HttpMethod::Get, "/_iceberg/v1/{warehouse}/namespaces", GET_TABLE_NAMESPACE);
assert_action(HttpMethod::Head, "/iceberg/v1/{warehouse}/namespaces/{namespace}", GET_TABLE_NAMESPACE);
assert_action(HttpMethod::Head, "/_iceberg/v1/{warehouse}/namespaces/{namespace}", GET_TABLE_NAMESPACE);
assert_action(HttpMethod::Post, "/iceberg/v1/{warehouse}/namespaces/{namespace}/tables", CREATE_TABLE);
assert_action(HttpMethod::Post, "/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables", CREATE_TABLE);
assert_action(
@@ -1058,6 +1084,16 @@ mod tests {
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}",
GET_TABLE_METADATA,
);
assert_action(
HttpMethod::Head,
"/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}",
GET_TABLE,
);
assert_action(
HttpMethod::Head,
"/_iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}",
GET_TABLE,
);
assert_action(
HttpMethod::Post,
"/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}",
@@ -288,6 +288,7 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
table_route_sample(Method::GET, "/{warehouse}/namespaces", "/analytics/namespaces"),
table_route_sample(Method::POST, "/{warehouse}/namespaces", "/analytics/namespaces"),
table_route_sample(Method::GET, "/{warehouse}/namespaces/{namespace}", "/analytics/namespaces/sales"),
table_route_sample(Method::HEAD, "/{warehouse}/namespaces/{namespace}", "/analytics/namespaces/sales"),
table_route_sample(Method::DELETE, "/{warehouse}/namespaces/{namespace}", "/analytics/namespaces/sales"),
table_route_sample(
Method::GET,
@@ -309,6 +310,11 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
"/{warehouse}/namespaces/{namespace}/tables/{table}",
"/analytics/namespaces/sales/tables/orders",
),
table_route_sample(
Method::HEAD,
"/{warehouse}/namespaces/{namespace}/tables/{table}",
"/analytics/namespaces/sales/tables/orders",
),
table_route_sample(
Method::POST,
"/{warehouse}/namespaces/{namespace}/tables/{table}",
@@ -375,6 +381,7 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
compat_table_route_sample(Method::GET, "/{warehouse}/namespaces", "/analytics/namespaces"),
compat_table_route_sample(Method::POST, "/{warehouse}/namespaces", "/analytics/namespaces"),
compat_table_route_sample(Method::GET, "/{warehouse}/namespaces/{namespace}", "/analytics/namespaces/sales"),
compat_table_route_sample(Method::HEAD, "/{warehouse}/namespaces/{namespace}", "/analytics/namespaces/sales"),
compat_table_route_sample(Method::DELETE, "/{warehouse}/namespaces/{namespace}", "/analytics/namespaces/sales"),
compat_table_route_sample(
Method::GET,
@@ -396,6 +403,11 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
"/{warehouse}/namespaces/{namespace}/tables/{table}",
"/analytics/namespaces/sales/tables/orders",
),
compat_table_route_sample(
Method::HEAD,
"/{warehouse}/namespaces/{namespace}/tables/{table}",
"/analytics/namespaces/sales/tables/orders",
),
compat_table_route_sample(
Method::POST,
"/{warehouse}/namespaces/{namespace}/tables/{table}",