From ba0276dd6f8942272c9a7b1efa41a567e647aaca Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Thu, 2 Jul 2026 14:11:07 +0800 Subject: [PATCH] feat(table-catalog): harden maintenance scheduler control (#4164) Co-authored-by: Henry Guo --- rustfs/src/table_catalog.rs | 218 +++++++++++++++++++++++++++++++----- 1 file changed, 193 insertions(+), 25 deletions(-) diff --git a/rustfs/src/table_catalog.rs b/rustfs/src/table_catalog.rs index d72f00613..174233654 100644 --- a/rustfs/src/table_catalog.rs +++ b/rustfs/src/table_catalog.rs @@ -944,6 +944,11 @@ struct TableMaintenanceWorkerControlReport<'a> { now: OffsetDateTime, } +enum TableMaintenanceWorkerPreflight { + Ready(TableMaintenanceEffectiveConfig), + Complete(Box), +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize)] pub(crate) struct TableCatalogExport { pub table_bucket: TableBucketEntry, @@ -4040,36 +4045,129 @@ where worker_id: String, now: OffsetDateTime, ) -> TableCatalogStoreResult { + let namespace = parse_namespace_for_store(namespace)?; + let table = parse_table_for_store(table)?; + let table_path = self.paths.table_entry_path(table_bucket, &namespace, &table); + let namespace_name = namespace.public_name(); + let table_name = table.as_str().to_string(); + + let effective = { + let _guard = self.backend.acquire_write_lock(self.catalog_bucket(), &table_path).await?; + match self + .table_metadata_maintenance_worker_preflight(table_bucket, &namespace_name, &table_name, &worker_id, now) + .await? + { + TableMaintenanceWorkerPreflight::Ready(effective) => effective, + TableMaintenanceWorkerPreflight::Complete(report) => return Ok(*report), + } + }; + + let mut report = self + .plan_table_metadata_maintenance( + table_bucket, + &namespace_name, + &table_name, + effective.config.retain_recent_metadata_files, + ) + .await?; + + let (report, effective, delete) = { + let _guard = self.backend.acquire_write_lock(self.catalog_bucket(), &table_path).await?; + let effective = match self + .table_metadata_maintenance_worker_preflight(table_bucket, &namespace_name, &table_name, &worker_id, now) + .await? + { + TableMaintenanceWorkerPreflight::Ready(effective) => effective, + TableMaintenanceWorkerPreflight::Complete(report) => return Ok(*report), + }; + + if report.job.retain_recent_metadata_files != effective.config.retain_recent_metadata_files { + return Err(TableCatalogStoreError::Conflict( + "maintenance config changed before worker claim".to_string(), + )); + } + let Some((entry, _)) = self.read_table_with_etag_unlocked(table_bucket, &namespace, &table).await? else { + return Err(TableCatalogStoreError::NotFound(format!( + "table {}/{}/{}", + table_bucket, namespace_name, table_name + ))); + }; + if entry.metadata_location != report.current_metadata_location { + return Err(TableCatalogStoreError::Conflict( + "current metadata location changed before maintenance worker claim".to_string(), + )); + } + + let started_at = maintenance_timestamp(now); + report.job.operation = if effective.config.delete_enabled { + TableMetadataMaintenanceOperation::Delete + } else { + TableMetadataMaintenanceOperation::DryRun + }; + report.job.status = TableMetadataMaintenanceJobStatus::Running; + report.job.failure_reason = None; + report.job.config_source = effective.source; + report.job.worker_id = Some(worker_id); + report.job.lease_id = Uuid::new_v4().to_string(); + report.job.attempt = 1; + report.job.max_retry_attempts = effective.config.max_retry_attempts; + report.job.next_retry_after = None; + report.job.quarantine_enabled = effective.config.quarantine_enabled; + report.job.quarantine_retention_seconds = effective.config.quarantine_retention_seconds; + report.job.heartbeat_at = Some(started_at.clone()); + report.job.started_at = Some(started_at); + report.job.finished_at = None; + refresh_table_maintenance_report_recommended_actions(&mut report); + self.put_table_metadata_maintenance_report(&report).await?; + + let delete = effective.config.delete_enabled; + (report, effective, delete) + }; + + self.finish_table_metadata_maintenance_run(table_bucket, &namespace_name, &table_name, delete, &effective, report) + .await + } + + async fn table_metadata_maintenance_worker_preflight( + &self, + table_bucket: &str, + namespace: &str, + table: &str, + worker_id: &str, + now: OffsetDateTime, + ) -> TableCatalogStoreResult { let effective = self .get_effective_table_maintenance_config(table_bucket, namespace, table) .await?; if !effective.config.background_enabled { - return self + let report = self .put_table_metadata_maintenance_worker_control_report(TableMaintenanceWorkerControlReport { table_bucket, namespace, table, - worker_id, + worker_id: worker_id.to_string(), effective: &effective, status: TableMetadataMaintenanceJobStatus::Disabled, reason: "background maintenance is disabled", now, }) - .await; + .await?; + return Ok(TableMaintenanceWorkerPreflight::Complete(Box::new(report))); } if effective.config.worker_paused { - return self + let report = self .put_table_metadata_maintenance_worker_control_report(TableMaintenanceWorkerControlReport { table_bucket, namespace, table, - worker_id, + worker_id: worker_id.to_string(), effective: &effective, status: TableMetadataMaintenanceJobStatus::Paused, reason: "background maintenance worker is paused", now, }) - .await; + .await?; + return Ok(TableMaintenanceWorkerPreflight::Complete(Box::new(report))); } if let Some(current) = self @@ -4078,27 +4176,20 @@ where { if matches!(current.job.status, TableMetadataMaintenanceJobStatus::Running) { if table_maintenance_job_lease_is_active(¤t.job, effective.config.worker_lease_timeout_seconds, now) { - return Ok(current); + return Ok(TableMaintenanceWorkerPreflight::Complete(Box::new(current))); } let mut expired = current; expired.job.status = TableMetadataMaintenanceJobStatus::Failed; expired.job.failure_reason = Some("maintenance worker lease expired".to_string()); expired.job.finished_at = Some(maintenance_timestamp(now)); + refresh_table_maintenance_report_recommended_actions(&mut expired); self.put_table_metadata_maintenance_report(&expired).await?; } else if table_maintenance_job_retry_is_pending(¤t.job, now) { - return Ok(current); + return Ok(TableMaintenanceWorkerPreflight::Complete(Box::new(current))); } } - self.run_table_metadata_maintenance_with_config( - table_bucket, - namespace, - table, - effective.config.delete_enabled, - Some(worker_id), - effective, - ) - .await + Ok(TableMaintenanceWorkerPreflight::Ready(effective)) } pub(crate) async fn heartbeat_table_metadata_maintenance_job( @@ -5066,15 +5157,27 @@ where refresh_table_maintenance_report_recommended_actions(&mut report); self.put_table_metadata_maintenance_report(&report).await?; + self.finish_table_metadata_maintenance_run(table_bucket, namespace, table, delete, &effective, report) + .await + } + + async fn finish_table_metadata_maintenance_run( + &self, + table_bucket: &str, + namespace: &str, + table: &str, + delete: bool, + effective: &TableMaintenanceEffectiveConfig, + mut report: TableMetadataMaintenanceReport, + ) -> TableCatalogStoreResult { if delete && !effective.config.delete_enabled { - let mut failed = report; - failed.job.status = TableMetadataMaintenanceJobStatus::Failed; - failed.job.failure_reason = Some(TABLE_MAINTENANCE_DELETE_DISABLED_REASON.to_string()); - apply_maintenance_retry_after(&mut failed.job, &effective.config, OffsetDateTime::now_utc()); - failed.job.finished_at = Some(maintenance_timestamp(OffsetDateTime::now_utc())); - refresh_table_maintenance_report_recommended_actions(&mut failed); - self.put_table_metadata_maintenance_report(&failed).await?; - return Ok(failed); + report.job.status = TableMetadataMaintenanceJobStatus::Failed; + report.job.failure_reason = Some(TABLE_MAINTENANCE_DELETE_DISABLED_REASON.to_string()); + apply_maintenance_retry_after(&mut report.job, &effective.config, OffsetDateTime::now_utc()); + report.job.finished_at = Some(maintenance_timestamp(OffsetDateTime::now_utc())); + refresh_table_maintenance_report_recommended_actions(&mut report); + self.put_table_metadata_maintenance_report(&report).await?; + return Ok(report); } if delete { @@ -10089,6 +10192,7 @@ mod tests { fail_delete_attempts: BTreeMap<(String, String), BTreeSet>, put_attempts: BTreeMap<(String, String), usize>, delete_attempts: BTreeMap<(String, String), usize>, + write_lock_acquisitions: BTreeMap<(String, String), usize>, read_calls: usize, list_calls: usize, next_etag: u64, @@ -10147,6 +10251,16 @@ mod tests { state.list_calls = 0; } + async fn write_lock_acquisition_count(&self, bucket: &str, object: &str) -> usize { + self.state + .lock() + .await + .write_lock_acquisitions + .get(&(bucket.to_string(), object.to_string())) + .copied() + .unwrap_or_default() + } + async fn fail_next_read(&self, bucket: &str, object: &str) { let mut state = self.state.lock().await; let key = (bucket.to_string(), object.to_string()); @@ -10525,6 +10639,13 @@ mod tests { } async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult> { + { + let mut state = self.state.lock().await; + *state + .write_lock_acquisitions + .entry((bucket.to_string(), object.to_string())) + .or_default() += 1; + } let lock = { let mut locks = self.locks.lock().await; locks @@ -12029,12 +12150,14 @@ mod tests { 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 table_path = store.paths.table_entry_path(bucket, &namespace, &table); 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; + let before_worker_run = backend.write_lock_acquisition_count(RUSTFS_META_BUCKET, &table_path).await; let report = store .run_table_metadata_maintenance_worker_once(bucket, "sales", "orders", "worker-a".to_string()) @@ -12048,6 +12171,10 @@ mod tests { vec![TableMaintenanceRecommendedAction::EnableBackgroundMaintenance] ); assert_eq!(report.job.deleted_metadata_file_count, 0); + assert_eq!( + backend.write_lock_acquisition_count(RUSTFS_META_BUCKET, &table_path).await, + before_worker_run + 1 + ); assert!(backend.object_exists(bucket, &old).await.unwrap()); } @@ -12098,6 +12225,47 @@ mod tests { assert!(backend.object_exists(bucket, &old).await.unwrap()); } + #[tokio::test] + async fn maintenance_worker_run_claims_table_entry_write_lock() { + let backend = TestCatalogObjectBackend::default(); + let store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").expect("namespace should parse"); + let table = IdentifierSegment::parse("orders").expect("table should parse"); + let current = default_table_metadata_file_path(&namespace, &table, "00002.metadata.json"); + let table_path = store.paths.table_entry_path(bucket, &namespace, &table); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend + .seed_object(bucket, ¤t, br#"{"metadata-log":[]}"#.to_vec()) + .await; + store + .put_table_maintenance_config( + bucket, + "sales", + "orders", + TableMaintenanceConfig { + version: TABLE_MAINTENANCE_CONFIG_VERSION, + background_enabled: true, + ..Default::default() + }, + ) + .await + .expect("background maintenance config should persist"); + let before_worker_run = backend.write_lock_acquisition_count(RUSTFS_META_BUCKET, &table_path).await; + + let report = store + .run_table_metadata_maintenance_worker_once("analytics", "sales", "orders", "worker-a".to_string()) + .await + .expect("worker run should complete"); + + assert_eq!(report.job.worker_id.as_deref(), Some("worker-a")); + assert_eq!( + backend.write_lock_acquisition_count(RUSTFS_META_BUCKET, &table_path).await, + before_worker_run + 2 + ); + } + #[tokio::test] async fn maintenance_scheduler_report_marks_disabled_default() { let backend = TestCatalogObjectBackend::default();