mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 08:18:18 +00:00
feat(table-catalog): add maintenance audit timeline (#4207)
Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
@@ -88,6 +88,7 @@ catalog extension.
|
||||
| Manifest/data/delete reachability cleanup | Supported | Reads manifest-list and manifest Avro references, reports reachable objects, and deletes only unreferenced table objects that pass the safety window. |
|
||||
| Maintenance worker run endpoint | Preview / controlled | Supports run-once execution, current-job backpressure, retry deferral, lease expiry recovery, and heartbeat updates. |
|
||||
| Maintenance scheduler guardrails | Preview / controlled | Exposes disabled, paused, ready, active-job backpressure, retry deferral, quarantine boundary, recommended actions, and recent maintenance job audit timeline state for external schedulers and operators. |
|
||||
| Maintenance audit events | Preview / controlled | Job reports and scheduler job summaries include structured audit events for planning, worker transitions, heartbeats, lease expiry recovery, and mutating quarantine operations. |
|
||||
| Maintenance quarantine operations | Preview / controlled | Lets operators inspect, release, retry, or abandon the current quarantined maintenance job without moving the table pointer. |
|
||||
| Compaction planning | Preview / controlled | Plans partition-local and sort-order-local binpack candidates for Parquet files and does not mix data files from different partition directories or sort orders in one rewrite group. |
|
||||
| Compaction commit | Preview / controlled | Can commit a safe partition-local Parquet rewrite through the catalog while preserving Iceberg data file sort order IDs in the rewritten manifest. |
|
||||
|
||||
+347
-11
@@ -552,6 +552,48 @@ pub(crate) struct TableMaintenanceSchedulerQuarantineBoundary {
|
||||
pub source_job_id: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
|
||||
pub(crate) enum TableMaintenanceAuditActor {
|
||||
Scheduler,
|
||||
Worker,
|
||||
Operator,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
|
||||
pub(crate) enum TableMaintenanceAuditAction {
|
||||
Planned,
|
||||
WorkerControl,
|
||||
WorkerStarted,
|
||||
WorkerHeartbeat,
|
||||
WorkerLeaseExpired,
|
||||
WorkerSucceeded,
|
||||
WorkerFailed,
|
||||
QuarantineRelease,
|
||||
QuarantineRetry,
|
||||
QuarantineAbandon,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub(crate) struct TableMaintenanceAuditEvent {
|
||||
pub timestamp: String,
|
||||
pub actor: TableMaintenanceAuditActor,
|
||||
pub action: TableMaintenanceAuditAction,
|
||||
#[serde(default)]
|
||||
pub reason: Option<String>,
|
||||
#[serde(default, rename = "before-status")]
|
||||
pub before_status: Option<TableMetadataMaintenanceJobStatus>,
|
||||
#[serde(default, rename = "after-status")]
|
||||
pub after_status: Option<TableMetadataMaintenanceJobStatus>,
|
||||
#[serde(default, rename = "before-quarantined-object-count")]
|
||||
pub before_quarantined_object_count: Option<usize>,
|
||||
#[serde(default, rename = "after-quarantined-object-count")]
|
||||
pub after_quarantined_object_count: Option<usize>,
|
||||
#[serde(default, rename = "recommended-actions")]
|
||||
pub recommended_actions: Vec<TableMaintenanceRecommendedAction>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
|
||||
pub(crate) enum TableMaintenanceQuarantineAction {
|
||||
@@ -588,6 +630,8 @@ pub(crate) struct TableMaintenanceSchedulerJobSummary {
|
||||
pub heartbeat_at: Option<String>,
|
||||
pub next_retry_after: Option<String>,
|
||||
pub recommended_actions: Vec<TableMaintenanceRecommendedAction>,
|
||||
#[serde(default, rename = "audit-events")]
|
||||
pub audit_events: Vec<TableMaintenanceAuditEvent>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
@@ -677,6 +721,8 @@ pub(crate) struct TableMetadataMaintenanceReport {
|
||||
pub snapshot_expiration: Option<TableSnapshotExpirationReport>,
|
||||
#[serde(default)]
|
||||
pub compaction: Option<TableCompactionPlanningReport>,
|
||||
#[serde(default, rename = "audit-events")]
|
||||
pub audit_events: Vec<TableMaintenanceAuditEvent>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
|
||||
@@ -4070,6 +4116,8 @@ where
|
||||
));
|
||||
}
|
||||
|
||||
let before_status = Some(report.job.status.clone());
|
||||
let before_quarantined_object_count = Some(report.job.quarantined_object_count);
|
||||
report.job.quarantined_object_count = 0;
|
||||
match &action {
|
||||
TableMaintenanceQuarantineAction::Inspect => unreachable!("inspect branch handled before mutation"),
|
||||
@@ -4091,6 +4139,21 @@ where
|
||||
}
|
||||
}
|
||||
refresh_table_maintenance_report_recommended_actions(&mut report);
|
||||
let audit_action = match &action {
|
||||
TableMaintenanceQuarantineAction::Inspect => unreachable!("inspect branch handled before mutation"),
|
||||
TableMaintenanceQuarantineAction::Release => TableMaintenanceAuditAction::QuarantineRelease,
|
||||
TableMaintenanceQuarantineAction::Retry => TableMaintenanceAuditAction::QuarantineRetry,
|
||||
TableMaintenanceQuarantineAction::Abandon => TableMaintenanceAuditAction::QuarantineAbandon,
|
||||
};
|
||||
push_table_maintenance_audit_event(
|
||||
&mut report,
|
||||
OffsetDateTime::now_utc(),
|
||||
TableMaintenanceAuditActor::Operator,
|
||||
audit_action,
|
||||
request.reason,
|
||||
before_status,
|
||||
before_quarantined_object_count,
|
||||
);
|
||||
self.put_table_metadata_maintenance_report(&report).await?;
|
||||
report
|
||||
};
|
||||
@@ -4228,6 +4291,15 @@ where
|
||||
report.job.started_at = Some(started_at);
|
||||
report.job.finished_at = None;
|
||||
refresh_table_maintenance_report_recommended_actions(&mut report);
|
||||
push_table_maintenance_audit_event(
|
||||
&mut report,
|
||||
now,
|
||||
TableMaintenanceAuditActor::Worker,
|
||||
TableMaintenanceAuditAction::WorkerStarted,
|
||||
None,
|
||||
Some(TableMetadataMaintenanceJobStatus::Successful),
|
||||
Some(0),
|
||||
);
|
||||
self.put_table_metadata_maintenance_report(&report).await?;
|
||||
|
||||
let delete = effective.config.delete_enabled;
|
||||
@@ -4289,10 +4361,21 @@ where
|
||||
return Ok(TableMaintenanceWorkerPreflight::Complete(Box::new(current)));
|
||||
}
|
||||
let mut expired = current;
|
||||
let before_status = Some(expired.job.status.clone());
|
||||
let before_quarantined_object_count = Some(expired.job.quarantined_object_count);
|
||||
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);
|
||||
push_table_maintenance_audit_event(
|
||||
&mut expired,
|
||||
now,
|
||||
TableMaintenanceAuditActor::Scheduler,
|
||||
TableMaintenanceAuditAction::WorkerLeaseExpired,
|
||||
Some("maintenance worker lease expired".to_string()),
|
||||
before_status,
|
||||
before_quarantined_object_count,
|
||||
);
|
||||
self.put_table_metadata_maintenance_report(&expired).await?;
|
||||
} else if table_maintenance_job_retry_is_pending(¤t.job, now) {
|
||||
return Ok(TableMaintenanceWorkerPreflight::Complete(Box::new(current)));
|
||||
@@ -4365,6 +4448,17 @@ where
|
||||
}
|
||||
|
||||
report.job.heartbeat_at = Some(maintenance_timestamp(now));
|
||||
refresh_table_maintenance_report_recommended_actions(&mut report);
|
||||
let before_quarantined_object_count = Some(report.job.quarantined_object_count);
|
||||
push_table_maintenance_audit_event(
|
||||
&mut report,
|
||||
now,
|
||||
TableMaintenanceAuditActor::Worker,
|
||||
TableMaintenanceAuditAction::WorkerHeartbeat,
|
||||
None,
|
||||
Some(TableMetadataMaintenanceJobStatus::Running),
|
||||
before_quarantined_object_count,
|
||||
);
|
||||
self.put_table_metadata_maintenance_report(&report).await?;
|
||||
Ok(report)
|
||||
}
|
||||
@@ -4438,8 +4532,18 @@ where
|
||||
reachability_graph: TableMaintenanceReachabilityGraphReport::default(),
|
||||
snapshot_expiration: None,
|
||||
compaction: None,
|
||||
audit_events: Vec::new(),
|
||||
};
|
||||
let report = table_maintenance_report_with_recommended_actions(report);
|
||||
let mut report = table_maintenance_report_with_recommended_actions(report);
|
||||
push_table_maintenance_audit_event(
|
||||
&mut report,
|
||||
control.now,
|
||||
TableMaintenanceAuditActor::Scheduler,
|
||||
TableMaintenanceAuditAction::WorkerControl,
|
||||
Some(control.reason.to_string()),
|
||||
None,
|
||||
None,
|
||||
);
|
||||
self.put_table_metadata_maintenance_report(&report).await?;
|
||||
Ok(report)
|
||||
}
|
||||
@@ -5134,7 +5238,7 @@ where
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(table_maintenance_report_with_recommended_actions(TableMetadataMaintenanceReport {
|
||||
let mut report = table_maintenance_report_with_recommended_actions(TableMetadataMaintenanceReport {
|
||||
job: TableMetadataMaintenanceJob {
|
||||
job_id: Uuid::new_v4().to_string(),
|
||||
table_bucket: table_bucket.to_string(),
|
||||
@@ -5185,7 +5289,18 @@ where
|
||||
reachability_graph,
|
||||
snapshot_expiration: None,
|
||||
compaction: None,
|
||||
}))
|
||||
audit_events: Vec::new(),
|
||||
});
|
||||
push_table_maintenance_audit_event(
|
||||
&mut report,
|
||||
now,
|
||||
TableMaintenanceAuditActor::Scheduler,
|
||||
TableMaintenanceAuditAction::Planned,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
);
|
||||
Ok(report)
|
||||
}
|
||||
|
||||
pub(crate) async fn delete_table_metadata_maintenance_candidates(
|
||||
@@ -5247,7 +5362,8 @@ where
|
||||
.plan_table_metadata_maintenance(table_bucket, namespace, table, effective.config.retain_recent_metadata_files)
|
||||
.await?;
|
||||
|
||||
let started_at = maintenance_timestamp(OffsetDateTime::now_utc());
|
||||
let started_at_time = OffsetDateTime::now_utc();
|
||||
let started_at = maintenance_timestamp(started_at_time);
|
||||
report.job.operation = if delete {
|
||||
TableMetadataMaintenanceOperation::Delete
|
||||
} else {
|
||||
@@ -5267,6 +5383,15 @@ where
|
||||
report.job.started_at = Some(started_at);
|
||||
report.job.finished_at = None;
|
||||
refresh_table_maintenance_report_recommended_actions(&mut report);
|
||||
push_table_maintenance_audit_event(
|
||||
&mut report,
|
||||
started_at_time,
|
||||
TableMaintenanceAuditActor::Worker,
|
||||
TableMaintenanceAuditAction::WorkerStarted,
|
||||
None,
|
||||
Some(TableMetadataMaintenanceJobStatus::Successful),
|
||||
Some(0),
|
||||
);
|
||||
self.put_table_metadata_maintenance_report(&report).await?;
|
||||
|
||||
self.finish_table_metadata_maintenance_run(table_bucket, namespace, table, delete, &effective, report)
|
||||
@@ -5283,11 +5408,23 @@ where
|
||||
mut report: TableMetadataMaintenanceReport,
|
||||
) -> TableCatalogStoreResult<TableMetadataMaintenanceReport> {
|
||||
if delete && !effective.config.delete_enabled {
|
||||
let finished_at = OffsetDateTime::now_utc();
|
||||
let before_status = Some(report.job.status.clone());
|
||||
let before_quarantined_object_count = Some(report.job.quarantined_object_count);
|
||||
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()));
|
||||
apply_maintenance_retry_after(&mut report.job, &effective.config, finished_at);
|
||||
report.job.finished_at = Some(maintenance_timestamp(finished_at));
|
||||
refresh_table_maintenance_report_recommended_actions(&mut report);
|
||||
push_table_maintenance_audit_event(
|
||||
&mut report,
|
||||
finished_at,
|
||||
TableMaintenanceAuditActor::Worker,
|
||||
TableMaintenanceAuditAction::WorkerFailed,
|
||||
Some(TABLE_MAINTENANCE_DELETE_DISABLED_REASON.to_string()),
|
||||
before_status,
|
||||
before_quarantined_object_count,
|
||||
);
|
||||
self.put_table_metadata_maintenance_report(&report).await?;
|
||||
return Ok(report);
|
||||
}
|
||||
@@ -5300,25 +5437,62 @@ where
|
||||
{
|
||||
Ok(report) => report,
|
||||
Err(err) => {
|
||||
let finished_at = OffsetDateTime::now_utc();
|
||||
let mut failed = running_report;
|
||||
let before_status = Some(failed.job.status.clone());
|
||||
let before_quarantined_object_count = Some(failed.job.quarantined_object_count);
|
||||
let reason = err.to_string();
|
||||
failed.job.status = TableMetadataMaintenanceJobStatus::Failed;
|
||||
failed.job.failure_reason = Some(err.to_string());
|
||||
apply_maintenance_retry_after(&mut failed.job, &effective.config, OffsetDateTime::now_utc());
|
||||
failed.job.finished_at = Some(maintenance_timestamp(OffsetDateTime::now_utc()));
|
||||
failed.job.failure_reason = Some(reason.clone());
|
||||
apply_maintenance_retry_after(&mut failed.job, &effective.config, finished_at);
|
||||
failed.job.finished_at = Some(maintenance_timestamp(finished_at));
|
||||
refresh_table_maintenance_report_recommended_actions(&mut failed);
|
||||
push_table_maintenance_audit_event(
|
||||
&mut failed,
|
||||
finished_at,
|
||||
TableMaintenanceAuditActor::Worker,
|
||||
TableMaintenanceAuditAction::WorkerFailed,
|
||||
Some(reason),
|
||||
before_status,
|
||||
before_quarantined_object_count,
|
||||
);
|
||||
self.put_table_metadata_maintenance_report(&failed).await?;
|
||||
return Err(err);
|
||||
}
|
||||
};
|
||||
deleted.job.finished_at = Some(maintenance_timestamp(OffsetDateTime::now_utc()));
|
||||
let finished_at = OffsetDateTime::now_utc();
|
||||
let before_status = Some(TableMetadataMaintenanceJobStatus::Running);
|
||||
let before_quarantined_object_count = Some(0);
|
||||
deleted.job.finished_at = Some(maintenance_timestamp(finished_at));
|
||||
refresh_table_maintenance_report_recommended_actions(&mut deleted);
|
||||
push_table_maintenance_audit_event(
|
||||
&mut deleted,
|
||||
finished_at,
|
||||
TableMaintenanceAuditActor::Worker,
|
||||
TableMaintenanceAuditAction::WorkerSucceeded,
|
||||
None,
|
||||
before_status,
|
||||
before_quarantined_object_count,
|
||||
);
|
||||
self.put_table_metadata_maintenance_report(&deleted).await?;
|
||||
return Ok(deleted);
|
||||
}
|
||||
|
||||
let finished_at = OffsetDateTime::now_utc();
|
||||
let before_status = Some(report.job.status.clone());
|
||||
let before_quarantined_object_count = Some(report.job.quarantined_object_count);
|
||||
report.job.status = TableMetadataMaintenanceJobStatus::Successful;
|
||||
report.job.finished_at = Some(maintenance_timestamp(OffsetDateTime::now_utc()));
|
||||
report.job.finished_at = Some(maintenance_timestamp(finished_at));
|
||||
refresh_table_maintenance_report_recommended_actions(&mut report);
|
||||
push_table_maintenance_audit_event(
|
||||
&mut report,
|
||||
finished_at,
|
||||
TableMaintenanceAuditActor::Worker,
|
||||
TableMaintenanceAuditAction::WorkerSucceeded,
|
||||
None,
|
||||
before_status,
|
||||
before_quarantined_object_count,
|
||||
);
|
||||
self.put_table_metadata_maintenance_report(&report).await?;
|
||||
Ok(report)
|
||||
}
|
||||
@@ -5509,6 +5683,7 @@ where
|
||||
reachability_graph: report.reachability_graph,
|
||||
snapshot_expiration: report.snapshot_expiration,
|
||||
compaction: report.compaction,
|
||||
audit_events: report.audit_events,
|
||||
}))
|
||||
}
|
||||
}
|
||||
@@ -8949,6 +9124,28 @@ fn table_maintenance_quarantine_operator_reason(action: &str, reason: Option<&st
|
||||
}
|
||||
}
|
||||
|
||||
fn push_table_maintenance_audit_event(
|
||||
report: &mut TableMetadataMaintenanceReport,
|
||||
timestamp: OffsetDateTime,
|
||||
actor: TableMaintenanceAuditActor,
|
||||
action: TableMaintenanceAuditAction,
|
||||
reason: Option<String>,
|
||||
before_status: Option<TableMetadataMaintenanceJobStatus>,
|
||||
before_quarantined_object_count: Option<usize>,
|
||||
) {
|
||||
report.audit_events.push(TableMaintenanceAuditEvent {
|
||||
timestamp: maintenance_timestamp(timestamp),
|
||||
actor,
|
||||
action,
|
||||
reason,
|
||||
before_status,
|
||||
after_status: Some(report.job.status.clone()),
|
||||
before_quarantined_object_count,
|
||||
after_quarantined_object_count: Some(report.job.quarantined_object_count),
|
||||
recommended_actions: report.job.recommended_actions.clone(),
|
||||
});
|
||||
}
|
||||
|
||||
fn table_maintenance_recommended_actions(job: &TableMetadataMaintenanceJob) -> Vec<TableMaintenanceRecommendedAction> {
|
||||
let mut actions = Vec::new();
|
||||
match job.status {
|
||||
@@ -9024,6 +9221,7 @@ fn table_maintenance_scheduler_job_summary(report: &TableMetadataMaintenanceRepo
|
||||
heartbeat_at: report.job.heartbeat_at.clone(),
|
||||
next_retry_after: report.job.next_retry_after.clone(),
|
||||
recommended_actions: report.job.recommended_actions.clone(),
|
||||
audit_events: report.audit_events.clone(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12738,6 +12936,18 @@ mod tests {
|
||||
assert_eq!(result.report.job.job_id, failed.job.job_id);
|
||||
assert_eq!(result.report.job.quarantined_object_count, 0);
|
||||
assert!(result.report.job.next_retry_after.is_none());
|
||||
let event = result
|
||||
.report
|
||||
.audit_events
|
||||
.last()
|
||||
.expect("quarantine retry should append an audit event");
|
||||
assert_eq!(event.action, TableMaintenanceAuditAction::QuarantineRetry);
|
||||
assert_eq!(event.actor, TableMaintenanceAuditActor::Operator);
|
||||
assert_eq!(event.reason.as_deref(), Some("operator reviewed retained candidates"));
|
||||
assert_eq!(event.before_status, Some(TableMetadataMaintenanceJobStatus::Failed));
|
||||
assert_eq!(event.after_status, Some(TableMetadataMaintenanceJobStatus::Failed));
|
||||
assert_eq!(event.before_quarantined_object_count, Some(2));
|
||||
assert_eq!(event.after_quarantined_object_count, Some(0));
|
||||
assert!(
|
||||
!result
|
||||
.report
|
||||
@@ -12747,6 +12957,16 @@ mod tests {
|
||||
);
|
||||
assert_eq!(result.scheduler.status, TableMaintenanceSchedulerStatus::Ready);
|
||||
assert!(!result.scheduler.quarantine.active);
|
||||
let summary = result
|
||||
.scheduler
|
||||
.audit_timeline
|
||||
.iter()
|
||||
.find(|summary| summary.job_id == failed.job.job_id)
|
||||
.expect("scheduler timeline should include the retried job");
|
||||
assert_eq!(
|
||||
summary.audit_events.last().map(|event| event.action.clone()),
|
||||
Some(TableMaintenanceAuditAction::QuarantineRetry)
|
||||
);
|
||||
|
||||
let current_report = store
|
||||
.get_table_metadata_maintenance_report(bucket, "sales", "orders", MAINTENANCE_JOB_ALIAS_CURRENT)
|
||||
@@ -12793,6 +13013,7 @@ mod tests {
|
||||
.expect("current maintenance report should load")
|
||||
.expect("current maintenance report should exist");
|
||||
assert_eq!(current_report.job.quarantined_object_count, 2);
|
||||
assert_eq!(current_report.audit_events, failed.audit_events);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -12899,6 +13120,113 @@ mod tests {
|
||||
assert_eq!(error, TableCatalogStoreError::Conflict("maintenance job is not current".to_string()));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn maintenance_worker_run_records_audit_timeline_events() {
|
||||
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");
|
||||
|
||||
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 report = store
|
||||
.run_table_metadata_maintenance_worker_once(bucket, "sales", "orders", "worker-a".to_string())
|
||||
.await
|
||||
.expect("maintenance worker should finish");
|
||||
|
||||
let actions = report
|
||||
.audit_events
|
||||
.iter()
|
||||
.map(|event| event.action.clone())
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(
|
||||
actions,
|
||||
vec![
|
||||
TableMaintenanceAuditAction::Planned,
|
||||
TableMaintenanceAuditAction::WorkerStarted,
|
||||
TableMaintenanceAuditAction::WorkerSucceeded,
|
||||
]
|
||||
);
|
||||
assert_eq!(report.audit_events[1].actor, TableMaintenanceAuditActor::Worker);
|
||||
assert_eq!(report.audit_events[1].before_status, Some(TableMetadataMaintenanceJobStatus::Successful));
|
||||
assert_eq!(report.audit_events[2].after_status, Some(TableMetadataMaintenanceJobStatus::Successful));
|
||||
|
||||
let scheduler = store
|
||||
.get_table_maintenance_scheduler_report(bucket, "sales", "orders")
|
||||
.await
|
||||
.expect("scheduler report should load");
|
||||
let summary = scheduler.current_job.expect("current job should be visible");
|
||||
assert_eq!(summary.job_id, report.job.job_id);
|
||||
assert_eq!(summary.audit_events, report.audit_events);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn maintenance_heartbeat_appends_worker_audit_event() {
|
||||
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 now = OffsetDateTime::UNIX_EPOCH + Duration::seconds(100);
|
||||
|
||||
seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await;
|
||||
backend
|
||||
.seed_object(bucket, ¤t, br#"{"metadata-log":[]}"#.to_vec())
|
||||
.await;
|
||||
let mut running = store
|
||||
.plan_table_metadata_maintenance(bucket, "sales", "orders", 0)
|
||||
.await
|
||||
.expect("maintenance report should be planned");
|
||||
running.job.status = TableMetadataMaintenanceJobStatus::Running;
|
||||
running.job.worker_id = Some("worker-a".to_string());
|
||||
running.job.lease_id = "lease-a".to_string();
|
||||
running.job.heartbeat_at = Some(maintenance_timestamp(now - Duration::seconds(10)));
|
||||
store
|
||||
.put_table_metadata_maintenance_report(&running)
|
||||
.await
|
||||
.expect("running maintenance report should be seeded");
|
||||
|
||||
let heartbeat = store
|
||||
.heartbeat_table_metadata_maintenance_job_at(
|
||||
TableMaintenanceHeartbeatRef {
|
||||
table_bucket: bucket,
|
||||
namespace: "sales",
|
||||
table: "orders",
|
||||
job_id: &running.job.job_id,
|
||||
lease_id: "lease-a",
|
||||
worker_id: "worker-a",
|
||||
},
|
||||
now,
|
||||
)
|
||||
.await
|
||||
.expect("heartbeat should update the running job");
|
||||
|
||||
let event = heartbeat.audit_events.last().expect("heartbeat should append an audit event");
|
||||
assert_eq!(event.action, TableMaintenanceAuditAction::WorkerHeartbeat);
|
||||
assert_eq!(event.actor, TableMaintenanceAuditActor::Worker);
|
||||
assert_eq!(event.before_status, Some(TableMetadataMaintenanceJobStatus::Running));
|
||||
assert_eq!(event.after_status, Some(TableMetadataMaintenanceJobStatus::Running));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn maintenance_worker_run_defers_until_retry_after() {
|
||||
let backend = TestCatalogObjectBackend::default();
|
||||
@@ -13089,6 +13417,14 @@ mod tests {
|
||||
expired.job.recommended_actions,
|
||||
vec![TableMaintenanceRecommendedAction::InvestigateFailure]
|
||||
);
|
||||
let event = expired
|
||||
.audit_events
|
||||
.last()
|
||||
.expect("expired lease recovery should append an audit event");
|
||||
assert_eq!(event.action, TableMaintenanceAuditAction::WorkerLeaseExpired);
|
||||
assert_eq!(event.actor, TableMaintenanceAuditActor::Scheduler);
|
||||
assert_eq!(event.before_status, Some(TableMetadataMaintenanceJobStatus::Running));
|
||||
assert_eq!(event.after_status, Some(TableMetadataMaintenanceJobStatus::Failed));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
||||
@@ -232,8 +232,10 @@ The smoke test also probes catalog-backed advanced Iceberg surfaces:
|
||||
drop routes with persisted view metadata and view-scoped authorization
|
||||
- metadata maintenance supports safe dry-run planning, controlled worker
|
||||
execution checks, and a scheduler status report with disabled, paused,
|
||||
backpressure, retry, quarantine, and audit-timeline state; quarantined jobs
|
||||
can be inspected, released, retried, or abandoned through an operator endpoint
|
||||
backpressure, retry, quarantine, and audit-timeline state; job reports expose
|
||||
structured audit events for planning, worker transitions, heartbeat updates,
|
||||
lease expiry, and mutating quarantine operations; quarantined jobs can be
|
||||
inspected, released, retried, or abandoned through an operator endpoint
|
||||
- catalog diagnostics exposes the table recovery and consistency state used by
|
||||
operators
|
||||
- catalog export and diagnostics expose the current catalog backing manifest,
|
||||
@@ -336,7 +338,7 @@ Unsupported behavior is documented instead of hidden behind internal errors. The
|
||||
current unsupported inventory is:
|
||||
|
||||
- credential vending: automated after table bootstrap with exact-prefix validation and a data-plane scope probe; full no-long-term-data-credential bootstrap is not claimed
|
||||
- background maintenance worker: controlled run-once, heartbeat, quarantine operation, and scheduler status endpoints are registered; disabled/paused/backpressure/retry/quarantine/audit-timeline state is machine-readable; continuous in-process scheduling is not claimed
|
||||
- background maintenance worker: controlled run-once, heartbeat, quarantine operation, and scheduler status endpoints are registered; disabled/paused/backpressure/retry/quarantine/audit-timeline state and per-job audit events are machine-readable; continuous in-process scheduling is not claimed
|
||||
- manifest/data reachability cleanup: metadata maintenance reads manifest-list and manifest Avro references, reports manifest/data/delete reachability, and deletes only unreferenced table objects that pass the safety window
|
||||
- snapshot expiration dry-run planning and manual catalog commit: supported through metadata maintenance reports
|
||||
- automatic maintenance scheduling: external scheduler hook supported through the worker run endpoint and scheduler status report; built-in periodic scheduling is not claimed
|
||||
|
||||
@@ -183,7 +183,7 @@ UNSUPPORTED_INVENTORY: list[dict[str, str]] = [
|
||||
"status": "controlled-run-once-supported",
|
||||
"roadmap_area": "maintenance-worker",
|
||||
"catalog_endpoint": "GET /v1/{prefix}/namespaces/{namespace}/tables/{table}/maintenance/scheduler",
|
||||
"expected_behavior": "background-enabled maintenance can be driven by the worker run endpoint and inspected through the scheduler status endpoint with disabled/paused state, current-job backpressure, retry deferral, quarantine boundary, audit timeline, lease expiry recovery, heartbeat updates, and operator quarantine inspect/release/retry/abandon actions; built-in periodic scheduling is not claimed",
|
||||
"expected_behavior": "background-enabled maintenance can be driven by the worker run endpoint and inspected through the scheduler status endpoint with disabled/paused state, current-job backpressure, retry deferral, quarantine boundary, audit timeline, per-job audit events, lease expiry recovery, heartbeat updates, and operator quarantine inspect/release/retry/abandon actions; built-in periodic scheduling is not claimed",
|
||||
},
|
||||
{
|
||||
"capability": "manifest-data-reachability-cleanup",
|
||||
@@ -1065,7 +1065,14 @@ def run_maintenance_probe(args: argparse.Namespace, deps: RuntimeDeps) -> None:
|
||||
job = report.get("job")
|
||||
if not isinstance(job, dict) or not job.get("job-id"):
|
||||
raise RuntimeError("maintenance metadata endpoint did not return a job id")
|
||||
signed_rest_request(args, deps, "GET", table_endpoint_path(args, f"/maintenance/jobs/{job['job-id']}"))
|
||||
audit_events = report.get("audit-events")
|
||||
if not isinstance(audit_events, list) or not audit_events:
|
||||
raise RuntimeError("maintenance metadata endpoint did not return audit events")
|
||||
if not any(isinstance(event, dict) and event.get("action") == "PLANNED" for event in audit_events):
|
||||
raise RuntimeError("maintenance metadata endpoint did not report a planning audit event")
|
||||
job_report = signed_rest_request(args, deps, "GET", table_endpoint_path(args, f"/maintenance/jobs/{job['job-id']}"))
|
||||
if not isinstance(job_report.get("audit-events"), list):
|
||||
raise RuntimeError("maintenance job endpoint did not return audit events")
|
||||
quarantine = signed_rest_request(
|
||||
args,
|
||||
deps,
|
||||
@@ -1079,14 +1086,19 @@ def run_maintenance_probe(args: argparse.Namespace, deps: RuntimeDeps) -> None:
|
||||
expected_scheduler_statuses = {"READY", "DISABLED", "PAUSED", "BACKPRESSURED", "RETRY_DEFERRED", "QUARANTINED"}
|
||||
if scheduler.get("status") not in expected_scheduler_statuses:
|
||||
raise RuntimeError("maintenance scheduler endpoint did not return a stable scheduler status")
|
||||
if "audit_timeline" not in scheduler:
|
||||
scheduler_timeline = scheduler.get("audit_timeline")
|
||||
if not isinstance(scheduler_timeline, list):
|
||||
raise RuntimeError("maintenance scheduler endpoint did not return an audit timeline")
|
||||
if scheduler_timeline and not isinstance(scheduler_timeline[0].get("audit-events"), list):
|
||||
raise RuntimeError("maintenance scheduler audit timeline did not include job audit events")
|
||||
|
||||
worker = signed_rest_request(args, deps, "POST", table_endpoint_path(args, "/maintenance/worker/run"), {})
|
||||
worker_job = worker.get("job")
|
||||
expected_statuses = {"DISABLED", "PAUSED", "FAILED", "SUCCESSFUL", "RUNNING"}
|
||||
if not isinstance(worker_job, dict) or worker_job.get("status") not in expected_statuses:
|
||||
raise RuntimeError("maintenance worker endpoint did not return a stable job status")
|
||||
if not isinstance(worker.get("audit-events"), list):
|
||||
raise RuntimeError("maintenance worker endpoint did not return audit events")
|
||||
|
||||
|
||||
def run_catalog_api_probes(args: argparse.Namespace, deps: RuntimeDeps) -> None:
|
||||
|
||||
@@ -310,15 +310,15 @@ class PyIcebergSmokeConfigTest(unittest.TestCase):
|
||||
if (method, path) == ("GET", config_path):
|
||||
return {"version": 1}
|
||||
if (method, path) == ("POST", maintenance_path):
|
||||
return {"job": {"job-id": "job-1"}}
|
||||
return {"job": {"job-id": "job-1"}, "audit-events": [{"action": "PLANNED"}]}
|
||||
if (method, path) == ("GET", job_path):
|
||||
return {"job": {"job-id": "job-1", "status": "SUCCESSFUL"}}
|
||||
return {"job": {"job-id": "job-1", "status": "SUCCESSFUL"}, "audit-events": [{"action": "PLANNED"}]}
|
||||
if (method, path) == ("POST", quarantine_path):
|
||||
return {"action": "INSPECT", "report": {"job": {"job-id": "job-1"}}}
|
||||
if (method, path) == ("GET", scheduler_path):
|
||||
return {"status": "DISABLED", "audit_timeline": [{"job_id": "job-1"}]}
|
||||
return {"status": "DISABLED", "audit_timeline": [{"job_id": "job-1", "audit-events": [{"action": "PLANNED"}]}]}
|
||||
if (method, path) == ("POST", worker_path):
|
||||
return {"job": {"status": "UNKNOWN"}}
|
||||
return {"job": {"status": "UNKNOWN"}, "audit-events": [{"action": "WORKER_CONTROL"}]}
|
||||
raise AssertionError(f"unexpected REST request: {method} {path}")
|
||||
|
||||
with mock.patch.object(pyiceberg_smoke, "signed_rest_request", side_effect=fake_signed_request):
|
||||
@@ -374,15 +374,15 @@ class PyIcebergSmokeConfigTest(unittest.TestCase):
|
||||
if (method, path) == ("GET", config_path):
|
||||
return {"version": 1}
|
||||
if (method, path) == ("POST", maintenance_path):
|
||||
return {"job": {"job-id": "job-1"}}
|
||||
return {"job": {"job-id": "job-1"}, "audit-events": [{"action": "PLANNED"}]}
|
||||
if (method, path) == ("GET", pyiceberg_smoke.table_endpoint_path(args, "/maintenance/jobs/job-1")):
|
||||
return {"job": {"job-id": "job-1", "status": "SUCCESSFUL"}}
|
||||
return {"job": {"job-id": "job-1", "status": "SUCCESSFUL"}, "audit-events": [{"action": "PLANNED"}]}
|
||||
if (method, path) == ("POST", quarantine_path):
|
||||
return {"action": "INSPECT", "report": {"job": {"job-id": "job-1"}}}
|
||||
if (method, path) == ("GET", scheduler_path):
|
||||
return {"status": "DISABLED", "audit_timeline": [{"job_id": "job-1"}]}
|
||||
return {"status": "DISABLED", "audit_timeline": [{"job_id": "job-1", "audit-events": [{"action": "PLANNED"}]}]}
|
||||
if (method, path) == ("POST", worker_path):
|
||||
return {"job": {"status": "PAUSED"}}
|
||||
return {"job": {"status": "PAUSED"}, "audit-events": [{"action": "WORKER_CONTROL"}]}
|
||||
if (method, path) == ("GET", diagnostics_path):
|
||||
return {"status": "ok"}
|
||||
raise AssertionError(f"unexpected REST request: {method} {path}")
|
||||
|
||||
Reference in New Issue
Block a user