diff --git a/docs/architecture/s3-tables-support-matrix.md b/docs/architecture/s3-tables-support-matrix.md index 5b818650a..0c561223d 100644 --- a/docs/architecture/s3-tables-support-matrix.md +++ b/docs/architecture/s3-tables-support-matrix.md @@ -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. | diff --git a/rustfs/src/table_catalog.rs b/rustfs/src/table_catalog.rs index 525ae7be6..4355d6829 100644 --- a/rustfs/src/table_catalog.rs +++ b/rustfs/src/table_catalog.rs @@ -552,6 +552,48 @@ pub(crate) struct TableMaintenanceSchedulerQuarantineBoundary { pub source_job_id: Option, } +#[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, + #[serde(default, rename = "before-status")] + pub before_status: Option, + #[serde(default, rename = "after-status")] + pub after_status: Option, + #[serde(default, rename = "before-quarantined-object-count")] + pub before_quarantined_object_count: Option, + #[serde(default, rename = "after-quarantined-object-count")] + pub after_quarantined_object_count: Option, + #[serde(default, rename = "recommended-actions")] + pub recommended_actions: Vec, +} + #[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, pub next_retry_after: Option, pub recommended_actions: Vec, + #[serde(default, rename = "audit-events")] + pub audit_events: Vec, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] @@ -677,6 +721,8 @@ pub(crate) struct TableMetadataMaintenanceReport { pub snapshot_expiration: Option, #[serde(default)] pub compaction: Option, + #[serde(default, rename = "audit-events")] + pub audit_events: Vec, } #[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 { 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, + before_status: Option, + before_quarantined_object_count: Option, +) { + 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 { 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::>(); + 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] diff --git a/scripts/table-catalog/README.md b/scripts/table-catalog/README.md index 187e101e5..947590df8 100644 --- a/scripts/table-catalog/README.md +++ b/scripts/table-catalog/README.md @@ -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 diff --git a/scripts/table-catalog/pyiceberg_smoke.py b/scripts/table-catalog/pyiceberg_smoke.py index 064a077a5..73096c98a 100755 --- a/scripts/table-catalog/pyiceberg_smoke.py +++ b/scripts/table-catalog/pyiceberg_smoke.py @@ -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: diff --git a/scripts/table-catalog/test_pyiceberg_smoke.py b/scripts/table-catalog/test_pyiceberg_smoke.py index ca57306fe..f530c945b 100644 --- a/scripts/table-catalog/test_pyiceberg_smoke.py +++ b/scripts/table-catalog/test_pyiceberg_smoke.py @@ -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}")