mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 08:18:18 +00:00
fix(table-catalog): validate committed metadata identity (#3327)
* fix(table-catalog): validate committed metadata identity * fix(table-catalog): preserve legacy metadata identity commits --------- Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
@@ -765,14 +765,63 @@ fn validate_table_location_in_bucket(bucket: &str, location: &str) -> S3Result<(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn validate_metadata_table_location_in_bucket(bucket: &str, metadata: &serde_json::Value) -> S3Result<()> {
|
||||
let location = metadata
|
||||
fn metadata_table_uuid(metadata: &serde_json::Value) -> S3Result<&str> {
|
||||
metadata
|
||||
.get("table-uuid")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.filter(|uuid| !uuid.is_empty())
|
||||
.ok_or_else(|| s3_error!(InvalidRequest, "table metadata is missing table-uuid"))
|
||||
}
|
||||
|
||||
fn metadata_format_version(metadata: &serde_json::Value) -> S3Result<u16> {
|
||||
let version = metadata
|
||||
.get("format-version")
|
||||
.and_then(serde_json::Value::as_u64)
|
||||
.filter(|version| *version > 0)
|
||||
.ok_or_else(|| s3_error!(InvalidRequest, "table metadata is missing format-version"))?;
|
||||
u16::try_from(version).map_err(|_| s3_error!(InvalidRequest, "table metadata format-version is too large"))
|
||||
}
|
||||
|
||||
fn metadata_table_location(metadata: &serde_json::Value) -> S3Result<&str> {
|
||||
metadata
|
||||
.get("location")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.ok_or_else(|| s3_error!(InvalidRequest, "table metadata is missing location"))?;
|
||||
.filter(|location| !location.is_empty())
|
||||
.ok_or_else(|| s3_error!(InvalidRequest, "table metadata is missing location"))
|
||||
}
|
||||
|
||||
fn validate_metadata_table_location_in_bucket(bucket: &str, metadata: &serde_json::Value) -> S3Result<()> {
|
||||
let location = metadata_table_location(metadata)?;
|
||||
validate_table_location_in_bucket(bucket, location)
|
||||
}
|
||||
|
||||
fn validate_metadata_matches_current_metadata(
|
||||
current_metadata: &serde_json::Value,
|
||||
target_metadata: &serde_json::Value,
|
||||
) -> S3Result<()> {
|
||||
let expected_table_uuid = metadata_table_uuid(current_metadata)?;
|
||||
metadata_format_version(current_metadata)?;
|
||||
let target_table_uuid = metadata_table_uuid(target_metadata)?;
|
||||
metadata_format_version(target_metadata)?;
|
||||
if target_table_uuid != expected_table_uuid {
|
||||
return Err(s3_error!(
|
||||
InvalidRequest,
|
||||
"table metadata table-uuid does not match current table metadata"
|
||||
));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn adopt_registered_metadata_identity(
|
||||
entry: &mut crate::table_catalog::TableEntry,
|
||||
metadata: &serde_json::Value,
|
||||
) -> S3Result<()> {
|
||||
entry.table_uuid = metadata_table_uuid(metadata)?.to_string();
|
||||
entry.format_version = metadata_format_version(metadata)?;
|
||||
entry.warehouse_location = metadata_table_location(metadata)?.to_string();
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn table_entry_from_register_request(
|
||||
bucket: &str,
|
||||
namespace: &crate::table_catalog::Namespace,
|
||||
@@ -1608,9 +1657,11 @@ async fn register_table_response<S>(
|
||||
where
|
||||
S: crate::table_catalog::TableCatalogStore + ?Sized,
|
||||
{
|
||||
let entry = table_entry_from_register_request(bucket, namespace, request)?;
|
||||
let mut 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?;
|
||||
validate_metadata_table_location_in_bucket(bucket, &metadata)?;
|
||||
adopt_registered_metadata_identity(&mut entry, &metadata)?;
|
||||
store.register_table(entry.clone()).await.map_err(catalog_store_error)?;
|
||||
Ok(load_table_response_from_entry(entry, metadata))
|
||||
}
|
||||
@@ -1741,8 +1792,11 @@ where
|
||||
if !crate::table_catalog::is_valid_table_metadata_location(namespace, &table_name, &request.metadata_location) {
|
||||
return Err(s3_error!(InvalidRequest, "metadata location must be inside the table metadata directory"));
|
||||
}
|
||||
let current_metadata = read_table_metadata_json(metadata_backend, bucket, ¤t.metadata_location).await?;
|
||||
validate_metadata_table_location_in_bucket(bucket, ¤t_metadata)?;
|
||||
let target_metadata = read_table_metadata_json(metadata_backend, bucket, &request.metadata_location).await?;
|
||||
validate_metadata_table_location_in_bucket(bucket, &target_metadata)?;
|
||||
validate_metadata_matches_current_metadata(¤t_metadata, &target_metadata)?;
|
||||
let commit_request = crate::table_catalog::TableCommitRequest {
|
||||
table_bucket: bucket.to_string(),
|
||||
namespace: namespace.public_name(),
|
||||
@@ -1776,9 +1830,25 @@ where
|
||||
}
|
||||
|
||||
let request = table_commit_request_from_rest_request(bucket, namespace, table, request)?;
|
||||
let Some(current) = store
|
||||
.load_table(bucket, &namespace.public_name(), table)
|
||||
.await
|
||||
.map_err(catalog_store_error)?
|
||||
else {
|
||||
return Err(s3_error!(InvalidRequest, "table not found"));
|
||||
};
|
||||
let table_name = crate::table_catalog::IdentifierSegment::parse(table.to_string())
|
||||
.map_err(|err| s3_error!(InvalidRequest, "invalid table name: {}", err))?;
|
||||
if !crate::table_catalog::is_valid_table_metadata_location(namespace, &table_name, &request.new_metadata_location) {
|
||||
return Err(s3_error!(InvalidRequest, "metadata location must be inside the table metadata directory"));
|
||||
}
|
||||
let current_metadata = read_table_metadata_json(metadata_backend, bucket, ¤t.metadata_location).await?;
|
||||
validate_metadata_table_location_in_bucket(bucket, ¤t_metadata)?;
|
||||
let target_metadata = read_table_metadata_json(metadata_backend, bucket, &request.new_metadata_location).await?;
|
||||
validate_metadata_table_location_in_bucket(bucket, &target_metadata)?;
|
||||
validate_metadata_matches_current_metadata(¤t_metadata, &target_metadata)?;
|
||||
let result = store.commit_table(request).await.map_err(catalog_store_error)?;
|
||||
let metadata = read_table_metadata_json(metadata_backend, bucket, &result.table.metadata_location).await?;
|
||||
Ok(commit_table_response_from_result(result, metadata))
|
||||
Ok(commit_table_response_from_result(result, target_metadata))
|
||||
}
|
||||
|
||||
async fn standard_commit_table_response<S>(
|
||||
@@ -1801,8 +1871,10 @@ where
|
||||
};
|
||||
let current_metadata = read_table_metadata_json(metadata_backend, bucket, ¤t.metadata_location).await?;
|
||||
validate_table_commit_requirements(¤t_metadata, &request.requirements)?;
|
||||
let expected_metadata = current_metadata.clone();
|
||||
let next_metadata = apply_table_commit_updates(current_metadata, &request.updates, ¤t.metadata_location)?;
|
||||
validate_metadata_table_location_in_bucket(bucket, &next_metadata)?;
|
||||
validate_metadata_matches_current_metadata(&expected_metadata, &next_metadata)?;
|
||||
let table_name = crate::table_catalog::IdentifierSegment::parse(table.to_string())
|
||||
.map_err(|err| s3_error!(InvalidRequest, "invalid table name: {}", err))?;
|
||||
let next_generation = current.generation.saturating_add(1);
|
||||
@@ -1893,9 +1965,10 @@ where
|
||||
B: crate::table_catalog::TableCatalogObjectBackend,
|
||||
{
|
||||
ensure_table_bucket_entry(store, bucket, table_bucket_enabled).await?;
|
||||
let entry = table_entry_from_import_request(bucket, namespace, table, request)?;
|
||||
let mut entry = table_entry_from_import_request(bucket, namespace, table, request)?;
|
||||
let metadata = read_table_metadata_json(metadata_backend, bucket, &entry.metadata_location).await?;
|
||||
validate_metadata_table_location_in_bucket(bucket, &metadata)?;
|
||||
adopt_registered_metadata_identity(&mut entry, &metadata)?;
|
||||
store.register_table(entry.clone()).await.map_err(catalog_store_error)?;
|
||||
Ok(load_table_response_from_entry(entry, metadata))
|
||||
}
|
||||
@@ -1923,8 +1996,11 @@ where
|
||||
if !crate::table_catalog::is_valid_table_metadata_location(namespace, &table_name, &request.metadata_location) {
|
||||
return Err(s3_error!(InvalidRequest, "metadata location must be inside the table metadata directory"));
|
||||
}
|
||||
let current_metadata = read_table_metadata_json(metadata_backend, bucket, ¤t.metadata_location).await?;
|
||||
validate_metadata_table_location_in_bucket(bucket, ¤t_metadata)?;
|
||||
let target_metadata = read_table_metadata_json(metadata_backend, bucket, &request.metadata_location).await?;
|
||||
validate_metadata_table_location_in_bucket(bucket, &target_metadata)?;
|
||||
validate_metadata_matches_current_metadata(¤t_metadata, &target_metadata)?;
|
||||
let commit_request = crate::table_catalog::TableCommitRequest {
|
||||
table_bucket: bucket.to_string(),
|
||||
namespace: namespace.public_name(),
|
||||
@@ -2990,6 +3066,154 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn standard_commit_accepts_legacy_catalog_uuid_when_current_metadata_matches() {
|
||||
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
|
||||
.expect("table bucket entry should be seeded");
|
||||
create_namespace_response(
|
||||
&store,
|
||||
"warehouse",
|
||||
CreateNamespaceRequest {
|
||||
namespace: vec!["analytics".to_string()],
|
||||
properties: BTreeMap::new(),
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect("namespace should be created");
|
||||
let current_location = ".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00001.metadata.json";
|
||||
let legacy_entry = table_entry_from_register_request(
|
||||
"warehouse",
|
||||
&namespace,
|
||||
RegisterTableRequest {
|
||||
name: "events".to_string(),
|
||||
metadata_location: current_location.to_string(),
|
||||
overwrite: false,
|
||||
},
|
||||
)
|
||||
.expect("table entry should build");
|
||||
assert_ne!(legacy_entry.table_uuid, "metadata-table-uuid");
|
||||
store
|
||||
.register_table(legacy_entry.clone())
|
||||
.await
|
||||
.expect("legacy table entry should register");
|
||||
metadata_backend
|
||||
.put_json(
|
||||
"warehouse",
|
||||
current_location,
|
||||
serde_json::json!({
|
||||
"format-version": 2,
|
||||
"table-uuid": "metadata-table-uuid",
|
||||
"location": "s3://warehouse/tables/table-id",
|
||||
"properties": {}
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
|
||||
let commit_request: RestCommitTableRequest = serde_json::from_value(serde_json::json!({
|
||||
"updates": [
|
||||
{
|
||||
"action": "set-properties",
|
||||
"updates": {
|
||||
"owner": "lakehouse"
|
||||
}
|
||||
}
|
||||
]
|
||||
}))
|
||||
.expect("standard commit table request should parse");
|
||||
let committed = commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", commit_request)
|
||||
.await
|
||||
.expect("legacy catalog uuid should not block standard commit");
|
||||
|
||||
assert_eq!(committed.metadata["table-uuid"], "metadata-table-uuid");
|
||||
assert_eq!(committed.metadata["properties"]["owner"], "lakehouse");
|
||||
assert_eq!(committed.generation, legacy_entry.generation + 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn metadata_location_api_accepts_legacy_catalog_uuid_when_target_matches_current_metadata() {
|
||||
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
|
||||
.expect("table bucket entry should be seeded");
|
||||
create_namespace_response(
|
||||
&store,
|
||||
"warehouse",
|
||||
CreateNamespaceRequest {
|
||||
namespace: vec!["analytics".to_string()],
|
||||
properties: BTreeMap::new(),
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect("namespace should be created");
|
||||
let current_location = ".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00001.metadata.json";
|
||||
let legacy_entry = table_entry_from_register_request(
|
||||
"warehouse",
|
||||
&namespace,
|
||||
RegisterTableRequest {
|
||||
name: "events".to_string(),
|
||||
metadata_location: current_location.to_string(),
|
||||
overwrite: false,
|
||||
},
|
||||
)
|
||||
.expect("table entry should build");
|
||||
assert_ne!(legacy_entry.table_uuid, "metadata-table-uuid");
|
||||
store
|
||||
.register_table(legacy_entry.clone())
|
||||
.await
|
||||
.expect("legacy table entry should register");
|
||||
metadata_backend
|
||||
.put_json(
|
||||
"warehouse",
|
||||
current_location,
|
||||
serde_json::json!({
|
||||
"format-version": 2,
|
||||
"table-uuid": "metadata-table-uuid",
|
||||
"location": "s3://warehouse/tables/table-id"
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
let next_location = ".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00002.metadata.json";
|
||||
metadata_backend
|
||||
.put_json(
|
||||
"warehouse",
|
||||
next_location,
|
||||
serde_json::json!({
|
||||
"format-version": 2,
|
||||
"table-uuid": "metadata-table-uuid",
|
||||
"location": "s3://warehouse/tables/table-id",
|
||||
"last-sequence-number": 2
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
|
||||
let updated = update_table_metadata_location_response(
|
||||
&store,
|
||||
&metadata_backend,
|
||||
"warehouse",
|
||||
&namespace,
|
||||
"events",
|
||||
UpdateTableMetadataLocationRequest {
|
||||
metadata_location: next_location.to_string(),
|
||||
version_token: legacy_entry.version_token,
|
||||
commit_id: Some("commit-1".to_string()),
|
||||
idempotency_key: None,
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("legacy catalog uuid should not block metadata-location update");
|
||||
|
||||
assert_eq!(updated.metadata_location, next_location);
|
||||
assert_eq!(updated.generation, legacy_entry.generation + 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn table_metadata_maintenance_helper_runs_dry_run_and_delete() {
|
||||
let backend = TestTableCatalogObjectBackend::default();
|
||||
@@ -3823,6 +4047,119 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn register_table_response_adopts_metadata_table_uuid() {
|
||||
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
|
||||
.expect("table bucket entry should be seeded");
|
||||
create_namespace_response(
|
||||
&store,
|
||||
"warehouse",
|
||||
CreateNamespaceRequest {
|
||||
namespace: vec!["analytics".to_string()],
|
||||
properties: BTreeMap::new(),
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect("namespace should be created");
|
||||
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": "metadata-table-uuid",
|
||||
"location": "s3://warehouse/tables/table-id"
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
|
||||
register_table_response(
|
||||
&store,
|
||||
&metadata_backend,
|
||||
"warehouse",
|
||||
&namespace,
|
||||
RegisterTableRequest {
|
||||
name: "events".to_string(),
|
||||
metadata_location: metadata_location.to_string(),
|
||||
overwrite: false,
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect("table should register");
|
||||
|
||||
let entry = store
|
||||
.load_table("warehouse", "analytics", "events")
|
||||
.await
|
||||
.expect("table lookup should succeed")
|
||||
.expect("table should exist");
|
||||
assert_eq!(entry.table_uuid, "metadata-table-uuid");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn register_table_response_rejects_metadata_without_format_version() {
|
||||
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
|
||||
.expect("table bucket entry should be seeded");
|
||||
create_namespace_response(
|
||||
&store,
|
||||
"warehouse",
|
||||
CreateNamespaceRequest {
|
||||
namespace: vec!["analytics".to_string()],
|
||||
properties: BTreeMap::new(),
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect("namespace should be created");
|
||||
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!({
|
||||
"table-uuid": "metadata-table-uuid",
|
||||
"location": "s3://warehouse/tables/table-id"
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(
|
||||
register_table_response(
|
||||
&store,
|
||||
&metadata_backend,
|
||||
"warehouse",
|
||||
&namespace,
|
||||
RegisterTableRequest {
|
||||
name: "events".to_string(),
|
||||
metadata_location: metadata_location.to_string(),
|
||||
overwrite: false,
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.is_err()
|
||||
);
|
||||
assert!(
|
||||
store
|
||||
.load_table("warehouse", "analytics", "events")
|
||||
.await
|
||||
.expect("table lookup should succeed")
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn metadata_location_api_loads_and_updates_current_pointer() {
|
||||
let store = TestTableCatalogStore::default();
|
||||
@@ -3843,21 +4180,29 @@ mod tests {
|
||||
.await
|
||||
.expect("namespace should be created");
|
||||
let current_location = ".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00001.metadata.json";
|
||||
store
|
||||
.register_table(
|
||||
table_entry_from_register_request(
|
||||
"warehouse",
|
||||
&namespace,
|
||||
RegisterTableRequest {
|
||||
name: "events".to_string(),
|
||||
metadata_location: current_location.to_string(),
|
||||
overwrite: false,
|
||||
},
|
||||
)
|
||||
.expect("table entry should build"),
|
||||
let entry = table_entry_from_register_request(
|
||||
"warehouse",
|
||||
&namespace,
|
||||
RegisterTableRequest {
|
||||
name: "events".to_string(),
|
||||
metadata_location: current_location.to_string(),
|
||||
overwrite: false,
|
||||
},
|
||||
)
|
||||
.expect("table entry should build");
|
||||
let table_uuid = entry.table_uuid.clone();
|
||||
store.register_table(entry).await.expect("table should register");
|
||||
metadata_backend
|
||||
.put_json(
|
||||
"warehouse",
|
||||
current_location,
|
||||
serde_json::json!({
|
||||
"format-version": 2,
|
||||
"table-uuid": table_uuid,
|
||||
"location": "s3://warehouse/tables/table-id"
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.expect("table should register");
|
||||
.await;
|
||||
let current = get_table_metadata_location_response(&store, "warehouse", &namespace, "events")
|
||||
.await
|
||||
.expect("metadata location should load");
|
||||
@@ -3868,7 +4213,7 @@ mod tests {
|
||||
next_location,
|
||||
serde_json::json!({
|
||||
"format-version": 2,
|
||||
"table-uuid": "table-uuid",
|
||||
"table-uuid": table_uuid,
|
||||
"location": "s3://warehouse/tables/table-id"
|
||||
}),
|
||||
)
|
||||
@@ -3970,6 +4315,92 @@ mod tests {
|
||||
assert_eq!(unchanged.generation, current.generation);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn metadata_location_api_rejects_mismatched_table_uuid_before_commit() {
|
||||
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
|
||||
.expect("table bucket entry should be seeded");
|
||||
create_namespace_response(
|
||||
&store,
|
||||
"warehouse",
|
||||
CreateNamespaceRequest {
|
||||
namespace: vec!["analytics".to_string()],
|
||||
properties: BTreeMap::new(),
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect("namespace should be created");
|
||||
let current_location = ".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00001.metadata.json";
|
||||
metadata_backend
|
||||
.put_json(
|
||||
"warehouse",
|
||||
current_location,
|
||||
serde_json::json!({
|
||||
"format-version": 2,
|
||||
"table-uuid": "table-uuid",
|
||||
"location": "s3://warehouse/tables/table-id"
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
register_table_response(
|
||||
&store,
|
||||
&metadata_backend,
|
||||
"warehouse",
|
||||
&namespace,
|
||||
RegisterTableRequest {
|
||||
name: "events".to_string(),
|
||||
metadata_location: current_location.to_string(),
|
||||
overwrite: false,
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect("table should register");
|
||||
let current = get_table_metadata_location_response(&store, "warehouse", &namespace, "events")
|
||||
.await
|
||||
.expect("metadata location should load");
|
||||
let mismatched_location =
|
||||
".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00002.metadata.json";
|
||||
metadata_backend
|
||||
.put_json(
|
||||
"warehouse",
|
||||
mismatched_location,
|
||||
serde_json::json!({
|
||||
"format-version": 2,
|
||||
"table-uuid": "other-table-uuid",
|
||||
"location": "s3://warehouse/tables/table-id"
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(
|
||||
update_table_metadata_location_response(
|
||||
&store,
|
||||
&metadata_backend,
|
||||
"warehouse",
|
||||
&namespace,
|
||||
"events",
|
||||
UpdateTableMetadataLocationRequest {
|
||||
metadata_location: mismatched_location.to_string(),
|
||||
version_token: current.version_token,
|
||||
commit_id: Some("commit-1".to_string()),
|
||||
idempotency_key: None,
|
||||
},
|
||||
)
|
||||
.await
|
||||
.is_err()
|
||||
);
|
||||
let unchanged = get_table_metadata_location_response(&store, "warehouse", &namespace, "events")
|
||||
.await
|
||||
.expect("metadata location should still load");
|
||||
assert_eq!(unchanged.metadata_location, current_location);
|
||||
assert_eq!(unchanged.generation, current.generation);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn catalog_import_and_rollback_use_register_and_commit_paths() {
|
||||
let backend = TestTableCatalogObjectBackend::default();
|
||||
@@ -4151,4 +4582,192 @@ mod tests {
|
||||
assert_eq!(unchanged.metadata_location, current_location);
|
||||
assert_eq!(unchanged.generation, current.generation);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rollback_rejects_mismatched_table_uuid_before_commit() {
|
||||
let backend = TestTableCatalogObjectBackend::default();
|
||||
let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone());
|
||||
let bucket = "warehouse";
|
||||
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
|
||||
let table = crate::table_catalog::IdentifierSegment::parse("events").expect("table should parse");
|
||||
ensure_table_bucket_entry(&store, bucket, true)
|
||||
.await
|
||||
.expect("table bucket entry should be seeded");
|
||||
create_namespace_response(
|
||||
&store,
|
||||
bucket,
|
||||
CreateNamespaceRequest {
|
||||
namespace: vec!["analytics".to_string()],
|
||||
properties: BTreeMap::new(),
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect("namespace should be created");
|
||||
let current_location = crate::table_catalog::default_table_metadata_file_path(&namespace, &table, "00001.metadata.json");
|
||||
backend
|
||||
.put_json(
|
||||
bucket,
|
||||
¤t_location,
|
||||
serde_json::json!({
|
||||
"format-version": 2,
|
||||
"table-uuid": "table-uuid",
|
||||
"location": "s3://warehouse/tables/table-id"
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
catalog_import_response(
|
||||
&store,
|
||||
&backend,
|
||||
bucket,
|
||||
&namespace,
|
||||
"events",
|
||||
CatalogImportRequest {
|
||||
metadata_location: current_location.clone(),
|
||||
properties: BTreeMap::new(),
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect("catalog import should register table");
|
||||
let current = store
|
||||
.load_table(bucket, "analytics", "events")
|
||||
.await
|
||||
.expect("table lookup should succeed")
|
||||
.expect("table should exist");
|
||||
|
||||
let mismatched_location =
|
||||
crate::table_catalog::default_table_metadata_file_path(&namespace, &table, "00002.metadata.json");
|
||||
backend
|
||||
.put_json(
|
||||
bucket,
|
||||
&mismatched_location,
|
||||
serde_json::json!({
|
||||
"format-version": 2,
|
||||
"table-uuid": "other-table-uuid",
|
||||
"location": "s3://warehouse/tables/table-id"
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(
|
||||
rollback_table_response(
|
||||
&store,
|
||||
&backend,
|
||||
bucket,
|
||||
&namespace,
|
||||
"events",
|
||||
RollbackTableRequest {
|
||||
metadata_location: mismatched_location,
|
||||
version_token: current.version_token,
|
||||
commit_id: Some("rollback-1".to_string()),
|
||||
idempotency_key: None,
|
||||
},
|
||||
)
|
||||
.await
|
||||
.is_err()
|
||||
);
|
||||
let unchanged = store
|
||||
.load_table(bucket, "analytics", "events")
|
||||
.await
|
||||
.expect("table lookup should succeed")
|
||||
.expect("table should still exist");
|
||||
assert_eq!(unchanged.metadata_location, current_location);
|
||||
assert_eq!(unchanged.generation, current.generation);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn legacy_commit_rejects_mismatched_table_uuid_before_commit() {
|
||||
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
|
||||
.expect("table bucket entry should be seeded");
|
||||
create_namespace_response(
|
||||
&store,
|
||||
"warehouse",
|
||||
CreateNamespaceRequest {
|
||||
namespace: vec!["analytics".to_string()],
|
||||
properties: BTreeMap::new(),
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect("namespace should be created");
|
||||
let current_location = ".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00001.metadata.json";
|
||||
metadata_backend
|
||||
.put_json(
|
||||
"warehouse",
|
||||
current_location,
|
||||
serde_json::json!({
|
||||
"format-version": 2,
|
||||
"table-uuid": "table-uuid",
|
||||
"location": "s3://warehouse/tables/table-id"
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
register_table_response(
|
||||
&store,
|
||||
&metadata_backend,
|
||||
"warehouse",
|
||||
&namespace,
|
||||
RegisterTableRequest {
|
||||
name: "events".to_string(),
|
||||
metadata_location: current_location.to_string(),
|
||||
overwrite: false,
|
||||
},
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.expect("table should register");
|
||||
let current = store
|
||||
.load_table("warehouse", "analytics", "events")
|
||||
.await
|
||||
.expect("table lookup should succeed")
|
||||
.expect("table should exist");
|
||||
let mismatched_location =
|
||||
".rustfs-table/warehouses/default/namespaces/analytics/tables/events/metadata/00002.metadata.json";
|
||||
metadata_backend
|
||||
.put_json(
|
||||
"warehouse",
|
||||
mismatched_location,
|
||||
serde_json::json!({
|
||||
"format-version": 2,
|
||||
"table-uuid": "other-table-uuid",
|
||||
"location": "s3://warehouse/tables/table-id"
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(
|
||||
commit_table_response(
|
||||
&store,
|
||||
&metadata_backend,
|
||||
"warehouse",
|
||||
&namespace,
|
||||
"events",
|
||||
RestCommitTableRequest {
|
||||
commit_id: Some("commit-1".to_string()),
|
||||
idempotency_key: None,
|
||||
operation: Some("append".to_string()),
|
||||
expected_version_token: Some(current.version_token.clone()),
|
||||
expected_metadata_location: Some(current.metadata_location.clone()),
|
||||
new_metadata_location: Some(mismatched_location.to_string()),
|
||||
requirements: Vec::new(),
|
||||
updates: Vec::new(),
|
||||
writer: Some("pyiceberg".to_string()),
|
||||
},
|
||||
)
|
||||
.await
|
||||
.is_err()
|
||||
);
|
||||
let unchanged = store
|
||||
.load_table("warehouse", "analytics", "events")
|
||||
.await
|
||||
.expect("table lookup should succeed")
|
||||
.expect("table should still exist");
|
||||
assert_eq!(unchanged.metadata_location, current_location);
|
||||
assert_eq!(unchanged.generation, current.generation);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user