feat(table-catalog): tighten REST load/register compatibility (#3245)

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
Co-authored-by: cxymds <Cxymds@qq.com>
This commit is contained in:
Henry Guo
2026-06-07 21:54:15 +08:00
committed by GitHub
parent 3c71d5ef1c
commit dd0ef1f766
+218 -28
View File
@@ -38,7 +38,6 @@ const TABLE_CATALOG_ENDPOINTS: &[&str] = &[
"GET /{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}",
"DELETE /{warehouse}/namespaces/{namespace}/tables/{table}",
@@ -78,7 +77,7 @@ struct RegisterTableRequest {
#[serde(rename = "metadata-location")]
metadata_location: String,
#[serde(default)]
properties: BTreeMap<String, String>,
overwrite: bool,
}
#[derive(Debug, Deserialize)]
@@ -277,11 +276,14 @@ fn table_name_from_params(params: &Params<'_, '_>) -> S3Result<String> {
Ok(table.to_string())
}
fn table_catalog_store() -> S3Result<crate::table_catalog::EcStoreTableCatalogStore<ECStore>> {
fn table_catalog_backend() -> S3Result<crate::table_catalog::EcStoreTableCatalogObjectBackend<ECStore>> {
let store = new_object_layer_fn().ok_or_else(|| s3_error!(InternalError, "object store not initialized"))?;
Ok(crate::table_catalog::ObjectTableCatalogStore::new(
crate::table_catalog::EcStoreTableCatalogObjectBackend::new(store),
))
Ok(crate::table_catalog::EcStoreTableCatalogObjectBackend::new(store))
}
fn table_catalog_store() -> S3Result<crate::table_catalog::EcStoreTableCatalogStore<ECStore>> {
let backend = table_catalog_backend()?;
Ok(crate::table_catalog::ObjectTableCatalogStore::new(backend))
}
async fn table_bucket_enabled_from_metadata(bucket: &str) -> S3Result<bool> {
@@ -375,10 +377,10 @@ fn list_tables_response_from_entries(entries: Vec<crate::table_catalog::TableEnt
Ok(RestListTablesResponse { identifiers })
}
fn load_table_response_from_entry(entry: crate::table_catalog::TableEntry) -> RestLoadTableResponse {
fn load_table_response_from_entry(entry: crate::table_catalog::TableEntry, metadata: serde_json::Value) -> RestLoadTableResponse {
RestLoadTableResponse {
metadata_location: entry.metadata_location,
metadata: serde_json::Value::Null,
metadata,
config: BTreeMap::from([("warehouse-location".to_string(), entry.warehouse_location)]),
}
}
@@ -418,6 +420,9 @@ fn table_entry_from_register_request(
namespace: &crate::table_catalog::Namespace,
request: RegisterTableRequest,
) -> S3Result<crate::table_catalog::TableEntry> {
if request.overwrite {
return Err(s3_error!(NotImplemented, "register table overwrite is not supported"));
}
let table = crate::table_catalog::IdentifierSegment::parse(request.name)
.map_err(|err| s3_error!(InvalidRequest, "invalid table name: {}", err))?;
if !crate::table_catalog::is_valid_table_metadata_location(namespace, &table, &request.metadata_location) {
@@ -439,7 +444,7 @@ fn table_entry_from_register_request(
version_token: format!("token-{}", Uuid::new_v4()),
generation: 1,
state: crate::table_catalog::TableCatalogEntryState::Active,
properties: request.properties,
properties: BTreeMap::new(),
created_at: None,
updated_at: None,
})
@@ -450,15 +455,18 @@ fn table_entry_from_create_table_request(
namespace: &crate::table_catalog::Namespace,
request: CreateTableRequest,
) -> S3Result<crate::table_catalog::TableEntry> {
table_entry_from_register_request(
let properties = request.properties;
let mut entry = table_entry_from_register_request(
bucket,
namespace,
RegisterTableRequest {
name: request.name,
metadata_location: request.metadata_location,
properties: request.properties,
overwrite: false,
},
)
)?;
entry.properties = properties;
Ok(entry)
}
fn namespace_entry_from_create_request(
@@ -545,6 +553,7 @@ where
async fn register_table_response<S>(
store: &S,
metadata_backend: &impl crate::table_catalog::TableCatalogObjectBackend,
bucket: &str,
namespace: &crate::table_catalog::Namespace,
request: RegisterTableRequest,
@@ -555,12 +564,14 @@ where
{
let entry = table_entry_from_register_request(bucket, namespace, request)?;
ensure_table_bucket_entry(store, bucket, table_bucket_enabled).await?;
let metadata = read_table_metadata_json(metadata_backend, bucket, &entry.metadata_location).await?;
store.register_table(entry.clone()).await.map_err(catalog_store_error)?;
Ok(load_table_response_from_entry(entry))
Ok(load_table_response_from_entry(entry, metadata))
}
async fn create_table_response<S>(
store: &S,
metadata_backend: &impl crate::table_catalog::TableCatalogObjectBackend,
bucket: &str,
namespace: &crate::table_catalog::Namespace,
request: CreateTableRequest,
@@ -571,8 +582,29 @@ where
{
let entry = table_entry_from_create_table_request(bucket, namespace, request)?;
ensure_table_bucket_entry(store, bucket, table_bucket_enabled).await?;
let metadata = read_table_metadata_json(metadata_backend, bucket, &entry.metadata_location).await?;
store.create_table(entry.clone()).await.map_err(catalog_store_error)?;
Ok(load_table_response_from_entry(entry))
Ok(load_table_response_from_entry(entry, metadata))
}
async fn read_table_metadata_json(
metadata_backend: &impl crate::table_catalog::TableCatalogObjectBackend,
bucket: &str,
metadata_location: &str,
) -> S3Result<serde_json::Value> {
let Some(object) = metadata_backend
.read_object(bucket, metadata_location)
.await
.map_err(catalog_store_error)?
else {
return Err(s3_error!(InvalidRequest, "table metadata object not found: {metadata_location}"));
};
let metadata = serde_json::from_slice::<serde_json::Value>(&object.data)
.map_err(|err| s3_error!(InvalidRequest, "failed to parse table metadata JSON: {}", err))?;
if !metadata.is_object() {
return Err(s3_error!(InvalidRequest, "table metadata JSON must be an object"));
}
Ok(metadata)
}
async fn list_tables_response<S>(
@@ -592,6 +624,7 @@ where
async fn load_table_response<S>(
store: &S,
metadata_backend: &impl crate::table_catalog::TableCatalogObjectBackend,
bucket: &str,
namespace: &crate::table_catalog::Namespace,
table: &str,
@@ -606,7 +639,8 @@ where
else {
return Err(s3_error!(InvalidRequest, "table not found"));
};
Ok(load_table_response_from_entry(entry))
let metadata = read_table_metadata_json(metadata_backend, bucket, &entry.metadata_location).await?;
Ok(load_table_response_from_entry(entry, metadata))
}
async fn commit_table_response<S>(
@@ -723,9 +757,11 @@ impl Operation for RestCreateTableHandler {
let warehouse = warehouse_from_params(&params)?;
let namespace = namespace_from_params(&params)?;
let request = read_json_body::<CreateTableRequest>(req.input).await?;
let store = table_catalog_store()?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
let table_bucket_enabled = table_bucket_enabled_from_metadata(&warehouse).await?;
let response = create_table_response(&store, &warehouse, &namespace, request, table_bucket_enabled).await?;
let response =
create_table_response(&store, &metadata_backend, &warehouse, &namespace, request, table_bucket_enabled).await?;
build_json_response(StatusCode::OK, &response)
}
}
@@ -739,9 +775,11 @@ impl Operation for RestRegisterTableHandler {
let warehouse = warehouse_from_params(&params)?;
let namespace = namespace_from_params(&params)?;
let request = read_json_body::<RegisterTableRequest>(req.input).await?;
let store = table_catalog_store()?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
let table_bucket_enabled = table_bucket_enabled_from_metadata(&warehouse).await?;
let response = register_table_response(&store, &warehouse, &namespace, request, table_bucket_enabled).await?;
let response =
register_table_response(&store, &metadata_backend, &warehouse, &namespace, request, table_bucket_enabled).await?;
build_json_response(StatusCode::OK, &response)
}
}
@@ -755,8 +793,9 @@ impl Operation for RestLoadTableHandler {
let warehouse = warehouse_from_params(&params)?;
let namespace = namespace_from_params(&params)?;
let table = table_name_from_params(&params)?;
let store = table_catalog_store()?;
let response = load_table_response(&store, &warehouse, &namespace, &table).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?;
build_json_response(StatusCode::OK, &response)
}
}
@@ -796,6 +835,7 @@ impl Operation for RestDropTableHandler {
mod tests {
use super::*;
use crate::table_catalog::TableCatalogStore;
use std::sync::Arc;
#[test]
fn catalog_config_response_lists_mvp_rest_endpoints() {
@@ -810,6 +850,11 @@ mod tests {
.endpoints
.contains(&"POST /{warehouse}/namespaces/{namespace}/register")
);
assert!(
!response
.endpoints
.contains(&"POST /{warehouse}/namespaces/{namespace}/tables")
);
assert!(
response
.endpoints
@@ -928,9 +973,7 @@ mod tests {
let request: RegisterTableRequest = serde_json::from_value(serde_json::json!({
"name": "events",
"metadata-location": ".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00001.metadata.json",
"properties": {
"write.format.default": "parquet"
}
"overwrite": false
}))
.expect("request should parse");
@@ -943,11 +986,44 @@ mod tests {
entry.metadata_location,
".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00001.metadata.json"
);
assert_eq!(entry.properties.get("write.format.default").map(String::as_str), Some("parquet"));
assert!(entry.properties.is_empty());
assert_eq!(entry.generation, 1);
assert!(!entry.version_token.is_empty());
}
#[test]
fn load_table_response_includes_rest_metadata_payload() {
let metadata = serde_json::json!({
"format-version": 2,
"table-uuid": "table-uuid",
"location": "s3://warehouse/tables/table-id"
});
let response = load_table_response_from_entry(
crate::table_catalog::TableEntry {
version: crate::table_catalog::TABLE_CATALOG_ENTRY_VERSION,
table_bucket: "warehouse".to_string(),
namespace: "analytics".to_string(),
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,
},
metadata.clone(),
);
assert_eq!(response.metadata, metadata);
}
#[test]
fn commit_table_request_uses_rest_commit_fields() {
let request: RestCommitTableRequest = serde_json::from_value(serde_json::json!({
@@ -987,6 +1063,89 @@ mod tests {
commits: tokio::sync::Mutex<Vec<crate::table_catalog::CommitLogEntry>>,
}
#[derive(Clone, Default)]
struct TestTableCatalogObjectBackend {
objects: Arc<tokio::sync::Mutex<BTreeMap<(String, String), crate::table_catalog::TableCatalogObject>>>,
}
impl TestTableCatalogObjectBackend {
async fn put_json(&self, bucket: &str, object: &str, value: serde_json::Value) {
let data = serde_json::to_vec(&value).expect("metadata JSON should serialize");
self.objects.lock().await.insert(
(bucket.to_string(), object.to_string()),
crate::table_catalog::TableCatalogObject {
data,
etag: Some("etag".to_string()),
},
);
}
}
#[async_trait::async_trait]
impl crate::table_catalog::TableCatalogObjectBackend for TestTableCatalogObjectBackend {
async fn read_object(
&self,
bucket: &str,
object: &str,
) -> crate::table_catalog::TableCatalogStoreResult<Option<crate::table_catalog::TableCatalogObject>> {
Ok(self
.objects
.lock()
.await
.get(&(bucket.to_string(), object.to_string()))
.cloned())
}
async fn object_exists(&self, bucket: &str, object: &str) -> crate::table_catalog::TableCatalogStoreResult<bool> {
Ok(self
.objects
.lock()
.await
.contains_key(&(bucket.to_string(), object.to_string())))
}
async fn put_object(
&self,
bucket: &str,
object: &str,
data: Vec<u8>,
_precondition: crate::table_catalog::TableCatalogPutPrecondition,
) -> crate::table_catalog::TableCatalogStoreResult<()> {
self.objects.lock().await.insert(
(bucket.to_string(), object.to_string()),
crate::table_catalog::TableCatalogObject {
data,
etag: Some("etag".to_string()),
},
);
Ok(())
}
async fn delete_object(&self, bucket: &str, object: &str) -> crate::table_catalog::TableCatalogStoreResult<()> {
self.objects.lock().await.remove(&(bucket.to_string(), object.to_string()));
Ok(())
}
async fn list_objects(&self, bucket: &str, prefix: &str) -> crate::table_catalog::TableCatalogStoreResult<Vec<String>> {
Ok(self
.objects
.lock()
.await
.keys()
.filter(|(object_bucket, object)| object_bucket == bucket && object.starts_with(prefix))
.map(|(_, object)| object.clone())
.collect())
}
async fn acquire_write_lock(
&self,
_bucket: &str,
_object: &str,
) -> crate::table_catalog::TableCatalogStoreResult<Box<dyn Send>> {
Ok(Box::new(()))
}
}
#[async_trait::async_trait]
impl crate::table_catalog::TableCatalogStore for TestTableCatalogStore {
async fn get_table_bucket(
@@ -1288,6 +1447,7 @@ mod tests {
#[tokio::test]
async fn table_helpers_call_catalog_store() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
ensure_table_bucket_entry(&store, "warehouse", true)
.await
@@ -1306,14 +1466,26 @@ mod tests {
let metadata_location =
".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00001.metadata.json";
metadata_backend
.put_json(
"warehouse",
metadata_location,
serde_json::json!({
"format-version": 2,
"table-uuid": "table-uuid",
"location": "s3://warehouse/tables/table-id"
}),
)
.await;
let register = register_table_response(
&store,
&metadata_backend,
"warehouse",
&namespace,
RegisterTableRequest {
name: "events".to_string(),
metadata_location: metadata_location.to_string(),
properties: BTreeMap::new(),
overwrite: false,
},
true,
)
@@ -1321,16 +1493,18 @@ mod tests {
.expect("table should register");
assert_eq!(register.metadata_location, metadata_location);
assert_eq!(register.metadata["format-version"], 2);
let list = list_tables_response(&store, "warehouse", &namespace)
.await
.expect("table list should load");
assert_eq!(list.identifiers[0].name, "events");
let load = load_table_response(&store, "warehouse", &namespace, "events")
let load = load_table_response(&store, &metadata_backend, "warehouse", &namespace, "events")
.await
.expect("table should load");
assert_eq!(load.metadata_location, metadata_location);
assert_eq!(load.metadata["table-uuid"], "table-uuid");
let current = store
.load_table("warehouse", "analytics", "events")
@@ -1339,6 +1513,18 @@ mod tests {
.expect("table should exist");
let next_metadata_location =
".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00002.metadata.json";
metadata_backend
.put_json(
"warehouse",
next_metadata_location,
serde_json::json!({
"format-version": 2,
"table-uuid": "table-uuid",
"location": "s3://warehouse/tables/table-id",
"last-sequence-number": 2
}),
)
.await;
let commit = commit_table_response(
&store,
"warehouse",
@@ -1372,6 +1558,10 @@ mod tests {
drop_table_in_store(&store, "warehouse", &namespace, "events")
.await
.expect("table should drop");
assert!(load_table_response(&store, "warehouse", &namespace, "events").await.is_err());
assert!(
load_table_response(&store, &metadata_backend, "warehouse", &namespace, "events")
.await
.is_err()
);
}
}