diff --git a/rustfs/src/admin/handlers/table_catalog.rs b/rustfs/src/admin/handlers/table_catalog.rs index e546c2349..c656a2f32 100644 --- a/rustfs/src/admin/handlers/table_catalog.rs +++ b/rustfs/src/admin/handlers/table_catalog.rs @@ -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 { + 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( 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( @@ -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); + } }