feat(table-catalog): harden maintenance scheduler control (#4164)

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
Henry Guo
2026-07-02 14:11:07 +08:00
committed by GitHub
parent 48f8443626
commit ba0276dd6f
+193 -25
View File
@@ -944,6 +944,11 @@ struct TableMaintenanceWorkerControlReport<'a> {
now: OffsetDateTime,
}
enum TableMaintenanceWorkerPreflight {
Ready(TableMaintenanceEffectiveConfig),
Complete(Box<TableMetadataMaintenanceReport>),
}
#[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<TableMetadataMaintenanceReport> {
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<TableMaintenanceWorkerPreflight> {
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(&current.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(&current.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<TableMetadataMaintenanceReport> {
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<usize>>,
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<Box<dyn Send>> {
{
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, &current, 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, &current, 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();