From 7ec16e197c1d027a40f8fadf683387a89f2ac44f Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Tue, 9 Jun 2026 22:08:28 +0800 Subject: [PATCH] feat(table-catalog): add metadata maintenance control plane (#3302) Co-authored-by: Henry Guo --- rustfs/src/admin/handlers/table_catalog.rs | 224 +++++++++++++++++++- rustfs/src/admin/route_policy.rs | 14 +- rustfs/src/admin/route_registration_test.rs | 10 + rustfs/src/table_catalog.rs | 84 +++++++- 4 files changed, 327 insertions(+), 5 deletions(-) diff --git a/rustfs/src/admin/handlers/table_catalog.rs b/rustfs/src/admin/handlers/table_catalog.rs index 77e883bf2..3ffb1576c 100644 --- a/rustfs/src/admin/handlers/table_catalog.rs +++ b/rustfs/src/admin/handlers/table_catalog.rs @@ -57,6 +57,7 @@ static REGISTER_TABLE_HANDLER: RestRegisterTableHandler = RestRegisterTableHandl static LOAD_TABLE_HANDLER: RestLoadTableHandler = RestLoadTableHandler {}; static COMMIT_TABLE_HANDLER: RestCommitTableHandler = RestCommitTableHandler {}; static DROP_TABLE_HANDLER: RestDropTableHandler = RestDropTableHandler {}; +static TABLE_METADATA_MAINTENANCE_HANDLER: RestTableMetadataMaintenanceHandler = RestTableMetadataMaintenanceHandler {}; #[derive(Debug, Serialize)] struct CatalogConfigResponse { @@ -123,6 +124,15 @@ struct RestCommitTableRequest { writer: Option, } +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct TableMetadataMaintenanceRequest { + #[serde(default, rename = "retain-recent-metadata-files")] + retain_recent_metadata_files: usize, + #[serde(default)] + delete: bool, +} + #[derive(Debug, Serialize)] struct RestNamespaceResponse { namespace: Vec, @@ -221,6 +231,11 @@ pub fn register_table_catalog_route(r: &mut S3Router) -> std::io format!("{TABLE_CATALOG_PREFIX}/{{warehouse}}/namespaces/{{namespace}}/tables/{{table}}").as_str(), AdminOperation(&DROP_TABLE_HANDLER), )?; + r.insert( + Method::POST, + format!("{TABLE_CATALOG_PREFIX}/{{warehouse}}/namespaces/{{namespace}}/tables/{{table}}/maintenance/metadata").as_str(), + AdminOperation(&TABLE_METADATA_MAINTENANCE_HANDLER), + )?; Ok(()) } @@ -1454,6 +1469,34 @@ where .map_err(catalog_store_error) } +async fn table_metadata_maintenance_response( + store: &crate::table_catalog::ObjectTableCatalogStore, + bucket: &str, + namespace: &crate::table_catalog::Namespace, + table: &str, + request: TableMetadataMaintenanceRequest, +) -> S3Result +where + B: crate::table_catalog::TableCatalogObjectBackend, +{ + if request.delete { + store + .delete_table_metadata_maintenance_candidates( + bucket, + &namespace.public_name(), + table, + request.retain_recent_metadata_files, + ) + .await + .map_err(catalog_store_error) + } else { + store + .plan_table_metadata_maintenance(bucket, &namespace.public_name(), table, request.retain_recent_metadata_files) + .await + .map_err(catalog_store_error) + } +} + pub struct GetCatalogConfigHandler {} #[async_trait::async_trait] @@ -1618,6 +1661,23 @@ impl Operation for RestDropTableHandler { } } +pub struct RestTableMetadataMaintenanceHandler {} + +#[async_trait::async_trait] +impl Operation for RestTableMetadataMaintenanceHandler { + async fn call(&self, req: S3Request, params: Params<'_, '_>) -> S3Result> { + let warehouse = warehouse_from_params(¶ms)?; + authorize_table_catalog_warehouse_request(&req, &warehouse, AdminAction::RunTableMaintenanceAction).await?; + let namespace = namespace_from_params(¶ms)?; + let table = table_name_from_params(¶ms)?; + let request = read_json_body::(req.input).await?; + let metadata_backend = table_catalog_backend()?; + let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend); + let response = table_metadata_maintenance_response(&store, &warehouse, &namespace, &table, request).await?; + build_json_response(StatusCode::OK, &response) + } +} + #[cfg(test)] mod tests { use super::*; @@ -1674,6 +1734,7 @@ mod tests { ("RestLoadTableHandler", "AdminAction::GetTableMetadataAction"), ("RestCommitTableHandler", "AdminAction::CommitTableAction"), ("RestDropTableHandler", "AdminAction::DeleteTableAction"), + ("RestTableMetadataMaintenanceHandler", "AdminAction::RunTableMaintenanceAction"), ] { let block = operation_block(src, handler); assert!( @@ -1711,6 +1772,7 @@ mod tests { let _: &RestLoadTableHandler = &LOAD_TABLE_HANDLER; let _: &RestCommitTableHandler = &COMMIT_TABLE_HANDLER; let _: &RestDropTableHandler = &DROP_TABLE_HANDLER; + let _: &RestTableMetadataMaintenanceHandler = &TABLE_METADATA_MAINTENANCE_HANDLER; assert_operation::(); assert_operation::(); @@ -1722,6 +1784,28 @@ mod tests { assert_operation::(); assert_operation::(); assert_operation::(); + assert_operation::(); + } + + #[test] + fn table_metadata_maintenance_request_uses_conservative_defaults() { + let request: TableMetadataMaintenanceRequest = + serde_json::from_value(serde_json::json!({})).expect("default maintenance request should parse"); + + assert_eq!(request.retain_recent_metadata_files, 0); + assert!(!request.delete); + } + + #[test] + fn table_metadata_maintenance_request_accepts_delete_mode() { + let request: TableMetadataMaintenanceRequest = serde_json::from_value(serde_json::json!({ + "retain-recent-metadata-files": 2, + "delete": true + })) + .expect("metadata maintenance request should parse"); + + assert_eq!(request.retain_recent_metadata_files, 2); + assert!(request.delete); } #[test] @@ -2046,6 +2130,74 @@ mod tests { ); } + #[tokio::test] + async fn table_metadata_maintenance_helper_runs_dry_run_and_delete() { + 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"); + let old = crate::table_catalog::default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let current = crate::table_catalog::default_table_metadata_file_path(&namespace, &table, "00002.metadata.json"); + + seed_object_table_for_metadata_maintenance(&store, &backend, bucket, &namespace, &table, current.clone()).await; + backend + .put_json_with_mod_time(bucket, &old, serde_json::json!({}), Some(OffsetDateTime::UNIX_EPOCH)) + .await; + backend + .put_json_with_mod_time( + bucket, + ¤t, + serde_json::json!({ + "metadata-log": [] + }), + Some(OffsetDateTime::UNIX_EPOCH), + ) + .await; + + let dry_run = table_metadata_maintenance_response( + &store, + bucket, + &namespace, + "events", + TableMetadataMaintenanceRequest { + retain_recent_metadata_files: 0, + delete: false, + }, + ) + .await + .expect("metadata maintenance dry-run should succeed"); + assert_eq!(dry_run.cleanup_candidate_locations, vec![old.clone()]); + assert_eq!(dry_run.deletable_metadata_locations, vec![old.clone()]); + assert!( + backend + .object_exists(bucket, &old) + .await + .expect("old metadata lookup should succeed") + ); + + let deleted = table_metadata_maintenance_response( + &store, + bucket, + &namespace, + "events", + TableMetadataMaintenanceRequest { + retain_recent_metadata_files: 0, + delete: true, + }, + ) + .await + .expect("metadata maintenance delete should succeed"); + assert_eq!(deleted.cleanup_candidate_locations, vec![old.clone()]); + assert_eq!(deleted.deletable_metadata_locations, vec![old.clone()]); + assert!( + !backend + .object_exists(bucket, &old) + .await + .expect("old metadata lookup should succeed after delete") + ); + } + #[test] fn commit_requirements_reject_mismatched_table_uuid() { let metadata = serde_json::json!({ @@ -2167,18 +2319,88 @@ mod tests { impl TestTableCatalogObjectBackend { async fn put_json(&self, bucket: &str, object: &str, value: serde_json::Value) { + self.put_json_with_mod_time(bucket, object, value, None).await; + } + + async fn put_json_with_mod_time( + &self, + bucket: &str, + object: &str, + value: serde_json::Value, + mod_time: Option, + ) { 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()), - mod_time: None, + mod_time, }, ); } } + async fn seed_object_table_for_metadata_maintenance( + store: &crate::table_catalog::ObjectTableCatalogStore, + backend: &TestTableCatalogObjectBackend, + bucket: &str, + namespace: &crate::table_catalog::Namespace, + table: &crate::table_catalog::IdentifierSegment, + current_metadata_location: String, + ) { + store + .put_table_bucket(crate::table_catalog::TableBucketEntry { + version: crate::table_catalog::TABLE_CATALOG_ENTRY_VERSION, + table_bucket: bucket.to_string(), + catalog_type: crate::table_catalog::TABLE_BUCKET_CATALOG_TYPE.to_string(), + warehouse_root: format!("s3://{bucket}/"), + state: crate::table_catalog::TableCatalogEntryState::Active, + properties: BTreeMap::new(), + created_at: None, + updated_at: None, + }) + .await + .expect("table bucket entry should seed"); + store + .create_namespace(crate::table_catalog::NamespaceEntry { + version: crate::table_catalog::TABLE_CATALOG_ENTRY_VERSION, + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + namespace_id: namespace.storage_id(), + state: crate::table_catalog::TableCatalogEntryState::Active, + properties: BTreeMap::new(), + created_at: None, + updated_at: None, + }) + .await + .expect("namespace entry should seed"); + store + .create_table(crate::table_catalog::TableEntry { + version: crate::table_catalog::TABLE_CATALOG_ENTRY_VERSION, + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + table_id: "table-id".to_string(), + table_uuid: "table-uuid".to_string(), + format: "ICEBERG".to_string(), + format_version: 2, + warehouse_location: format!("s3://{bucket}/tables/table-id"), + metadata_location: current_metadata_location, + 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 entry should seed"); + backend + .put_json(bucket, "unrelated/ignored.json", serde_json::json!({})) + .await; + } + #[async_trait::async_trait] impl crate::table_catalog::TableCatalogObjectBackend for TestTableCatalogObjectBackend { async fn read_object( diff --git a/rustfs/src/admin/route_policy.rs b/rustfs/src/admin/route_policy.rs index fe362863e..14d834c90 100644 --- a/rustfs/src/admin/route_policy.rs +++ b/rustfs/src/admin/route_policy.rs @@ -61,6 +61,7 @@ const PROFILING: AdminActionRef = AdminActionRef::new("ProfilingAdminAction"); const REBALANCE: AdminActionRef = AdminActionRef::new("RebalanceAdminAction"); const REGISTER_TABLE: AdminActionRef = AdminActionRef::new("RegisterTableAction"); const REMOVE_USER_FROM_GROUP: AdminActionRef = AdminActionRef::new("RemoveUserFromGroupAdminAction"); +const RUN_TABLE_MAINTENANCE: AdminActionRef = AdminActionRef::new("RunTableMaintenanceAction"); const SERVER_INFO: AdminActionRef = AdminActionRef::new("ServerInfoAdminAction"); const SET_BUCKET_QUOTA: AdminActionRef = AdminActionRef::new("SetBucketQuotaAdminAction"); const SET_BUCKET_TARGET: AdminActionRef = AdminActionRef::new("SetBucketTargetAction"); @@ -649,6 +650,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ DELETE_TABLE, RouteRiskLevel::High, ), + admin( + HttpMethod::Post, + "/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/maintenance/metadata", + RUN_TABLE_MAINTENANCE, + RouteRiskLevel::High, + ), ]; pub const DEFERRED_ADMIN_ROUTE_POLICIES: &[DeferredAdminRoutePolicy] = &[ @@ -822,7 +829,7 @@ mod tests { let table_specs = ADMIN_ROUTE_POLICY_SPECS .iter() .filter(|spec| spec.path().starts_with("/iceberg/v1")); - assert_eq!(table_specs.count(), 11); + assert_eq!(table_specs.count(), 12); assert_action(HttpMethod::Get, "/iceberg/v1/{warehouse}/namespaces", GET_TABLE_NAMESPACE); assert_action(HttpMethod::Post, "/iceberg/v1/{warehouse}/namespaces/{namespace}/tables", CREATE_TABLE); assert_action( @@ -835,6 +842,11 @@ mod tests { "/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}", COMMIT_TABLE, ); + assert_action( + HttpMethod::Post, + "/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/maintenance/metadata", + RUN_TABLE_MAINTENANCE, + ); } #[test] diff --git a/rustfs/src/admin/route_registration_test.rs b/rustfs/src/admin/route_registration_test.rs index a5c25a195..dfc10e1cc 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -303,6 +303,11 @@ fn expected_admin_route_matrix() -> Vec { "/{warehouse}/namespaces/{namespace}/tables/{table}", "/analytics/namespaces/sales/tables/orders", ), + table_route_sample( + Method::POST, + "/{warehouse}/namespaces/{namespace}/tables/{table}/maintenance/metadata", + "/analytics/namespaces/sales/tables/orders/maintenance/metadata", + ), ] } @@ -420,6 +425,11 @@ fn test_register_routes_cover_representative_admin_paths() { assert_route(&router, Method::GET, &table_catalog_path("/analytics/namespaces/sales/tables/orders")); assert_route(&router, Method::POST, &table_catalog_path("/analytics/namespaces/sales/tables/orders")); assert_route(&router, Method::DELETE, &table_catalog_path("/analytics/namespaces/sales/tables/orders")); + assert_route( + &router, + Method::POST, + &table_catalog_path("/analytics/namespaces/sales/tables/orders/maintenance/metadata"), + ); assert_route(&router, Method::POST, &admin_path("/v3/service")); assert_route(&router, Method::GET, &admin_path("/v3/info")); diff --git a/rustfs/src/table_catalog.rs b/rustfs/src/table_catalog.rs index 84fd1bd52..2b7dcefd0 100644 --- a/rustfs/src/table_catalog.rs +++ b/rustfs/src/table_catalog.rs @@ -245,11 +245,26 @@ pub(crate) struct TableCommitResult { pub commit_log: CommitLogEntry, } +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub(crate) struct TableMetadataMaintenanceJob { + pub job_id: String, + pub table_bucket: String, + pub namespace: String, + pub table: String, + pub current_metadata_location: String, + pub current_generation: u64, + pub retain_recent_metadata_files: usize, + pub safety_window_seconds: i64, + pub cleanup_watermark_unix_seconds: i64, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize)] pub(crate) struct TableMetadataMaintenanceReport { + pub job: TableMetadataMaintenanceJob, pub current_metadata_location: String, pub retained_metadata_locations: Vec, pub cleanup_candidate_locations: Vec, + pub deletable_metadata_locations: Vec, } #[derive(Debug, Clone, PartialEq, Eq, Serialize)] @@ -724,12 +739,37 @@ where let cleanup_candidate_locations = metadata_locations .into_iter() .filter(|metadata_location| !retained.contains(metadata_location)) - .collect(); + .collect::>(); + + let now = OffsetDateTime::now_utc(); + let mut deletable_metadata_locations = Vec::new(); + for metadata_location in &cleanup_candidate_locations { + let Some(candidate_object) = self.backend.read_object(table_bucket, metadata_location).await? else { + continue; + }; + if metadata_candidate_is_past_safety_window(candidate_object.mod_time, now) { + deletable_metadata_locations.push(metadata_location.clone()); + } + } + let current_metadata_location = entry.metadata_location; Ok(TableMetadataMaintenanceReport { - current_metadata_location: entry.metadata_location, + job: TableMetadataMaintenanceJob { + job_id: Uuid::new_v4().to_string(), + table_bucket: table_bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + current_metadata_location: current_metadata_location.clone(), + current_generation: entry.generation, + retain_recent_metadata_files, + safety_window_seconds: TABLE_METADATA_CLEANUP_SAFETY_WINDOW_SECONDS, + cleanup_watermark_unix_seconds: (now - Duration::seconds(TABLE_METADATA_CLEANUP_SAFETY_WINDOW_SECONDS)) + .unix_timestamp(), + }, + current_metadata_location, retained_metadata_locations: retained.into_iter().collect(), cleanup_candidate_locations, + deletable_metadata_locations, }) } @@ -826,9 +866,11 @@ where } Ok(TableMetadataMaintenanceReport { + job: report.job, current_metadata_location: entry.metadata_location, retained_metadata_locations: protected.into_iter().collect(), - cleanup_candidate_locations, + cleanup_candidate_locations: cleanup_candidate_locations.clone(), + deletable_metadata_locations: cleanup_candidate_locations, }) } } @@ -2330,6 +2372,42 @@ mod tests { assert_eq!(report.cleanup_candidate_locations, vec![v1, v2]); } + #[tokio::test] + async fn maintenance_dry_run_reports_job_context_and_deletable_candidates() { + let backend = TestCatalogObjectBackend::default(); + let store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + let table = IdentifierSegment::parse("orders").unwrap(); + let old = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let current = default_table_metadata_file_path(&namespace, &table, "00002.metadata.json"); + let fresh = default_table_metadata_file_path(&namespace, &table, "00003.metadata.json"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend.seed_object(bucket, &old, b"{}".to_vec()).await; + backend + .seed_object(bucket, ¤t, br#"{"metadata-log":[]}"#.to_vec()) + .await; + backend + .seed_object_with_mod_time(bucket, &fresh, b"{}".to_vec(), Some(OffsetDateTime::now_utc())) + .await; + + let report = store + .plan_table_metadata_maintenance(bucket, "sales", "orders", 0) + .await + .unwrap(); + + assert_eq!(report.job.table_bucket, bucket); + assert_eq!(report.job.namespace, "sales"); + assert_eq!(report.job.table, "orders"); + assert_eq!(report.job.current_generation, 1); + assert_eq!(report.job.safety_window_seconds, TABLE_METADATA_CLEANUP_SAFETY_WINDOW_SECONDS); + assert!(!report.job.job_id.is_empty()); + assert!(report.job.cleanup_watermark_unix_seconds <= OffsetDateTime::now_utc().unix_timestamp()); + assert_eq!(report.cleanup_candidate_locations, vec![old.clone(), fresh]); + assert_eq!(report.deletable_metadata_locations, vec![old]); + } + #[tokio::test] async fn maintenance_dry_run_keeps_metadata_log_references() { let backend = TestCatalogObjectBackend::default();