From 941b8afee8619c6a038c20aac51565a65c959ee4 Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Mon, 15 Jun 2026 14:08:28 +0800 Subject: [PATCH] feat(table-catalog): add snapshot maintenance foundation (#3470) Co-authored-by: Henry Guo --- rustfs/src/admin/handlers/table_catalog.rs | 564 +++++++++++- rustfs/src/table_catalog.rs | 985 +++++++++++++++++++++ scripts/table-catalog/README.md | 4 +- scripts/table-catalog/pyiceberg_smoke.py | 12 +- 4 files changed, 1556 insertions(+), 9 deletions(-) diff --git a/rustfs/src/admin/handlers/table_catalog.rs b/rustfs/src/admin/handlers/table_catalog.rs index 9ed8e2a65..dff5cdb56 100644 --- a/rustfs/src/admin/handlers/table_catalog.rs +++ b/rustfs/src/admin/handlers/table_catalog.rs @@ -216,6 +216,12 @@ struct TableMetadataMaintenanceRequest { retain_recent_metadata_files: usize, #[serde(default)] delete: bool, + #[serde(default, rename = "snapshot-expiration")] + snapshot_expiration: Option, + #[serde(default, rename = "commit-snapshot-expiration")] + commit_snapshot_expiration: bool, + #[serde(default)] + compaction: Option, } #[derive(Debug, Deserialize)] @@ -2456,6 +2462,7 @@ where async fn table_metadata_maintenance_response( store: &crate::table_catalog::ObjectTableCatalogStore, + metadata_backend: &B, bucket: &str, namespace: &crate::table_catalog::Namespace, table: &str, @@ -2464,7 +2471,35 @@ async fn table_metadata_maintenance_response( where B: crate::table_catalog::TableCatalogObjectBackend, { - store + if request.delete && request.commit_snapshot_expiration { + return Err(s3_error!( + InvalidRequest, + "snapshot expiration commit cannot be combined with metadata deletion" + )); + } + + let snapshot_expiration_request = request.snapshot_expiration; + let commit_snapshot_expiration = request.commit_snapshot_expiration; + let compaction_request = request.compaction; + let compaction = match compaction_request { + Some(config) => Some( + store + .plan_table_compaction(bucket, &namespace.public_name(), table, config) + .await + .map_err(catalog_store_error)?, + ), + None => None, + }; + let snapshot_expiration_plan = match snapshot_expiration_request { + Some(config) => Some( + store + .plan_table_snapshot_expiration(bucket, &namespace.public_name(), table, config) + .await + .map_err(catalog_store_error)?, + ), + None => None, + }; + let mut report = store .run_table_metadata_maintenance_with_retention( bucket, &namespace.public_name(), @@ -2474,7 +2509,118 @@ where request.retain_recent_metadata_files, ) .await - .map_err(catalog_store_error) + .map_err(catalog_store_error)?; + let snapshot_expiration = match (snapshot_expiration_plan, commit_snapshot_expiration) { + (Some(plan), true) => { + Some(commit_table_snapshot_expiration_response(store, metadata_backend, bucket, namespace, table, plan).await?) + } + (Some(plan), false) => Some(plan), + (None, _) => None, + }; + report.snapshot_expiration = snapshot_expiration; + report.compaction = compaction; + if report.snapshot_expiration.is_some() || report.compaction.is_some() { + let committed_snapshot_expiration = report + .snapshot_expiration + .as_ref() + .is_some_and(|snapshot_expiration| snapshot_expiration.committed_metadata_location.is_some()); + match store.put_table_metadata_maintenance_report(&report).await { + Ok(()) => {} + Err(err) if committed_snapshot_expiration => { + tracing::warn!( + error = %err, + warehouse = bucket, + namespace = namespace.public_name(), + table, + "failed to persist table maintenance report after snapshot expiration commit" + ); + } + Err(err) => return Err(catalog_store_error(err)), + } + } + Ok(report) +} + +async fn commit_table_snapshot_expiration_response( + store: &crate::table_catalog::ObjectTableCatalogStore, + metadata_backend: &B, + bucket: &str, + namespace: &crate::table_catalog::Namespace, + table: &str, + mut report: crate::table_catalog::TableSnapshotExpirationReport, +) -> S3Result +where + B: crate::table_catalog::TableCatalogObjectBackend, +{ + let Some(current) = store + .load_table(bucket, &namespace.public_name(), table) + .await + .map_err(catalog_store_error)? + else { + return Err(s3_error!(InvalidRequest, "table not found")); + }; + if report.table_id != current.table_id || report.current_metadata_location != current.metadata_location { + return Err(s3_error!(PreconditionFailed, "snapshot expiration plan is stale")); + } + if report.manual_review_count > 0 { + return Err(s3_error!(InvalidRequest, "snapshot expiration plan requires manual review before commit")); + } + let expired_snapshot_ids = report + .snapshot_reports + .iter() + .filter(|snapshot| snapshot.state == crate::table_catalog::TableSnapshotExpirationSnapshotState::ExpirationCandidate) + .filter_map(|snapshot| snapshot.snapshot_id) + .collect::>(); + if expired_snapshot_ids.is_empty() { + return Ok(report); + } + + let current_metadata = read_table_metadata_json(metadata_backend, bucket, ¤t.metadata_location).await?; + let updates = [serde_json::json!({ + "action": "remove-snapshots", + "snapshot-ids": expired_snapshot_ids.clone() + })]; + let next_metadata = apply_table_commit_updates(current_metadata.clone(), &updates, ¤t.metadata_location)?; + validate_metadata_matches_current_metadata(¤t_metadata, &next_metadata)?; + validate_metadata_table_location_in_bucket(bucket, &next_metadata)?; + let table_name = crate::table_catalog::IdentifierSegment::parse(table.to_string()) + .map_err(|err| s3_error!(InvalidRequest, "invalid table name: {}", err))?; + let (commit_id, metadata_file_token) = standard_commit_ids(None); + let next_generation = current.generation.saturating_add(1); + let next_metadata_location = crate::table_catalog::default_table_metadata_file_path( + namespace, + &table_name, + &next_metadata_file_name(next_generation, &metadata_file_token), + ); + let next_metadata_data = serde_json::to_vec(&next_metadata) + .map_err(|err| s3_error!(InternalError, "failed to serialize snapshot expiration metadata: {}", err))?; + metadata_backend + .put_object( + bucket, + &next_metadata_location, + next_metadata_data, + crate::table_catalog::TableCatalogPutPrecondition::IfAbsent, + ) + .await + .map_err(catalog_store_error)?; + + let commit_request = crate::table_catalog::TableCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.to_string(), + commit_id, + idempotency_key: None, + operation: "expire-snapshots".to_string(), + expected_version_token: current.version_token, + expected_metadata_location: current.metadata_location, + new_metadata_location: next_metadata_location, + requirements: Vec::new(), + writer: Some("rustfs-maintenance".to_string()), + }; + let result = store.commit_table(commit_request).await.map_err(catalog_store_error)?; + report.expired_snapshot_ids = expired_snapshot_ids; + report.committed_metadata_location = Some(result.table.metadata_location); + Ok(report) } async fn catalog_import_response( @@ -2868,8 +3014,9 @@ impl Operation for RestTableMetadataMaintenanceHandler { authorize_table_catalog_resource_request(&req, &resource, AdminAction::RunTableMaintenanceAction).await?; 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?; + let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone()); + let response = + table_metadata_maintenance_response(&store, &metadata_backend, &warehouse, &namespace, &table, request).await?; build_json_response(StatusCode::OK, &response) } } @@ -3315,6 +3462,9 @@ mod tests { assert_eq!(request.retain_recent_metadata_files, 0); assert!(!request.delete); + assert!(request.snapshot_expiration.is_none()); + assert!(!request.commit_snapshot_expiration); + assert!(request.compaction.is_none()); } #[test] @@ -3327,6 +3477,39 @@ mod tests { assert_eq!(request.retain_recent_metadata_files, 2); assert!(request.delete); + assert!(request.snapshot_expiration.is_none()); + assert!(!request.commit_snapshot_expiration); + assert!(request.compaction.is_none()); + } + + #[test] + fn table_metadata_maintenance_request_accepts_snapshot_and_compaction_plans() { + let request: TableMetadataMaintenanceRequest = serde_json::from_value(serde_json::json!({ + "commit-snapshot-expiration": true, + "snapshot-expiration": { + "min-snapshots-to-keep": 2, + "max-snapshot-age-ms": 3600000 + }, + "compaction": { + "target-file-size-bytes": 536870912, + "small-file-threshold-bytes": 67108864, + "min-input-files": 5, + "max-rewrite-bytes-per-job": 10737418240u64 + } + })) + .expect("metadata maintenance request should parse maintenance planning config"); + + let snapshot_expiration = request + .snapshot_expiration + .expect("snapshot expiration config should be present"); + assert_eq!(snapshot_expiration.min_snapshots_to_keep, 2); + assert_eq!(snapshot_expiration.max_snapshot_age_ms, 3_600_000); + assert!(request.commit_snapshot_expiration); + let compaction = request.compaction.expect("compaction config should be present"); + assert_eq!(compaction.target_file_size_bytes, 536_870_912); + assert_eq!(compaction.small_file_threshold_bytes, 67_108_864); + assert_eq!(compaction.min_input_files, 5); + assert_eq!(compaction.max_rewrite_bytes_per_job, 10_737_418_240); } #[tokio::test] @@ -4184,7 +4367,26 @@ mod tests { bucket, ¤t, serde_json::json!({ - "metadata-log": [] + "current-snapshot-id": 20, + "metadata-log": [], + "snapshots": [ + { + "snapshot-id": 10, + "timestamp-ms": 1000, + "manifest-list": "s3://warehouse/tables/table-id/metadata/snap-10.avro" + }, + { + "snapshot-id": 20, + "timestamp-ms": 2000, + "manifest-list": "s3://warehouse/tables/table-id/metadata/snap-20.avro" + } + ], + "refs": { + "main": { + "snapshot-id": 20, + "type": "branch" + } + } }), Some(OffsetDateTime::UNIX_EPOCH), ) @@ -4230,18 +4432,45 @@ mod tests { let dry_run = table_metadata_maintenance_response( &store, + &backend, bucket, &namespace, "events", TableMetadataMaintenanceRequest { retain_recent_metadata_files: 0, delete: false, + snapshot_expiration: Some(crate::table_catalog::TableSnapshotExpirationConfig { + min_snapshots_to_keep: 1, + max_snapshot_age_ms: 1, + }), + commit_snapshot_expiration: false, + compaction: Some(crate::table_catalog::TableCompactionPlanningConfig { + target_file_size_bytes: 512 * 1024 * 1024, + small_file_threshold_bytes: 64 * 1024 * 1024, + min_input_files: 2, + max_rewrite_bytes_per_job: 1024 * 1024 * 1024, + }), }, ) .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()]); + let snapshot_expiration = dry_run + .snapshot_expiration + .as_ref() + .expect("dry-run report should include snapshot expiration planning"); + assert_eq!(snapshot_expiration.expiration_candidate_count, 1); + assert_eq!(snapshot_expiration.current_snapshot_id, Some(20)); + let compaction = dry_run + .compaction + .as_ref() + .expect("dry-run report should include compaction planning"); + assert_eq!( + compaction.status, + crate::table_catalog::TableCompactionPlanningStatus::ManualReviewRequired + ); + assert_eq!(compaction.manual_review_count, 1); let stored_dry_run = store .get_table_metadata_maintenance_report(bucket, "analytics", "events", &dry_run.job.job_id) .await @@ -4257,12 +4486,16 @@ mod tests { let deleted = table_metadata_maintenance_response( &store, + &backend, bucket, &namespace, "events", TableMetadataMaintenanceRequest { retain_recent_metadata_files: 0, delete: true, + snapshot_expiration: None, + commit_snapshot_expiration: false, + compaction: None, }, ) .await @@ -4277,6 +4510,327 @@ mod tests { ); } + #[tokio::test] + async fn table_metadata_maintenance_helper_commits_snapshot_expiration() { + 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 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, + ¤t, + serde_json::json!({ + "format-version": 2, + "table-uuid": "table-uuid", + "location": "s3://warehouse/tables/table-id", + "last-sequence-number": 2, + "last-updated-ms": 2000, + "last-column-id": 1, + "schemas": [], + "current-schema-id": 0, + "partition-specs": [], + "default-spec-id": 0, + "sort-orders": [], + "default-sort-order-id": 0, + "current-snapshot-id": 20, + "metadata-log": [], + "snapshot-log": [ + { + "timestamp-ms": 1000, + "snapshot-id": 10 + }, + { + "timestamp-ms": 2000, + "snapshot-id": 20 + } + ], + "snapshots": [ + { + "snapshot-id": 10, + "timestamp-ms": 1000, + "manifest-list": "s3://warehouse/tables/table-id/metadata/snap-10.avro" + }, + { + "snapshot-id": 20, + "timestamp-ms": 2000, + "manifest-list": "s3://warehouse/tables/table-id/metadata/snap-20.avro" + } + ], + "refs": { + "main": { + "snapshot-id": 20, + "type": "branch" + } + } + }), + Some(OffsetDateTime::UNIX_EPOCH), + ) + .await; + + let report = table_metadata_maintenance_response( + &store, + &backend, + bucket, + &namespace, + "events", + TableMetadataMaintenanceRequest { + retain_recent_metadata_files: 0, + delete: false, + snapshot_expiration: Some(crate::table_catalog::TableSnapshotExpirationConfig { + min_snapshots_to_keep: 1, + max_snapshot_age_ms: 1, + }), + commit_snapshot_expiration: true, + compaction: None, + }, + ) + .await + .expect("snapshot expiration commit should succeed"); + + let snapshot_expiration = report + .snapshot_expiration + .as_ref() + .expect("maintenance report should include snapshot expiration"); + assert_eq!(snapshot_expiration.expired_snapshot_ids, vec![10]); + let committed_location = snapshot_expiration + .committed_metadata_location + .as_ref() + .expect("snapshot expiration commit should report committed metadata") + .clone(); + assert_ne!(committed_location, current); + + let entry = store + .load_table(bucket, "analytics", "events") + .await + .expect("table lookup should succeed") + .expect("table should exist"); + assert_eq!(entry.metadata_location, committed_location); + assert_eq!(entry.generation, 2); + + let committed_object = backend + .read_object(bucket, &entry.metadata_location) + .await + .expect("committed metadata lookup should succeed") + .expect("committed metadata object should exist"); + let committed_metadata = + serde_json::from_slice::(&committed_object.data).expect("committed metadata should be valid JSON"); + let snapshots = committed_metadata + .get("snapshots") + .and_then(serde_json::Value::as_array) + .expect("committed metadata should contain snapshots"); + assert_eq!(snapshots.len(), 1); + assert_eq!(snapshots[0].get("snapshot-id").and_then(serde_json::Value::as_i64), Some(20)); + let snapshot_log = committed_metadata + .get("snapshot-log") + .and_then(serde_json::Value::as_array) + .expect("committed metadata should contain snapshot-log"); + assert_eq!(snapshot_log.len(), 1); + assert_eq!( + committed_metadata["metadata-log"][0]["metadata-file"], + serde_json::Value::String(current.clone()) + ); + assert!( + backend + .object_exists(bucket, ¤t) + .await + .expect("previous metadata lookup should succeed") + ); + } + + #[tokio::test] + async fn table_metadata_maintenance_helper_rejects_snapshot_expiration_manual_review_commit() { + 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 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, + ¤t, + serde_json::json!({ + "format-version": 2, + "table-uuid": "table-uuid", + "location": "s3://warehouse/tables/table-id", + "current-snapshot-id": 20, + "metadata-log": [], + "snapshot-log": [], + "snapshots": [ + { + "snapshot-id": 10, + "timestamp-ms": 1000, + "manifest-list": "s3://warehouse/tables/table-id/metadata/snap-10.avro" + }, + { + "snapshot-id": 20, + "timestamp-ms": 2000, + "manifest-list": "s3://warehouse/tables/table-id/metadata/snap-20.avro" + } + ], + "refs": { + "main": { + "snapshot-id": 20, + "type": "branch" + }, + "audit": { + "snapshot-id": 10, + "type": "tag" + } + } + }), + Some(OffsetDateTime::UNIX_EPOCH), + ) + .await; + + let result = table_metadata_maintenance_response( + &store, + &backend, + bucket, + &namespace, + "events", + TableMetadataMaintenanceRequest { + retain_recent_metadata_files: 0, + delete: false, + snapshot_expiration: Some(crate::table_catalog::TableSnapshotExpirationConfig { + min_snapshots_to_keep: 1, + max_snapshot_age_ms: 1, + }), + commit_snapshot_expiration: true, + compaction: None, + }, + ) + .await; + + assert!(result.is_err()); + let entry = store + .load_table(bucket, "analytics", "events") + .await + .expect("table lookup should succeed") + .expect("table should exist"); + assert_eq!(entry.metadata_location, current); + assert_eq!(entry.generation, 1); + } + + #[tokio::test] + async fn table_metadata_maintenance_helper_rejects_stale_snapshot_expiration_plan() { + 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 current = crate::table_catalog::default_table_metadata_file_path(&namespace, &table, "00002.metadata.json"); + let next = crate::table_catalog::default_table_metadata_file_path(&namespace, &table, "00003.metadata.json"); + + let metadata = serde_json::json!({ + "format-version": 2, + "table-uuid": "table-uuid", + "location": "s3://warehouse/tables/table-id", + "current-snapshot-id": 20, + "metadata-log": [], + "snapshot-log": [], + "snapshots": [ + { + "snapshot-id": 10, + "timestamp-ms": 1000, + "manifest-list": "s3://warehouse/tables/table-id/metadata/snap-10.avro" + }, + { + "snapshot-id": 20, + "timestamp-ms": 2000, + "manifest-list": "s3://warehouse/tables/table-id/metadata/snap-20.avro" + } + ], + "refs": { + "main": { + "snapshot-id": 20, + "type": "branch" + } + } + }); + seed_object_table_for_metadata_maintenance(&store, &backend, bucket, &namespace, &table, current.clone()).await; + backend + .put_json_with_mod_time(bucket, ¤t, metadata.clone(), Some(OffsetDateTime::UNIX_EPOCH)) + .await; + backend + .put_json_with_mod_time(bucket, &next, metadata, Some(OffsetDateTime::UNIX_EPOCH)) + .await; + + let stale_plan = store + .plan_table_snapshot_expiration( + bucket, + "analytics", + "events", + crate::table_catalog::TableSnapshotExpirationConfig { + min_snapshots_to_keep: 1, + max_snapshot_age_ms: 1, + }, + ) + .await + .expect("snapshot expiration plan should build"); + store + .commit_table(crate::table_catalog::TableCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + commit_id: "advance-pointer".to_string(), + idempotency_key: None, + operation: "append".to_string(), + expected_version_token: "token-v1".to_string(), + expected_metadata_location: current, + new_metadata_location: next, + requirements: Vec::new(), + writer: Some("test".to_string()), + }) + .await + .expect("pointer advance should succeed"); + + let result = commit_table_snapshot_expiration_response(&store, &backend, bucket, &namespace, "events", stale_plan).await; + + assert!(result.is_err()); + let entry = store + .load_table(bucket, "analytics", "events") + .await + .expect("table lookup should succeed") + .expect("table should exist"); + assert_eq!(entry.generation, 2); + } + + #[tokio::test] + async fn table_metadata_maintenance_helper_rejects_delete_with_snapshot_expiration_commit() { + let backend = TestTableCatalogObjectBackend::default(); + let store = crate::table_catalog::ObjectTableCatalogStore::new(backend.clone()); + let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse"); + + let result = table_metadata_maintenance_response( + &store, + &backend, + "warehouse", + &namespace, + "events", + TableMetadataMaintenanceRequest { + retain_recent_metadata_files: 0, + delete: true, + snapshot_expiration: Some(crate::table_catalog::TableSnapshotExpirationConfig { + min_snapshots_to_keep: 1, + max_snapshot_age_ms: 1, + }), + commit_snapshot_expiration: true, + compaction: None, + }, + ) + .await; + + assert!(result.is_err()); + } + #[test] fn commit_requirements_reject_mismatched_table_uuid() { let metadata = serde_json::json!({ diff --git a/rustfs/src/table_catalog.rs b/rustfs/src/table_catalog.rs index 1be3ace19..aaceba812 100644 --- a/rustfs/src/table_catalog.rs +++ b/rustfs/src/table_catalog.rs @@ -84,6 +84,13 @@ const MAINTENANCE_JOB_ALIAS_CURRENT: &str = "current"; const TABLE_CATALOG_LIST_MAX_KEYS: i32 = 1000; const TABLE_METADATA_CLEANUP_SAFETY_WINDOW_SECONDS: i64 = 15 * 60; const TABLE_COMMIT_SLOW_LOG_THRESHOLD: StdDuration = StdDuration::from_secs(2); +const ICEBERG_MAIN_REF: &str = "main"; +const ICEBERG_MIN_SNAPSHOTS_TO_KEEP_PROPERTY: &str = "history.expire.min-snapshots-to-keep"; +const ICEBERG_MAX_SNAPSHOT_AGE_MS_PROPERTY: &str = "history.expire.max-snapshot-age-ms"; +const ICEBERG_MAX_REF_AGE_MS_PROPERTY: &str = "history.expire.max-ref-age-ms"; +const ICEBERG_REF_MIN_SNAPSHOTS_TO_KEEP_FIELD: &str = "min-snapshots-to-keep"; +const ICEBERG_REF_MAX_SNAPSHOT_AGE_MS_FIELD: &str = "max-snapshot-age-ms"; +const ICEBERG_REF_MAX_REF_AGE_MS_FIELD: &str = "max-ref-age-ms"; #[derive(Debug, Clone, PartialEq, Eq)] pub enum CatalogIdentifierError { @@ -348,6 +355,125 @@ pub(crate) struct TableMetadataMaintenanceReport { pub object_reports: Vec, #[serde(default)] pub referenced_object_reports: Vec, + #[serde(default, rename = "snapshot-expiration")] + pub snapshot_expiration: Option, + #[serde(default)] + pub compaction: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct TableSnapshotExpirationConfig { + #[serde(rename = "min-snapshots-to-keep")] + pub min_snapshots_to_keep: usize, + #[serde(rename = "max-snapshot-age-ms")] + pub max_snapshot_age_ms: i64, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub(crate) struct TableSnapshotExpirationReport { + pub table_bucket: String, + pub namespace: String, + pub table: String, + pub table_id: String, + pub current_metadata_location: String, + pub current_snapshot_id: Option, + pub config: TableSnapshotExpirationConfig, + pub expiration_watermark_ms: i64, + pub retained_snapshot_count: usize, + pub expiration_candidate_count: usize, + pub manual_review_count: usize, + #[serde(default)] + pub expired_snapshot_ids: Vec, + #[serde(default)] + pub committed_metadata_location: Option, + pub snapshot_reports: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub(crate) struct TableSnapshotExpirationSnapshotReport { + pub snapshot_id: Option, + pub sequence_number: Option, + pub timestamp_ms: Option, + pub manifest_list: Option, + pub state: TableSnapshotExpirationSnapshotState, + pub reasons: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "SCREAMING_SNAKE_CASE")] +pub(crate) enum TableSnapshotExpirationSnapshotState { + Retained, + ExpirationCandidate, + ManualReviewRequired, +} + +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)] +#[serde(rename_all = "SCREAMING_SNAKE_CASE")] +pub(crate) enum TableSnapshotExpirationReason { + CurrentSnapshot, + MinSnapshotsToKeep, + ProtectedSnapshotRef, + UserDefinedSnapshotRef, + SnapshotRefRetentionConflict, + TableRetentionPropertyConflict, + MissingSnapshotId, + MissingSnapshotTimestamp, + SnapshotAgeWithinRetention, + SnapshotAgeExpired, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct TableCompactionPlanningConfig { + #[serde(rename = "target-file-size-bytes")] + pub target_file_size_bytes: u64, + #[serde(rename = "small-file-threshold-bytes")] + pub small_file_threshold_bytes: u64, + #[serde(rename = "min-input-files")] + pub min_input_files: usize, + #[serde(rename = "max-rewrite-bytes-per-job")] + pub max_rewrite_bytes_per_job: u64, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub(crate) struct TableCompactionPlanningReport { + pub table_bucket: String, + pub namespace: String, + pub table: String, + pub table_id: String, + pub current_metadata_location: String, + pub current_snapshot_id: Option, + pub config: TableCompactionPlanningConfig, + pub status: TableCompactionPlanningStatus, + pub candidate_file_count: usize, + pub rewrite_group_count: usize, + pub manual_review_count: usize, + pub snapshot_reports: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub(crate) struct TableCompactionSnapshotReport { + pub snapshot_id: Option, + pub manifest_list: Option, + pub status: TableCompactionPlanningStatus, + pub reasons: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "SCREAMING_SNAKE_CASE")] +pub(crate) enum TableCompactionPlanningStatus { + NoCandidates, + ManualReviewRequired, +} + +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)] +#[serde(rename_all = "SCREAMING_SNAKE_CASE")] +pub(crate) enum TableCompactionPlanningReason { + ManifestList, + ManifestAvroReaderUnavailable, + MissingCurrentSnapshot, + MissingManifestList, } #[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] @@ -1403,6 +1529,109 @@ where .map(|entry| entry.map(|(report, _)| report)) } + pub(crate) async fn plan_table_snapshot_expiration( + &self, + table_bucket: &str, + namespace: &str, + table: &str, + config: TableSnapshotExpirationConfig, + ) -> TableCatalogStoreResult { + validate_table_snapshot_expiration_config(&config)?; + 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 Some((entry, _)) = self.read_entry::(self.catalog_bucket(), &table_path).await? else { + return Err(TableCatalogStoreError::NotFound(format!( + "table {}/{}/{}", + table_bucket, + namespace.public_name(), + table.as_str() + ))); + }; + if !is_valid_table_metadata_location(&namespace, &table, &entry.metadata_location) { + return Err(TableCatalogStoreError::Invalid( + "current metadata location must be inside the table metadata directory".to_string(), + )); + } + + let Some(current_metadata_object) = self.backend.read_object(table_bucket, &entry.metadata_location).await? else { + return Err(TableCatalogStoreError::NotFound(format!( + "current metadata object {}", + entry.metadata_location + ))); + }; + let current_metadata = serde_json::from_slice::(¤t_metadata_object.data).map_err(|err| { + TableCatalogStoreError::Invalid(format!("failed to parse current metadata {}: {err}", entry.metadata_location)) + })?; + if !current_metadata.is_object() { + return Err(TableCatalogStoreError::Invalid(format!( + "current metadata {} must be a JSON object", + entry.metadata_location + ))); + } + + Ok(table_snapshot_expiration_report( + table_bucket, + &namespace, + &table, + &entry, + ¤t_metadata, + config, + OffsetDateTime::now_utc(), + )) + } + + pub(crate) async fn plan_table_compaction( + &self, + table_bucket: &str, + namespace: &str, + table: &str, + config: TableCompactionPlanningConfig, + ) -> TableCatalogStoreResult { + validate_table_compaction_planning_config(&config)?; + 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 Some((entry, _)) = self.read_entry::(self.catalog_bucket(), &table_path).await? else { + return Err(TableCatalogStoreError::NotFound(format!( + "table {}/{}/{}", + table_bucket, + namespace.public_name(), + table.as_str() + ))); + }; + if !is_valid_table_metadata_location(&namespace, &table, &entry.metadata_location) { + return Err(TableCatalogStoreError::Invalid( + "current metadata location must be inside the table metadata directory".to_string(), + )); + } + + let Some(current_metadata_object) = self.backend.read_object(table_bucket, &entry.metadata_location).await? else { + return Err(TableCatalogStoreError::NotFound(format!( + "current metadata object {}", + entry.metadata_location + ))); + }; + let current_metadata = serde_json::from_slice::(¤t_metadata_object.data).map_err(|err| { + TableCatalogStoreError::Invalid(format!("failed to parse current metadata {}: {err}", entry.metadata_location)) + })?; + if !current_metadata.is_object() { + return Err(TableCatalogStoreError::Invalid(format!( + "current metadata {} must be a JSON object", + entry.metadata_location + ))); + } + + Ok(table_compaction_planning_report( + table_bucket, + &namespace, + &table, + &entry, + ¤t_metadata, + config, + )) + } + pub(crate) async fn export_table_catalog_entry( &self, table_bucket: &str, @@ -1699,6 +1928,8 @@ where deletable_metadata_locations, object_reports, referenced_object_reports, + snapshot_expiration: None, + compaction: None, }) } @@ -1928,6 +2159,8 @@ where deletable_metadata_locations: cleanup_candidate_locations, object_reports, referenced_object_reports, + snapshot_expiration: report.snapshot_expiration, + compaction: report.compaction, }) } } @@ -2726,6 +2959,357 @@ fn metadata_candidate_is_past_safety_window(mod_time: Option, no mod_time <= now - Duration::seconds(TABLE_METADATA_CLEANUP_SAFETY_WINDOW_SECONDS) } +#[derive(Debug)] +struct TableSnapshotExpirationDraft { + snapshot_id: Option, + sequence_number: Option, + timestamp_ms: Option, + manifest_list: Option, + reasons: BTreeSet, +} + +fn table_snapshot_expiration_report( + table_bucket: &str, + namespace: &Namespace, + table: &IdentifierSegment, + entry: &TableEntry, + current_metadata: &serde_json::Value, + config: TableSnapshotExpirationConfig, + now: OffsetDateTime, +) -> TableSnapshotExpirationReport { + let current_snapshot_id = current_metadata + .get("current-snapshot-id") + .and_then(serde_json::Value::as_i64); + let expiration_watermark_ms = unix_timestamp_millis(now).saturating_sub(config.max_snapshot_age_ms); + let (protected_ref_snapshot_ids, user_defined_ref_snapshot_ids, ref_retention_conflict_snapshot_ids) = + snapshot_expiration_ref_state(current_metadata, current_snapshot_id); + let table_retention_property_conflict = snapshot_expiration_table_property_conflicts(current_metadata, &config); + + let mut drafts = snapshot_expiration_drafts(current_metadata, current_snapshot_id); + mark_recent_snapshots_to_keep(&mut drafts, config.min_snapshots_to_keep); + + let mut snapshot_reports = Vec::with_capacity(drafts.len()); + for mut draft in drafts { + if let Some(snapshot_id) = draft.snapshot_id { + if protected_ref_snapshot_ids.contains(&snapshot_id) { + draft.reasons.insert(TableSnapshotExpirationReason::ProtectedSnapshotRef); + } + if user_defined_ref_snapshot_ids.contains(&snapshot_id) { + draft.reasons.insert(TableSnapshotExpirationReason::UserDefinedSnapshotRef); + } + if ref_retention_conflict_snapshot_ids.contains(&snapshot_id) { + draft + .reasons + .insert(TableSnapshotExpirationReason::SnapshotRefRetentionConflict); + } + } + if table_retention_property_conflict { + draft + .reasons + .insert(TableSnapshotExpirationReason::TableRetentionPropertyConflict); + } + + let state = if snapshot_expiration_requires_manual_review(&draft.reasons) { + TableSnapshotExpirationSnapshotState::ManualReviewRequired + } else if snapshot_expiration_is_retained(&draft.reasons) { + TableSnapshotExpirationSnapshotState::Retained + } else if let Some(timestamp_ms) = draft.timestamp_ms { + if timestamp_ms <= expiration_watermark_ms { + draft.reasons.insert(TableSnapshotExpirationReason::SnapshotAgeExpired); + TableSnapshotExpirationSnapshotState::ExpirationCandidate + } else { + draft + .reasons + .insert(TableSnapshotExpirationReason::SnapshotAgeWithinRetention); + TableSnapshotExpirationSnapshotState::Retained + } + } else { + draft.reasons.insert(TableSnapshotExpirationReason::MissingSnapshotTimestamp); + TableSnapshotExpirationSnapshotState::ManualReviewRequired + }; + + snapshot_reports.push(TableSnapshotExpirationSnapshotReport { + snapshot_id: draft.snapshot_id, + sequence_number: draft.sequence_number, + timestamp_ms: draft.timestamp_ms, + manifest_list: draft.manifest_list, + state, + reasons: draft.reasons.into_iter().collect(), + }); + } + + let retained_snapshot_count = snapshot_reports + .iter() + .filter(|snapshot| snapshot.state == TableSnapshotExpirationSnapshotState::Retained) + .count(); + let expiration_candidate_count = snapshot_reports + .iter() + .filter(|snapshot| snapshot.state == TableSnapshotExpirationSnapshotState::ExpirationCandidate) + .count(); + let manual_review_count = snapshot_reports + .iter() + .filter(|snapshot| snapshot.state == TableSnapshotExpirationSnapshotState::ManualReviewRequired) + .count(); + + TableSnapshotExpirationReport { + table_bucket: table_bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + table_id: entry.table_id.clone(), + current_metadata_location: entry.metadata_location.clone(), + current_snapshot_id, + config, + expiration_watermark_ms, + retained_snapshot_count, + expiration_candidate_count, + manual_review_count, + expired_snapshot_ids: Vec::new(), + committed_metadata_location: None, + snapshot_reports, + } +} + +fn table_compaction_planning_report( + table_bucket: &str, + namespace: &Namespace, + table: &IdentifierSegment, + entry: &TableEntry, + current_metadata: &serde_json::Value, + config: TableCompactionPlanningConfig, +) -> TableCompactionPlanningReport { + let current_snapshot_id = current_metadata + .get("current-snapshot-id") + .and_then(serde_json::Value::as_i64); + let snapshot_reports = match current_snapshot_id { + Some(current_snapshot_id) => { + let current_snapshot = current_metadata + .get("snapshots") + .and_then(serde_json::Value::as_array) + .into_iter() + .flatten() + .find(|snapshot| { + snapshot + .get("snapshot-id") + .and_then(serde_json::Value::as_i64) + .is_some_and(|snapshot_id| snapshot_id == current_snapshot_id) + }); + match current_snapshot { + Some(snapshot) => { + let manifest_list = snapshot + .get("manifest-list") + .and_then(serde_json::Value::as_str) + .map(ToString::to_string); + let mut reasons = BTreeSet::new(); + if manifest_list.is_some() { + reasons.insert(TableCompactionPlanningReason::ManifestList); + reasons.insert(TableCompactionPlanningReason::ManifestAvroReaderUnavailable); + } else { + reasons.insert(TableCompactionPlanningReason::MissingManifestList); + } + + vec![TableCompactionSnapshotReport { + snapshot_id: Some(current_snapshot_id), + manifest_list, + status: TableCompactionPlanningStatus::ManualReviewRequired, + reasons: reasons.into_iter().collect(), + }] + } + None => vec![TableCompactionSnapshotReport { + snapshot_id: Some(current_snapshot_id), + manifest_list: None, + status: TableCompactionPlanningStatus::ManualReviewRequired, + reasons: vec![TableCompactionPlanningReason::MissingCurrentSnapshot], + }], + } + } + None => Vec::new(), + }; + + let (status, manual_review_count) = if snapshot_reports.is_empty() { + (TableCompactionPlanningStatus::NoCandidates, 0) + } else { + ( + TableCompactionPlanningStatus::ManualReviewRequired, + snapshot_reports + .iter() + .filter(|snapshot| snapshot.status == TableCompactionPlanningStatus::ManualReviewRequired) + .count(), + ) + }; + + TableCompactionPlanningReport { + table_bucket: table_bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + table_id: entry.table_id.clone(), + current_metadata_location: entry.metadata_location.clone(), + current_snapshot_id, + config, + status, + candidate_file_count: 0, + rewrite_group_count: 0, + manual_review_count, + snapshot_reports, + } +} + +fn snapshot_expiration_drafts( + current_metadata: &serde_json::Value, + current_snapshot_id: Option, +) -> Vec { + let Some(snapshots) = current_metadata.get("snapshots").and_then(serde_json::Value::as_array) else { + return Vec::new(); + }; + + snapshots + .iter() + .map(|snapshot| { + let snapshot_id = snapshot.get("snapshot-id").and_then(serde_json::Value::as_i64); + let timestamp_ms = snapshot.get("timestamp-ms").and_then(serde_json::Value::as_i64); + let mut reasons = BTreeSet::new(); + if snapshot_id.is_none() { + reasons.insert(TableSnapshotExpirationReason::MissingSnapshotId); + } + if timestamp_ms.is_none() { + reasons.insert(TableSnapshotExpirationReason::MissingSnapshotTimestamp); + } + if snapshot_id.is_some() && snapshot_id == current_snapshot_id { + reasons.insert(TableSnapshotExpirationReason::CurrentSnapshot); + } + + TableSnapshotExpirationDraft { + snapshot_id, + sequence_number: snapshot.get("sequence-number").and_then(serde_json::Value::as_i64), + timestamp_ms, + manifest_list: snapshot + .get("manifest-list") + .and_then(serde_json::Value::as_str) + .map(ToString::to_string), + reasons, + } + }) + .collect() +} + +fn mark_recent_snapshots_to_keep(drafts: &mut [TableSnapshotExpirationDraft], min_snapshots_to_keep: usize) { + let mut snapshots_by_time = drafts + .iter() + .enumerate() + .filter_map(|(index, draft)| Some((draft.timestamp_ms?, index))) + .collect::>(); + snapshots_by_time.sort_by(|(left_timestamp, left_index), (right_timestamp, right_index)| { + right_timestamp.cmp(left_timestamp).then_with(|| left_index.cmp(right_index)) + }); + + for (_, index) in snapshots_by_time.into_iter().take(min_snapshots_to_keep) { + drafts[index] + .reasons + .insert(TableSnapshotExpirationReason::MinSnapshotsToKeep); + } +} + +fn snapshot_expiration_ref_state( + current_metadata: &serde_json::Value, + current_snapshot_id: Option, +) -> (BTreeSet, BTreeSet, BTreeSet) { + let mut protected_ref_snapshot_ids = BTreeSet::new(); + let mut user_defined_ref_snapshot_ids = BTreeSet::new(); + let mut ref_retention_conflict_snapshot_ids = BTreeSet::new(); + let Some(refs) = current_metadata.get("refs").and_then(serde_json::Value::as_object) else { + return ( + protected_ref_snapshot_ids, + user_defined_ref_snapshot_ids, + ref_retention_conflict_snapshot_ids, + ); + }; + + for (name, reference) in refs { + let Some(snapshot_id) = reference.get("snapshot-id").and_then(serde_json::Value::as_i64) else { + continue; + }; + if name != ICEBERG_MAIN_REF || Some(snapshot_id) != current_snapshot_id { + protected_ref_snapshot_ids.insert(snapshot_id); + } + if name != ICEBERG_MAIN_REF { + user_defined_ref_snapshot_ids.insert(snapshot_id); + } + if snapshot_ref_has_retention_policy(reference) { + ref_retention_conflict_snapshot_ids.insert(snapshot_id); + } + } + + ( + protected_ref_snapshot_ids, + user_defined_ref_snapshot_ids, + ref_retention_conflict_snapshot_ids, + ) +} + +fn snapshot_ref_has_retention_policy(reference: &serde_json::Value) -> bool { + reference.get(ICEBERG_REF_MIN_SNAPSHOTS_TO_KEEP_FIELD).is_some() + || reference.get(ICEBERG_REF_MAX_SNAPSHOT_AGE_MS_FIELD).is_some() + || reference.get(ICEBERG_REF_MAX_REF_AGE_MS_FIELD).is_some() +} + +fn snapshot_expiration_table_property_conflicts( + current_metadata: &serde_json::Value, + config: &TableSnapshotExpirationConfig, +) -> bool { + let Some(properties) = current_metadata.get("properties").and_then(serde_json::Value::as_object) else { + return false; + }; + + if properties.contains_key(ICEBERG_MAX_REF_AGE_MS_PROPERTY) { + return true; + } + if retention_property_conflicts_usize(properties, ICEBERG_MIN_SNAPSHOTS_TO_KEEP_PROPERTY, config.min_snapshots_to_keep) { + return true; + } + retention_property_conflicts_i64(properties, ICEBERG_MAX_SNAPSHOT_AGE_MS_PROPERTY, config.max_snapshot_age_ms) +} + +fn retention_property_conflicts_usize( + properties: &serde_json::Map, + key: &str, + expected: usize, +) -> bool { + let Some(value) = properties.get(key) else { + return false; + }; + serde_json_i64(value).and_then(|value| usize::try_from(value).ok()) != Some(expected) +} + +fn retention_property_conflicts_i64(properties: &serde_json::Map, key: &str, expected: i64) -> bool { + let Some(value) = properties.get(key) else { + return false; + }; + serde_json_i64(value) != Some(expected) +} + +fn serde_json_i64(value: &serde_json::Value) -> Option { + value.as_i64().or_else(|| value.as_str()?.parse::().ok()) +} + +fn snapshot_expiration_requires_manual_review(reasons: &BTreeSet) -> bool { + reasons.contains(&TableSnapshotExpirationReason::MissingSnapshotId) + || reasons.contains(&TableSnapshotExpirationReason::MissingSnapshotTimestamp) + || reasons.contains(&TableSnapshotExpirationReason::UserDefinedSnapshotRef) + || reasons.contains(&TableSnapshotExpirationReason::SnapshotRefRetentionConflict) + || reasons.contains(&TableSnapshotExpirationReason::TableRetentionPropertyConflict) +} + +fn snapshot_expiration_is_retained(reasons: &BTreeSet) -> bool { + reasons.contains(&TableSnapshotExpirationReason::CurrentSnapshot) + || reasons.contains(&TableSnapshotExpirationReason::MinSnapshotsToKeep) + || reasons.contains(&TableSnapshotExpirationReason::ProtectedSnapshotRef) +} + +fn unix_timestamp_millis(now: OffsetDateTime) -> i64 { + now.unix_timestamp() + .saturating_mul(1000) + .saturating_add(i64::from(now.millisecond())) +} + fn maintenance_timestamp(now: OffsetDateTime) -> String { now.format(&time::format_description::well_known::Rfc3339) .unwrap_or_else(|_| now.unix_timestamp().to_string()) @@ -2757,6 +3341,45 @@ fn validate_table_maintenance_config(config: &TableMaintenanceConfig) -> TableCa Ok(()) } +fn validate_table_snapshot_expiration_config(config: &TableSnapshotExpirationConfig) -> TableCatalogStoreResult<()> { + if config.min_snapshots_to_keep == 0 { + return Err(TableCatalogStoreError::Invalid( + "min-snapshots-to-keep must be greater than zero".to_string(), + )); + } + if config.max_snapshot_age_ms < 0 { + return Err(TableCatalogStoreError::Invalid("max-snapshot-age-ms cannot be negative".to_string())); + } + Ok(()) +} + +fn validate_table_compaction_planning_config(config: &TableCompactionPlanningConfig) -> TableCatalogStoreResult<()> { + if config.target_file_size_bytes == 0 { + return Err(TableCatalogStoreError::Invalid( + "target-file-size-bytes must be greater than zero".to_string(), + )); + } + if config.small_file_threshold_bytes == 0 { + return Err(TableCatalogStoreError::Invalid( + "small-file-threshold-bytes must be greater than zero".to_string(), + )); + } + if config.small_file_threshold_bytes > config.target_file_size_bytes { + return Err(TableCatalogStoreError::Invalid( + "small-file-threshold-bytes cannot exceed target-file-size-bytes".to_string(), + )); + } + if config.min_input_files < 2 { + return Err(TableCatalogStoreError::Invalid("min-input-files must be at least two".to_string())); + } + if config.max_rewrite_bytes_per_job < config.target_file_size_bytes { + return Err(TableCatalogStoreError::Invalid( + "max-rewrite-bytes-per-job must be at least target-file-size-bytes".to_string(), + )); + } + Ok(()) +} + fn commit_log_matches_request(commit_log: &CommitLogEntry, request: &TableCommitRequest, table_id: &str) -> bool { commit_log.version == TABLE_CATALOG_ENTRY_VERSION && commit_log.commit_id == request.commit_id @@ -3842,6 +4465,25 @@ mod tests { .expect("metadata maintenance object report should exist") } + fn snapshot_expiration_report( + report: &TableSnapshotExpirationReport, + snapshot_id: i64, + ) -> &TableSnapshotExpirationSnapshotReport { + report + .snapshot_reports + .iter() + .find(|snapshot| snapshot.snapshot_id == Some(snapshot_id)) + .expect("snapshot expiration report should include the snapshot") + } + + fn compaction_snapshot_report(report: &TableCompactionPlanningReport, snapshot_id: i64) -> &TableCompactionSnapshotReport { + report + .snapshot_reports + .iter() + .find(|snapshot| snapshot.snapshot_id == Some(snapshot_id)) + .expect("compaction planning report should include the snapshot") + } + #[async_trait::async_trait] impl TableCatalogObjectBackend for TestCatalogObjectBackend { async fn read_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult> { @@ -4752,6 +5394,349 @@ mod tests { assert_eq!(report.cleanup_candidate_locations, vec![orphan, unreferenced]); } + #[tokio::test] + async fn snapshot_expiration_plan_retains_current_recent_and_protected_refs() { + 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 current = default_table_metadata_file_path(&namespace, &table, "00004.metadata.json"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend + .seed_object( + bucket, + ¤t, + serde_json::to_vec(&serde_json::json!({ + "current-snapshot-id": 30, + "metadata-log": [], + "snapshots": [ + { + "snapshot-id": 10, + "timestamp-ms": 1000, + "manifest-list": "s3://analytics/tables/table-id/metadata/snap-10.avro" + }, + { + "snapshot-id": 20, + "timestamp-ms": 2000, + "manifest-list": "s3://analytics/tables/table-id/metadata/snap-20.avro" + }, + { + "snapshot-id": 30, + "timestamp-ms": 3000, + "manifest-list": "s3://analytics/tables/table-id/metadata/snap-30.avro" + } + ], + "refs": { + "main": { + "snapshot-id": 30, + "type": "branch" + }, + "audit": { + "snapshot-id": 10, + "type": "tag" + } + } + })) + .unwrap(), + ) + .await; + + let report = store + .plan_table_snapshot_expiration( + bucket, + "sales", + "orders", + TableSnapshotExpirationConfig { + min_snapshots_to_keep: 1, + max_snapshot_age_ms: 1, + }, + ) + .await + .expect("snapshot expiration planning should succeed"); + + let current_snapshot = snapshot_expiration_report(&report, 30); + assert_eq!(current_snapshot.state, TableSnapshotExpirationSnapshotState::Retained); + assert!( + current_snapshot + .reasons + .contains(&TableSnapshotExpirationReason::CurrentSnapshot) + ); + + let protected_snapshot = snapshot_expiration_report(&report, 10); + assert_eq!(protected_snapshot.state, TableSnapshotExpirationSnapshotState::ManualReviewRequired); + assert!( + protected_snapshot + .reasons + .contains(&TableSnapshotExpirationReason::ProtectedSnapshotRef) + ); + assert!( + protected_snapshot + .reasons + .contains(&TableSnapshotExpirationReason::UserDefinedSnapshotRef) + ); + + let expired_snapshot = snapshot_expiration_report(&report, 20); + assert_eq!(expired_snapshot.state, TableSnapshotExpirationSnapshotState::ExpirationCandidate); + assert!( + expired_snapshot + .reasons + .contains(&TableSnapshotExpirationReason::SnapshotAgeExpired) + ); + assert_eq!(report.expiration_candidate_count, 1); + assert_eq!(report.manual_review_count, 1); + } + + #[tokio::test] + async fn snapshot_expiration_plan_fails_closed_for_table_retention_property_conflicts() { + 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 current = default_table_metadata_file_path(&namespace, &table, "00004.metadata.json"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend + .seed_object( + bucket, + ¤t, + serde_json::to_vec(&serde_json::json!({ + "current-snapshot-id": 20, + "metadata-log": [], + "properties": { + "history.expire.min-snapshots-to-keep": "5" + }, + "snapshots": [ + { + "snapshot-id": 10, + "timestamp-ms": 1000, + "manifest-list": "s3://analytics/tables/table-id/metadata/snap-10.avro" + }, + { + "snapshot-id": 20, + "timestamp-ms": 2000, + "manifest-list": "s3://analytics/tables/table-id/metadata/snap-20.avro" + } + ], + "refs": { + "main": { + "snapshot-id": 20, + "type": "branch" + } + } + })) + .unwrap(), + ) + .await; + + let report = store + .plan_table_snapshot_expiration( + bucket, + "sales", + "orders", + TableSnapshotExpirationConfig { + min_snapshots_to_keep: 1, + max_snapshot_age_ms: 1, + }, + ) + .await + .expect("snapshot expiration planning should succeed"); + + assert_eq!(report.expiration_candidate_count, 0); + assert_eq!(report.manual_review_count, 2); + for snapshot in &report.snapshot_reports { + assert_eq!(snapshot.state, TableSnapshotExpirationSnapshotState::ManualReviewRequired); + assert!( + snapshot + .reasons + .contains(&TableSnapshotExpirationReason::TableRetentionPropertyConflict) + ); + } + } + + #[tokio::test] + async fn snapshot_expiration_plan_requires_snapshot_timestamps() { + 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 current = default_table_metadata_file_path(&namespace, &table, "00004.metadata.json"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend + .seed_object( + bucket, + ¤t, + serde_json::to_vec(&serde_json::json!({ + "current-snapshot-id": 20, + "metadata-log": [], + "snapshots": [ + { + "snapshot-id": 10, + "manifest-list": "s3://analytics/tables/table-id/metadata/snap-10.avro" + }, + { + "snapshot-id": 20, + "timestamp-ms": 2000, + "manifest-list": "s3://analytics/tables/table-id/metadata/snap-20.avro" + } + ], + "refs": { + "main": { + "snapshot-id": 20, + "type": "branch" + } + } + })) + .unwrap(), + ) + .await; + + let report = store + .plan_table_snapshot_expiration( + bucket, + "sales", + "orders", + TableSnapshotExpirationConfig { + min_snapshots_to_keep: 1, + max_snapshot_age_ms: 1, + }, + ) + .await + .expect("snapshot expiration planning should succeed"); + + let missing_timestamp = snapshot_expiration_report(&report, 10); + assert_eq!(missing_timestamp.state, TableSnapshotExpirationSnapshotState::ManualReviewRequired); + assert!( + missing_timestamp + .reasons + .contains(&TableSnapshotExpirationReason::MissingSnapshotTimestamp) + ); + assert_eq!(report.expiration_candidate_count, 0); + assert_eq!(report.manual_review_count, 1); + } + + #[tokio::test] + async fn compaction_plan_reports_manifest_reader_gap_without_rewrite_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 current = default_table_metadata_file_path(&namespace, &table, "00004.metadata.json"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend + .seed_object( + bucket, + ¤t, + serde_json::to_vec(&serde_json::json!({ + "current-snapshot-id": 20, + "metadata-log": [], + "snapshots": [ + { + "snapshot-id": 20, + "timestamp-ms": 2000, + "manifest-list": "s3://analytics/tables/table-id/metadata/snap-20.avro" + } + ], + "refs": { + "main": { + "snapshot-id": 20, + "type": "branch" + } + } + })) + .unwrap(), + ) + .await; + + let report = store + .plan_table_compaction( + bucket, + "sales", + "orders", + TableCompactionPlanningConfig { + target_file_size_bytes: 512 * 1024 * 1024, + small_file_threshold_bytes: 64 * 1024 * 1024, + min_input_files: 2, + max_rewrite_bytes_per_job: 1024 * 1024 * 1024, + }, + ) + .await + .expect("compaction planning should succeed"); + + assert_eq!(report.status, TableCompactionPlanningStatus::ManualReviewRequired); + assert_eq!(report.candidate_file_count, 0); + assert_eq!(report.rewrite_group_count, 0); + assert_eq!(report.manual_review_count, 1); + let snapshot = compaction_snapshot_report(&report, 20); + assert_eq!(snapshot.status, TableCompactionPlanningStatus::ManualReviewRequired); + assert!(snapshot.reasons.contains(&TableCompactionPlanningReason::ManifestList)); + assert!( + snapshot + .reasons + .contains(&TableCompactionPlanningReason::ManifestAvroReaderUnavailable) + ); + } + + #[tokio::test] + async fn compaction_plan_requires_current_snapshot_metadata() { + 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 current = default_table_metadata_file_path(&namespace, &table, "00004.metadata.json"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend + .seed_object( + bucket, + ¤t, + serde_json::to_vec(&serde_json::json!({ + "current-snapshot-id": 30, + "metadata-log": [], + "snapshots": [ + { + "snapshot-id": 20, + "timestamp-ms": 2000, + "manifest-list": "s3://analytics/tables/table-id/metadata/snap-20.avro" + } + ] + })) + .unwrap(), + ) + .await; + + let report = store + .plan_table_compaction( + bucket, + "sales", + "orders", + TableCompactionPlanningConfig { + target_file_size_bytes: 512 * 1024 * 1024, + small_file_threshold_bytes: 64 * 1024 * 1024, + min_input_files: 2, + max_rewrite_bytes_per_job: 1024 * 1024 * 1024, + }, + ) + .await + .expect("compaction planning should succeed"); + + assert_eq!(report.status, TableCompactionPlanningStatus::ManualReviewRequired); + assert_eq!(report.manual_review_count, 1); + let snapshot = compaction_snapshot_report(&report, 30); + assert!( + snapshot + .reasons + .contains(&TableCompactionPlanningReason::MissingCurrentSnapshot) + ); + } + #[tokio::test] async fn maintenance_dry_run_keeps_recent_metadata_files_and_ignores_non_metadata_objects() { let backend = TestCatalogObjectBackend::default(); diff --git a/scripts/table-catalog/README.md b/scripts/table-catalog/README.md index fd7718eca..0c7ce735e 100644 --- a/scripts/table-catalog/README.md +++ b/scripts/table-catalog/README.md @@ -151,7 +151,9 @@ 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: unsupported - manifest/data reachability cleanup: unsupported -- snapshot expiration and compaction: unsupported +- snapshot expiration dry-run planning and manual catalog commit: supported through metadata maintenance reports +- automatic maintenance scheduling: unsupported +- compaction rewrite: unsupported; planning reports fail closed until manifest Avro reading and rewrite support are implemented - Iceberg views: unsupported - multi-table transactions: not a short-term production claim diff --git a/scripts/table-catalog/pyiceberg_smoke.py b/scripts/table-catalog/pyiceberg_smoke.py index 9e0eca5d5..f058c2483 100755 --- a/scripts/table-catalog/pyiceberg_smoke.py +++ b/scripts/table-catalog/pyiceberg_smoke.py @@ -144,10 +144,16 @@ UNSUPPORTED_INVENTORY: list[dict[str, str]] = [ "expected_behavior": "metadata-only cleanup must not delete manifest, data, or delete files", }, { - "capability": "snapshot-expiration-and-compaction", - "status": "unsupported", + "capability": "snapshot-expiration", + "status": "manual-maintenance-supported", "roadmap_area": "snapshot-maintenance", - "expected_behavior": "no automatic snapshot expiration or data rewrite is claimed", + "expected_behavior": "metadata maintenance can report snapshot expiration plans and manually commit safe snapshot expiration through the catalog; background scheduling and manifest/data cleanup are not claimed", + }, + { + "capability": "compaction-rewrite", + "status": "planning-only", + "roadmap_area": "snapshot-maintenance", + "expected_behavior": "compaction planning fails closed when manifest Avro reading is required; no data file rewrite or automatic compaction is claimed", }, { "capability": "iceberg-views",