From 65c724c786c06bdb56818cf318b964471c374627 Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Mon, 15 Jun 2026 22:29:23 +0800 Subject: [PATCH] feat(table-catalog): add manifest reachability cleanup (#3484) Co-authored-by: Henry Guo --- Cargo.lock | 55 +- Cargo.toml | 1 + rustfs/Cargo.toml | 1 + rustfs/src/table_catalog.rs | 1319 +++++++++++++++++++++- scripts/table-catalog/README.md | 4 +- scripts/table-catalog/pyiceberg_smoke.py | 6 +- 6 files changed, 1355 insertions(+), 31 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index dcb1967e2..9209ef19b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -275,6 +275,30 @@ version = "1.0.102" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7f202df86484c868dbad7eaa557ef785d5c66295e41b460ef922eca0723b842c" +[[package]] +name = "apache-avro" +version = "0.21.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "36fa98bc79671c7981272d91a8753a928ff6a1cd8e4f20a44c45bd5d313840bf" +dependencies = [ + "bigdecimal", + "bon", + "digest 0.10.7", + "log", + "miniz_oxide", + "num-bigint", + "quad-rand", + "rand 0.9.4", + "regex-lite", + "serde", + "serde_bytes", + "serde_json", + "strum 0.27.2", + "strum_macros 0.27.2", + "thiserror 2.0.18", + "uuid", +] + [[package]] name = "ar_archive_writer" version = "0.5.2" @@ -1441,6 +1465,7 @@ dependencies = [ "num-bigint", "num-integer", "num-traits", + "serde", ] [[package]] @@ -6640,6 +6665,7 @@ checksum = "a5e44f723f1133c9deac646763579fdb3ac745e418f2a7af9cd0c431da1f20b9" dependencies = [ "num-integer", "num-traits", + "serde", ] [[package]] @@ -8172,6 +8198,12 @@ dependencies = [ "uuid", ] +[[package]] +name = "quad-rand" +version = "0.2.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a651516ddc9168ebd67b24afd085a718be02f8858fe406591b013d101ce2f40" + [[package]] name = "quanta" version = "0.12.6" @@ -9046,6 +9078,7 @@ version = "1.0.0-beta.8" dependencies = [ "aes-gcm", "anyhow", + "apache-avro", "astral-tokio-tar", "async-trait", "async_zip", @@ -9766,7 +9799,7 @@ dependencies = [ "rustfs-crypto", "serde", "serde_json", - "strum", + "strum 0.28.0", "temp-env", "test-case", "thiserror 2.0.18", @@ -11264,13 +11297,31 @@ version = "0.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f" +[[package]] +name = "strum" +version = "0.27.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "af23d6f6c1a224baef9d3f61e287d2761385a5b88fdab4eb4c6f11aeb54c4bcf" + [[package]] name = "strum" version = "0.28.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9628de9b8791db39ceda2b119bbe13134770b56c138ec1d3af810d045c04f9bd" dependencies = [ - "strum_macros", + "strum_macros 0.28.0", +] + +[[package]] +name = "strum_macros" +version = "0.27.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7695ce3845ea4b33927c055a39dc438a45b059f7c1b3d91d38d10355fb8cbca7" +dependencies = [ + "heck", + "proc-macro2", + "quote", + "syn 2.0.117", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index a56d81fdd..fdba38b87 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -160,6 +160,7 @@ tower = { version = "0.5.3", features = ["timeout"] } tower-http = { version = "0.6.11", features = ["cors"] } # Serialization and Data Formats +apache-avro = "0.21.0" bytes = { version = "1.11.1", features = ["serde"] } bytesize = "2.4.0" byteorder = "1.5.0" diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 5d44790a6..200280776 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -123,6 +123,7 @@ tower.workspace = true tower-http = { workspace = true, features = ["trace", "compression-full", "cors", "catch-panic", "timeout", "limit", "request-id", "add-extension"] } # Serialization and Data Formats +apache-avro = { workspace = true } bytes = { workspace = true } flatbuffers.workspace = true rmp-serde.workspace = true diff --git a/rustfs/src/table_catalog.rs b/rustfs/src/table_catalog.rs index 4a7567127..3f66151e2 100644 --- a/rustfs/src/table_catalog.rs +++ b/rustfs/src/table_catalog.rs @@ -67,6 +67,8 @@ const TABLE_MARKER_FILE: &str = "table.json"; const CURRENT_POINTER_FILE: &str = "current.json"; const LIFECYCLE_FILE: &str = "lifecycle.json"; const METADATA_DIR: &str = "metadata"; +const DATA_DIR: &str = "data"; +const DELETE_DIR: &str = "delete"; const TABLE_BUCKET_ENTRY_FILE: &str = "table-bucket.json"; const NAMESPACE_ENTRY_FILE: &str = "namespace-entry.json"; const TABLE_ENTRY_FILE: &str = "table-entry.json"; @@ -380,6 +382,14 @@ pub(crate) struct TableMetadataMaintenanceJob { #[serde(default)] pub deleted_metadata_file_count: usize, #[serde(default)] + pub planned_object_file_count: usize, + #[serde(default)] + pub cleanup_candidate_object_count: usize, + #[serde(default)] + pub deletable_object_count: usize, + #[serde(default)] + pub deleted_object_count: usize, + #[serde(default)] pub quarantined_object_count: usize, } @@ -390,8 +400,14 @@ pub(crate) struct TableMetadataMaintenanceReport { pub retained_metadata_locations: Vec, pub cleanup_candidate_locations: Vec, pub deletable_metadata_locations: Vec, + #[serde(default, rename = "cleanup-object-candidate-locations")] + pub cleanup_object_candidate_locations: Vec, + #[serde(default, rename = "deletable-object-locations")] + pub deletable_object_locations: Vec, #[serde(default)] pub object_reports: Vec, + #[serde(default, rename = "object-cleanup-reports")] + pub object_cleanup_reports: Vec, #[serde(default)] pub referenced_object_reports: Vec, #[serde(default, rename = "reachability-graph")] @@ -589,7 +605,11 @@ pub(crate) enum TableMetadataMaintenanceReason { SafetyWindowSatisfied, DeletedByMaintenance, ManifestList, + ManifestFile, + DataFile, + DeleteFile, UnsupportedManifestAvro, + UnreadableMetadata, QuarantineEnabled, RetryScheduled, } @@ -611,6 +631,14 @@ pub(crate) struct TableMetadataMaintenanceObjectReport { pub reasons: Vec, } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub(crate) struct TableMetadataMaintenanceObjectCleanupReport { + pub object_location: String, + pub object_kind: TableMetadataMaintenanceObjectKind, + pub state: TableMetadataMaintenanceObjectState, + pub reasons: Vec, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub(crate) struct TableMetadataMaintenanceReferencedObjectReport { pub object_location: String, @@ -1826,13 +1854,20 @@ where cleanup_candidate_count: 0, deletable_metadata_file_count: 0, deleted_metadata_file_count: 0, + planned_object_file_count: 0, + cleanup_candidate_object_count: 0, + deletable_object_count: 0, + deleted_object_count: 0, quarantined_object_count: 0, }, current_metadata_location, retained_metadata_locations: Vec::new(), cleanup_candidate_locations: Vec::new(), deletable_metadata_locations: Vec::new(), + cleanup_object_candidate_locations: Vec::new(), + deletable_object_locations: Vec::new(), object_reports: Vec::new(), + object_cleanup_reports: Vec::new(), referenced_object_reports: Vec::new(), reachability_graph: TableMaintenanceReachabilityGraphReport::default(), snapshot_expiration: None, @@ -2202,12 +2237,33 @@ where ); } } + let warehouse_object_prefix = table_warehouse_object_prefix(&entry).ok(); let current_metadata_location = entry.metadata_location; let retained_metadata_locations = retained.into_iter().collect::>(); let object_reports = metadata_maintenance_object_reports(maintenance_reasons); - let referenced_object_reports = metadata_maintenance_referenced_object_reports(¤t_metadata); + let referenced_object_reports = metadata_maintenance_referenced_object_reports( + &self.backend, + table_bucket, + &namespace, + &table, + warehouse_object_prefix.as_deref(), + ¤t_metadata, + &retained_metadata_locations, + ) + .await?; let reachability_graph = metadata_maintenance_reachability_graph_report(planned_metadata_file_count, &referenced_object_reports); + let (planned_object_file_count, cleanup_object_candidate_locations, deletable_object_locations, object_cleanup_reports) = + metadata_maintenance_object_cleanup_reports( + &self.backend, + table_bucket, + &namespace, + &table, + warehouse_object_prefix.as_deref(), + &referenced_object_reports, + now, + ) + .await?; Ok(TableMetadataMaintenanceReport { job: TableMetadataMaintenanceJob { @@ -2241,13 +2297,20 @@ where cleanup_candidate_count: cleanup_candidate_locations.len(), deletable_metadata_file_count: deletable_metadata_locations.len(), deleted_metadata_file_count: 0, + planned_object_file_count, + cleanup_candidate_object_count: cleanup_object_candidate_locations.len(), + deletable_object_count: deletable_object_locations.len(), + deleted_object_count: 0, quarantined_object_count: 0, }, current_metadata_location, retained_metadata_locations, cleanup_candidate_locations, deletable_metadata_locations, + cleanup_object_candidate_locations, + deletable_object_locations, object_reports, + object_cleanup_reports, referenced_object_reports, reachability_graph, snapshot_expiration: None, @@ -2403,6 +2466,7 @@ where "current metadata location changed before maintenance delete".to_string(), )); } + let warehouse_object_prefix = table_warehouse_object_prefix(&entry).ok(); let Some(current_metadata_object) = self.backend.read_object(table_bucket, &entry.metadata_location).await? else { return Err(TableCatalogStoreError::NotFound(format!( @@ -2467,6 +2531,66 @@ where for metadata_location in &cleanup_candidate_locations { self.backend.delete_object(table_bucket, metadata_location).await?; } + + let referenced_object_reports = metadata_maintenance_referenced_object_reports( + &self.backend, + table_bucket, + &namespace, + &table, + warehouse_object_prefix.as_deref(), + ¤t_metadata, + &report.retained_metadata_locations, + ) + .await?; + let referenced_object_locations = if referenced_object_reports + .iter() + .any(|report| report.state == TableMetadataMaintenanceObjectState::ManualReviewRequired) + { + BTreeSet::new() + } else { + referenced_object_reports + .iter() + .filter_map(|report| table_catalog_object_key_from_location(table_bucket, &report.object_location)) + .collect::>() + }; + let planned_deletable_object_locations = report.deletable_object_locations.iter().cloned().collect::>(); + let mut cleanup_object_candidate_locations = BTreeSet::new(); + if !referenced_object_reports + .iter() + .any(|report| report.state == TableMetadataMaintenanceObjectState::ManualReviewRequired) + { + for object_location in &report.cleanup_object_candidate_locations { + if table_maintenance_object_kind(&namespace, &table, warehouse_object_prefix.as_deref(), object_location) + .is_none() + { + return Err(TableCatalogStoreError::Invalid(format!( + "cleanup object candidate {object_location} must be inside table metadata, data, or delete directories" + ))); + } + if referenced_object_locations.contains(object_location.as_str()) { + return Err(TableCatalogStoreError::Conflict(format!( + "cleanup object candidate {object_location} is retained by current metadata" + ))); + } + let Some(candidate_object) = self.backend.read_object(table_bucket, object_location).await? else { + continue; + }; + if !planned_deletable_object_locations.contains(object_location.as_str()) { + continue; + } + if !metadata_candidate_is_past_safety_window(candidate_object.mod_time, now) { + continue; + } + cleanup_object_candidate_locations.insert(object_location.clone()); + } + } + + let cleanup_object_candidate_locations = cleanup_object_candidate_locations.into_iter().collect::>(); + let deleted_object_locations = cleanup_object_candidate_locations.iter().cloned().collect::>(); + for object_location in &cleanup_object_candidate_locations { + self.backend.delete_object(table_bucket, object_location).await?; + } + let retained_metadata_locations = protected.into_iter().collect::>(); let mut job = report.job; job.operation = TableMetadataMaintenanceOperation::Delete; @@ -2476,9 +2600,13 @@ where job.cleanup_candidate_count = cleanup_candidate_count; job.deletable_metadata_file_count = planned_deletable_locations.len(); job.deleted_metadata_file_count = cleanup_candidate_locations.len(); + job.cleanup_candidate_object_count = report.cleanup_object_candidate_locations.len(); + job.deletable_object_count = planned_deletable_object_locations.len(); + job.deleted_object_count = cleanup_object_candidate_locations.len(); let mut object_reports = report.object_reports; mark_deleted_metadata_object_reports(&mut object_reports, &deleted_locations); - let referenced_object_reports = report.referenced_object_reports; + let mut object_cleanup_reports = report.object_cleanup_reports; + mark_deleted_object_cleanup_reports(&mut object_cleanup_reports, &deleted_object_locations); Ok(TableMetadataMaintenanceReport { job, @@ -2486,7 +2614,10 @@ where retained_metadata_locations, cleanup_candidate_locations: cleanup_candidate_locations.clone(), deletable_metadata_locations: cleanup_candidate_locations, + cleanup_object_candidate_locations: cleanup_object_candidate_locations.clone(), + deletable_object_locations: cleanup_object_candidate_locations, object_reports, + object_cleanup_reports, referenced_object_reports, reachability_graph: report.reachability_graph, snapshot_expiration: report.snapshot_expiration, @@ -3140,31 +3271,524 @@ fn metadata_maintenance_object_reports( .collect() } -fn metadata_maintenance_referenced_object_reports( +#[derive(Debug, Clone)] +struct TableMetadataMaintenanceReferencedObjectAccumulator { + object_kind: TableMetadataMaintenanceObjectKind, + state: TableMetadataMaintenanceObjectState, + reasons: BTreeSet, +} + +fn insert_referenced_object_report( + reports: &mut BTreeMap, + object_location: String, + object_kind: TableMetadataMaintenanceObjectKind, + state: TableMetadataMaintenanceObjectState, + reason: TableMetadataMaintenanceReason, +) { + let report = reports + .entry(object_location) + .or_insert_with(|| TableMetadataMaintenanceReferencedObjectAccumulator { + object_kind, + state: TableMetadataMaintenanceObjectState::Retained, + reasons: BTreeSet::new(), + }); + if state == TableMetadataMaintenanceObjectState::ManualReviewRequired { + report.state = TableMetadataMaintenanceObjectState::ManualReviewRequired; + } + report.reasons.insert(reason); +} + +async fn metadata_maintenance_referenced_object_reports( + backend: &B, + table_bucket: &str, + namespace: &Namespace, + table: &IdentifierSegment, + warehouse_object_prefix: Option<&str>, current_metadata: &serde_json::Value, -) -> Vec { - let mut reasons_by_location = BTreeMap::>::new(); - if let Some(snapshots) = current_metadata.get("snapshots").and_then(serde_json::Value::as_array) { - for snapshot in snapshots { - let Some(manifest_list) = snapshot.get("manifest-list").and_then(serde_json::Value::as_str) else { + retained_metadata_locations: &[String], +) -> TableCatalogStoreResult> +where + B: TableCatalogObjectBackend, +{ + let mut reports = BTreeMap::::new(); + metadata_maintenance_referenced_object_reports_for_metadata( + backend, + table_bucket, + namespace, + table, + warehouse_object_prefix, + current_metadata, + &mut reports, + ) + .await?; + + for metadata_location in retained_metadata_locations { + let Some(metadata_object) = backend.read_object(table_bucket, metadata_location).await? else { + insert_referenced_object_report( + &mut reports, + metadata_location.clone(), + TableMetadataMaintenanceObjectKind::MetadataFile, + TableMetadataMaintenanceObjectState::ManualReviewRequired, + TableMetadataMaintenanceReason::UnreadableMetadata, + ); + continue; + }; + let Ok(metadata) = serde_json::from_slice::(&metadata_object.data) else { + insert_referenced_object_report( + &mut reports, + metadata_location.clone(), + TableMetadataMaintenanceObjectKind::MetadataFile, + TableMetadataMaintenanceObjectState::ManualReviewRequired, + TableMetadataMaintenanceReason::UnreadableMetadata, + ); + continue; + }; + if !metadata.is_object() { + insert_referenced_object_report( + &mut reports, + metadata_location.clone(), + TableMetadataMaintenanceObjectKind::MetadataFile, + TableMetadataMaintenanceObjectState::ManualReviewRequired, + TableMetadataMaintenanceReason::UnreadableMetadata, + ); + continue; + } + metadata_maintenance_referenced_object_reports_for_metadata( + backend, + table_bucket, + namespace, + table, + warehouse_object_prefix, + &metadata, + &mut reports, + ) + .await?; + } + + Ok(reports + .into_iter() + .map(|(object_location, report)| TableMetadataMaintenanceReferencedObjectReport { + object_location, + object_kind: report.object_kind, + state: report.state, + reasons: report.reasons.into_iter().collect(), + }) + .collect()) +} + +async fn metadata_maintenance_referenced_object_reports_for_metadata( + backend: &B, + table_bucket: &str, + namespace: &Namespace, + table: &IdentifierSegment, + warehouse_object_prefix: Option<&str>, + metadata: &serde_json::Value, + reports: &mut BTreeMap, +) -> TableCatalogStoreResult<()> +where + B: TableCatalogObjectBackend, +{ + let Some(snapshots) = metadata.get("snapshots").and_then(serde_json::Value::as_array) else { + return Ok(()); + }; + + for snapshot in snapshots { + if let Some(manifest_list_location) = snapshot.get("manifest-list").and_then(serde_json::Value::as_str) { + metadata_maintenance_referenced_manifest_list( + backend, + table_bucket, + namespace, + table, + warehouse_object_prefix, + manifest_list_location, + reports, + ) + .await?; + continue; + } + + let Some(manifests) = snapshot.get("manifests").and_then(serde_json::Value::as_array) else { + continue; + }; + for manifest in manifests { + let Some(manifest_location) = manifest.as_str() else { + insert_referenced_object_report( + reports, + "snapshots[].manifests".to_string(), + TableMetadataMaintenanceObjectKind::ManifestFile, + TableMetadataMaintenanceObjectState::ManualReviewRequired, + TableMetadataMaintenanceReason::UnsupportedManifestAvro, + ); continue; }; - reasons_by_location.entry(manifest_list.to_string()).or_default().extend([ - TableMetadataMaintenanceReason::ManifestList, - TableMetadataMaintenanceReason::UnsupportedManifestAvro, - ]); + metadata_maintenance_referenced_manifest_file( + backend, + table_bucket, + namespace, + table, + warehouse_object_prefix, + manifest_location, + reports, + ) + .await?; } } - reasons_by_location - .into_iter() - .map(|(object_location, reasons)| TableMetadataMaintenanceReferencedObjectReport { - object_location, - object_kind: TableMetadataMaintenanceObjectKind::ManifestList, - state: TableMetadataMaintenanceObjectState::ManualReviewRequired, - reasons: reasons.into_iter().collect(), - }) - .collect() + Ok(()) +} + +async fn metadata_maintenance_referenced_manifest_list( + backend: &B, + table_bucket: &str, + namespace: &Namespace, + table: &IdentifierSegment, + warehouse_object_prefix: Option<&str>, + manifest_list_location: &str, + reports: &mut BTreeMap, +) -> TableCatalogStoreResult<()> +where + B: TableCatalogObjectBackend, +{ + let Some(manifest_list_key) = table_catalog_object_key_from_location(table_bucket, manifest_list_location) else { + insert_referenced_object_report( + reports, + manifest_list_location.to_string(), + TableMetadataMaintenanceObjectKind::ManifestList, + TableMetadataMaintenanceObjectState::ManualReviewRequired, + TableMetadataMaintenanceReason::UnsupportedManifestAvro, + ); + return Ok(()); + }; + if table_maintenance_object_kind(namespace, table, warehouse_object_prefix, &manifest_list_key) + != Some(TableMetadataMaintenanceObjectKind::ManifestList) + { + insert_referenced_object_report( + reports, + manifest_list_key, + TableMetadataMaintenanceObjectKind::ManifestList, + TableMetadataMaintenanceObjectState::ManualReviewRequired, + TableMetadataMaintenanceReason::UnsupportedManifestAvro, + ); + return Ok(()); + } + insert_referenced_object_report( + reports, + manifest_list_key.clone(), + TableMetadataMaintenanceObjectKind::ManifestList, + TableMetadataMaintenanceObjectState::Retained, + TableMetadataMaintenanceReason::ManifestList, + ); + + let Some(manifest_list_object) = backend.read_object(table_bucket, &manifest_list_key).await? else { + mark_referenced_object_manual_review( + reports, + &manifest_list_key, + TableMetadataMaintenanceReason::UnsupportedManifestAvro, + ); + return Ok(()); + }; + let Ok(manifest_paths) = manifest_paths_from_manifest_list_avro(&manifest_list_object.data) else { + mark_referenced_object_manual_review( + reports, + &manifest_list_key, + TableMetadataMaintenanceReason::UnsupportedManifestAvro, + ); + return Ok(()); + }; + for manifest_location in manifest_paths { + metadata_maintenance_referenced_manifest_file( + backend, + table_bucket, + namespace, + table, + warehouse_object_prefix, + &manifest_location, + reports, + ) + .await?; + } + + Ok(()) +} + +async fn metadata_maintenance_referenced_manifest_file( + backend: &B, + table_bucket: &str, + namespace: &Namespace, + table: &IdentifierSegment, + warehouse_object_prefix: Option<&str>, + manifest_location: &str, + reports: &mut BTreeMap, +) -> TableCatalogStoreResult<()> +where + B: TableCatalogObjectBackend, +{ + let Some(manifest_key) = table_catalog_object_key_from_location(table_bucket, manifest_location) else { + insert_referenced_object_report( + reports, + manifest_location.to_string(), + TableMetadataMaintenanceObjectKind::ManifestFile, + TableMetadataMaintenanceObjectState::ManualReviewRequired, + TableMetadataMaintenanceReason::UnsupportedManifestAvro, + ); + return Ok(()); + }; + if table_maintenance_object_kind(namespace, table, warehouse_object_prefix, &manifest_key) + != Some(TableMetadataMaintenanceObjectKind::ManifestFile) + { + insert_referenced_object_report( + reports, + manifest_key, + TableMetadataMaintenanceObjectKind::ManifestFile, + TableMetadataMaintenanceObjectState::ManualReviewRequired, + TableMetadataMaintenanceReason::UnsupportedManifestAvro, + ); + return Ok(()); + } + insert_referenced_object_report( + reports, + manifest_key.clone(), + TableMetadataMaintenanceObjectKind::ManifestFile, + TableMetadataMaintenanceObjectState::Retained, + TableMetadataMaintenanceReason::ManifestFile, + ); + + let Some(manifest_object) = backend.read_object(table_bucket, &manifest_key).await? else { + mark_referenced_object_manual_review(reports, &manifest_key, TableMetadataMaintenanceReason::UnsupportedManifestAvro); + return Ok(()); + }; + let Ok(file_references) = file_references_from_manifest_avro(&manifest_object.data) else { + mark_referenced_object_manual_review(reports, &manifest_key, TableMetadataMaintenanceReason::UnsupportedManifestAvro); + return Ok(()); + }; + for (file_location, object_kind) in file_references { + let Some(file_key) = table_catalog_object_key_from_location(table_bucket, &file_location) else { + insert_referenced_object_report( + reports, + file_location, + object_kind, + TableMetadataMaintenanceObjectState::ManualReviewRequired, + TableMetadataMaintenanceReason::UnsupportedManifestAvro, + ); + continue; + }; + if table_maintenance_object_kind(namespace, table, warehouse_object_prefix, &file_key) != Some(object_kind.clone()) { + insert_referenced_object_report( + reports, + file_key, + object_kind, + TableMetadataMaintenanceObjectState::ManualReviewRequired, + TableMetadataMaintenanceReason::UnsupportedManifestAvro, + ); + continue; + } + insert_referenced_object_report( + reports, + file_key, + object_kind.clone(), + TableMetadataMaintenanceObjectState::Retained, + table_metadata_maintenance_reason_for_object_kind(&object_kind), + ); + } + + Ok(()) +} + +fn mark_referenced_object_manual_review( + reports: &mut BTreeMap, + object_location: &str, + reason: TableMetadataMaintenanceReason, +) { + if let Some(report) = reports.get_mut(object_location) { + report.state = TableMetadataMaintenanceObjectState::ManualReviewRequired; + report.reasons.insert(reason); + } +} + +fn manifest_paths_from_manifest_list_avro(data: &[u8]) -> TableCatalogStoreResult> { + let reader = apache_avro::Reader::new(data) + .map_err(|err| TableCatalogStoreError::Invalid(format!("failed to read manifest list Avro: {err}")))?; + let mut manifest_paths = Vec::new(); + for value in reader { + let value = + value.map_err(|err| TableCatalogStoreError::Invalid(format!("failed to read manifest list record: {err}")))?; + let manifest_path = avro_record_field(&value, "manifest_path") + .and_then(avro_string_value) + .ok_or_else(|| TableCatalogStoreError::Invalid("manifest list entry missing manifest_path".to_string()))?; + manifest_paths.push(manifest_path.to_string()); + } + Ok(manifest_paths) +} + +fn file_references_from_manifest_avro(data: &[u8]) -> TableCatalogStoreResult> { + let reader = apache_avro::Reader::new(data) + .map_err(|err| TableCatalogStoreError::Invalid(format!("failed to read manifest Avro: {err}")))?; + let mut files = Vec::new(); + for value in reader { + let value = value.map_err(|err| TableCatalogStoreError::Invalid(format!("failed to read manifest record: {err}")))?; + let data_file = avro_record_field(&value, "data_file") + .ok_or_else(|| TableCatalogStoreError::Invalid("manifest entry missing data_file".to_string()))?; + let file_path = avro_record_field(data_file, "file_path") + .and_then(avro_string_value) + .ok_or_else(|| TableCatalogStoreError::Invalid("manifest data file missing file_path".to_string()))?; + let content = avro_record_field(data_file, "content") + .and_then(avro_i32_value) + .ok_or_else(|| TableCatalogStoreError::Invalid("manifest data file missing content".to_string()))?; + let object_kind = match content { + 0 => TableMetadataMaintenanceObjectKind::DataFile, + 1 | 2 => TableMetadataMaintenanceObjectKind::DeleteFile, + _ => continue, + }; + files.push((file_path.to_string(), object_kind)); + } + Ok(files) +} + +fn avro_record_field<'a>(value: &'a apache_avro::types::Value, name: &str) -> Option<&'a apache_avro::types::Value> { + let value = avro_non_union_value(value); + let apache_avro::types::Value::Record(fields) = value else { + return None; + }; + fields + .iter() + .find_map(|(field_name, field_value)| (field_name == name).then_some(avro_non_union_value(field_value))) +} + +fn avro_non_union_value(value: &apache_avro::types::Value) -> &apache_avro::types::Value { + match value { + apache_avro::types::Value::Union(_, inner) => avro_non_union_value(inner), + value => value, + } +} + +fn avro_string_value(value: &apache_avro::types::Value) -> Option<&str> { + match avro_non_union_value(value) { + apache_avro::types::Value::String(value) => Some(value), + _ => None, + } +} + +fn avro_i32_value(value: &apache_avro::types::Value) -> Option { + match avro_non_union_value(value) { + apache_avro::types::Value::Int(value) => Some(*value), + _ => None, + } +} + +fn table_catalog_object_key_from_location(table_bucket: &str, location: &str) -> Option { + let object = if let Some(location) = location.strip_prefix("s3://") { + let (bucket, object) = location.split_once('/')?; + if bucket != table_bucket { + return None; + } + object + } else { + location + }; + + if object.is_empty() + || object.starts_with('/') + || object.contains("..") + || object.contains('\\') + || object.bytes().any(|byte| byte.is_ascii_control()) + { + return None; + } + + Some(object.to_string()) +} + +fn table_maintenance_object_kind( + namespace: &Namespace, + table: &IdentifierSegment, + warehouse_object_prefix: Option<&str>, + object_location: &str, +) -> Option { + let metadata_prefix = format!("{}/", default_table_metadata_dir_path(namespace, table)); + if let Some(kind) = table_maintenance_metadata_object_kind(&metadata_prefix, object_location) { + return Some(kind); + } + + let data_prefix = format!("{}/", default_table_data_dir_path(namespace, table)); + if object_location + .strip_prefix(&data_prefix) + .is_some_and(is_valid_table_maintenance_nested_object) + { + return Some(TableMetadataMaintenanceObjectKind::DataFile); + } + + let delete_prefix = format!("{}/", default_table_delete_dir_path(namespace, table)); + if object_location + .strip_prefix(&delete_prefix) + .is_some_and(is_valid_table_maintenance_nested_object) + { + return Some(TableMetadataMaintenanceObjectKind::DeleteFile); + } + + if let Some(warehouse_object_prefix) = warehouse_object_prefix { + let metadata_prefix = format!("{warehouse_object_prefix}{METADATA_DIR}/"); + if let Some(kind) = table_maintenance_metadata_object_kind(&metadata_prefix, object_location) { + return Some(kind); + } + + let data_prefix = format!("{warehouse_object_prefix}{DATA_DIR}/"); + if object_location + .strip_prefix(&data_prefix) + .is_some_and(is_valid_table_maintenance_nested_object) + { + return Some(TableMetadataMaintenanceObjectKind::DataFile); + } + + let delete_prefix = format!("{warehouse_object_prefix}{DELETE_DIR}/"); + if object_location + .strip_prefix(&delete_prefix) + .is_some_and(is_valid_table_maintenance_nested_object) + { + return Some(TableMetadataMaintenanceObjectKind::DeleteFile); + } + } + + None +} + +fn table_maintenance_metadata_object_kind( + metadata_prefix: &str, + object_location: &str, +) -> Option { + let file_name = object_location.strip_prefix(metadata_prefix)?; + if file_name.is_empty() + || file_name.contains('/') + || file_name.contains('\\') + || file_name.contains("..") + || file_name.bytes().any(|byte| byte.is_ascii_control()) + || !file_name.ends_with(".avro") + { + return None; + } + if file_name.starts_with("snap-") { + return Some(TableMetadataMaintenanceObjectKind::ManifestList); + } + Some(TableMetadataMaintenanceObjectKind::ManifestFile) +} + +fn is_valid_table_maintenance_nested_object(suffix: &str) -> bool { + !suffix.is_empty() + && !suffix.starts_with('/') + && !suffix.contains("..") + && !suffix.contains('\\') + && !suffix.bytes().any(|byte| byte.is_ascii_control()) +} + +fn table_metadata_maintenance_reason_for_object_kind( + object_kind: &TableMetadataMaintenanceObjectKind, +) -> TableMetadataMaintenanceReason { + match object_kind { + TableMetadataMaintenanceObjectKind::MetadataFile => TableMetadataMaintenanceReason::CurrentMetadata, + TableMetadataMaintenanceObjectKind::ManifestList => TableMetadataMaintenanceReason::ManifestList, + TableMetadataMaintenanceObjectKind::ManifestFile => TableMetadataMaintenanceReason::ManifestFile, + TableMetadataMaintenanceObjectKind::DataFile => TableMetadataMaintenanceReason::DataFile, + TableMetadataMaintenanceObjectKind::DeleteFile => TableMetadataMaintenanceReason::DeleteFile, + } } fn metadata_maintenance_reachability_graph_report( @@ -3194,10 +3818,14 @@ fn metadata_maintenance_reachability_graph_report( let mut reasons = BTreeSet::from([TableMaintenanceReachabilityGraphReason::MetadataJsonParsed]); if manifest_list_count > 0 { reasons.insert(TableMaintenanceReachabilityGraphReason::ManifestListAvroReferenced); + } + if referenced_object_reports.iter().any(|report| { + report + .reasons + .contains(&TableMetadataMaintenanceReason::UnsupportedManifestAvro) + }) { reasons.insert(TableMaintenanceReachabilityGraphReason::ManifestAvroReaderUnavailable); } - reasons.insert(TableMaintenanceReachabilityGraphReason::DataFileCleanupDeferred); - reasons.insert(TableMaintenanceReachabilityGraphReason::DeleteFileCleanupDeferred); TableMaintenanceReachabilityGraphReport { status: if manual_review_count == 0 { @@ -3215,6 +3843,134 @@ fn metadata_maintenance_reachability_graph_report( } } +async fn metadata_maintenance_object_cleanup_reports( + backend: &B, + table_bucket: &str, + namespace: &Namespace, + table: &IdentifierSegment, + warehouse_object_prefix: Option<&str>, + referenced_object_reports: &[TableMetadataMaintenanceReferencedObjectReport], + now: OffsetDateTime, +) -> TableCatalogStoreResult<(usize, Vec, Vec, Vec)> +where + B: TableCatalogObjectBackend, +{ + let scanned_objects = + table_maintenance_cleanup_objects(backend, table_bucket, namespace, table, warehouse_object_prefix).await?; + if referenced_object_reports + .iter() + .any(|report| report.state == TableMetadataMaintenanceObjectState::ManualReviewRequired) + { + return Ok((scanned_objects.len(), Vec::new(), Vec::new(), Vec::new())); + } + + let referenced_locations = referenced_object_reports + .iter() + .filter_map(|report| table_catalog_object_key_from_location(table_bucket, &report.object_location)) + .collect::>(); + let mut cleanup_candidate_locations = Vec::new(); + let mut deletable_object_locations = Vec::new(); + let mut cleanup_reports = Vec::new(); + + for (object_location, object_kind) in scanned_objects { + if referenced_locations.contains(&object_location) { + continue; + } + let mut reasons = BTreeSet::from([ + table_metadata_maintenance_reason_for_object_kind(&object_kind), + TableMetadataMaintenanceReason::NoCurrentReachability, + ]); + let state = match backend.read_object(table_bucket, &object_location).await? { + Some(object) if metadata_candidate_is_past_safety_window(object.mod_time, now) => { + reasons.insert(TableMetadataMaintenanceReason::SafetyWindowSatisfied); + cleanup_candidate_locations.push(object_location.clone()); + deletable_object_locations.push(object_location.clone()); + TableMetadataMaintenanceObjectState::Deletable + } + _ => { + reasons.insert(TableMetadataMaintenanceReason::SafetyWindowPending); + cleanup_candidate_locations.push(object_location.clone()); + TableMetadataMaintenanceObjectState::PendingSafetyWindow + } + }; + cleanup_reports.push(TableMetadataMaintenanceObjectCleanupReport { + object_location, + object_kind, + state, + reasons: reasons.into_iter().collect(), + }); + } + + Ok(( + referenced_locations.len() + cleanup_reports.len(), + cleanup_candidate_locations, + deletable_object_locations, + cleanup_reports, + )) +} + +async fn table_maintenance_cleanup_objects( + backend: &B, + table_bucket: &str, + namespace: &Namespace, + table: &IdentifierSegment, + warehouse_object_prefix: Option<&str>, +) -> TableCatalogStoreResult> +where + B: TableCatalogObjectBackend, +{ + let mut objects = BTreeMap::new(); + let mut metadata_prefixes = vec![format!("{}/", default_table_metadata_dir_path(namespace, table))]; + let mut data_prefixes = vec![format!("{}/", default_table_data_dir_path(namespace, table))]; + let mut delete_prefixes = vec![format!("{}/", default_table_delete_dir_path(namespace, table))]; + if let Some(warehouse_object_prefix) = warehouse_object_prefix { + metadata_prefixes.push(format!("{warehouse_object_prefix}{METADATA_DIR}/")); + data_prefixes.push(format!("{warehouse_object_prefix}{DATA_DIR}/")); + delete_prefixes.push(format!("{warehouse_object_prefix}{DELETE_DIR}/")); + } + metadata_prefixes.sort(); + metadata_prefixes.dedup(); + data_prefixes.sort(); + data_prefixes.dedup(); + delete_prefixes.sort(); + delete_prefixes.dedup(); + + for metadata_prefix in metadata_prefixes { + for object in backend.list_objects(table_bucket, &metadata_prefix).await? { + if let Some(kind) = table_maintenance_object_kind(namespace, table, warehouse_object_prefix, &object) + && matches!( + kind, + TableMetadataMaintenanceObjectKind::ManifestList | TableMetadataMaintenanceObjectKind::ManifestFile + ) + { + objects.insert(object, kind); + } + } + } + + for data_prefix in data_prefixes { + for object in backend.list_objects(table_bucket, &data_prefix).await? { + if table_maintenance_object_kind(namespace, table, warehouse_object_prefix, &object) + == Some(TableMetadataMaintenanceObjectKind::DataFile) + { + objects.insert(object, TableMetadataMaintenanceObjectKind::DataFile); + } + } + } + + for delete_prefix in delete_prefixes { + for object in backend.list_objects(table_bucket, &delete_prefix).await? { + if table_maintenance_object_kind(namespace, table, warehouse_object_prefix, &object) + == Some(TableMetadataMaintenanceObjectKind::DeleteFile) + { + objects.insert(object, TableMetadataMaintenanceObjectKind::DeleteFile); + } + } + } + + Ok(objects) +} + fn mark_deleted_metadata_object_reports( object_reports: &mut [TableMetadataMaintenanceObjectReport], deleted_locations: &BTreeSet, @@ -3235,6 +3991,26 @@ fn mark_deleted_metadata_object_reports( } } +fn mark_deleted_object_cleanup_reports( + object_reports: &mut [TableMetadataMaintenanceObjectCleanupReport], + deleted_locations: &BTreeSet, +) { + for object_report in object_reports { + if !deleted_locations.contains(&object_report.object_location) { + continue; + } + object_report.state = TableMetadataMaintenanceObjectState::Deleted; + if !object_report + .reasons + .contains(&TableMetadataMaintenanceReason::DeletedByMaintenance) + { + object_report + .reasons + .push(TableMetadataMaintenanceReason::DeletedByMaintenance); + } + } +} + fn metadata_log_locations( current_metadata: &serde_json::Value, namespace: &Namespace, @@ -4375,6 +5151,14 @@ pub(crate) fn default_table_metadata_dir_path(namespace: &Namespace, table: &Ide format!("{}{}/{}", default_table_root_prefix(namespace), table.as_str(), METADATA_DIR) } +pub(crate) fn default_table_data_dir_path(namespace: &Namespace, table: &IdentifierSegment) -> String { + format!("{}{}/{}", default_table_root_prefix(namespace), table.as_str(), DATA_DIR) +} + +pub(crate) fn default_table_delete_dir_path(namespace: &Namespace, table: &IdentifierSegment) -> String { + format!("{}{}/{}", default_table_root_prefix(namespace), table.as_str(), DELETE_DIR) +} + pub(crate) fn default_table_metadata_file_path( namespace: &Namespace, table: &IdentifierSegment, @@ -4940,6 +5724,80 @@ mod tests { .expect("compaction planning report should include the snapshot") } + fn object_cleanup_report<'a>( + report: &'a TableMetadataMaintenanceReport, + object_location: &str, + ) -> &'a TableMetadataMaintenanceObjectCleanupReport { + report + .object_cleanup_reports + .iter() + .find(|object| object.object_location == object_location) + .expect("metadata maintenance object cleanup report should exist") + } + + fn manifest_list_avro_bytes(manifest_paths: &[&str]) -> Vec { + let schema = apache_avro::Schema::parse_str( + r#" + { + "type": "record", + "name": "manifest_file", + "fields": [ + {"name": "manifest_path", "type": "string"} + ] + } + "#, + ) + .expect("manifest list avro schema should parse"); + let mut writer = apache_avro::Writer::new(&schema, Vec::new()); + for manifest_path in manifest_paths { + writer + .append(apache_avro::types::Value::Record(vec![( + "manifest_path".to_string(), + apache_avro::types::Value::String((*manifest_path).to_string()), + )])) + .expect("manifest list record should append"); + } + writer.into_inner().expect("manifest list avro bytes should flush") + } + + fn manifest_avro_bytes(files: &[(&str, i32)]) -> Vec { + let schema = apache_avro::Schema::parse_str( + r#" + { + "type": "record", + "name": "manifest_entry", + "fields": [ + { + "name": "data_file", + "type": { + "type": "record", + "name": "data_file", + "fields": [ + {"name": "content", "type": "int"}, + {"name": "file_path", "type": "string"} + ] + } + } + ] + } + "#, + ) + .expect("manifest avro schema should parse"); + let mut writer = apache_avro::Writer::new(&schema, Vec::new()); + for (file_path, content) in files { + writer + .append(apache_avro::types::Value::Record(vec![( + "data_file".to_string(), + apache_avro::types::Value::Record(vec![ + ("content".to_string(), apache_avro::types::Value::Int(*content)), + ("file_path".to_string(), apache_avro::types::Value::String((*file_path).to_string())), + ]), + )])) + .expect("manifest record should append"); + } + writer.into_inner().expect("manifest avro bytes should flush") + } + #[async_trait::async_trait] impl TableCatalogObjectBackend for TestCatalogObjectBackend { async fn read_object(&self, bucket: &str, object: &str) -> TableCatalogStoreResult> { @@ -6216,6 +7074,412 @@ mod tests { ); } + #[tokio::test] + async fn maintenance_reachability_expands_manifest_avro_references() { + 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, "00002.metadata.json"); + let metadata_dir = default_table_metadata_dir_path(&namespace, &table); + let table_root = format!("{}{}/", default_table_root_prefix(&namespace), table.as_str()); + let manifest_list = format!("{metadata_dir}/snap-10.avro"); + let manifest = format!("{metadata_dir}/manifest-10.avro"); + let data_file = format!("{table_root}data/part-00001.parquet"); + let delete_file = format!("{table_root}delete/pos-00001.parquet"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend + .seed_object(bucket, &manifest_list, manifest_list_avro_bytes(&[&manifest])) + .await; + backend + .seed_object(bucket, &manifest, manifest_avro_bytes(&[(&data_file, 0), (&delete_file, 1)])) + .await; + backend.seed_object(bucket, &data_file, b"data".to_vec()).await; + backend.seed_object(bucket, &delete_file, b"delete".to_vec()).await; + backend + .seed_object( + bucket, + ¤t, + serde_json::to_vec(&serde_json::json!({ + "metadata-log": [], + "snapshots": [ + { + "snapshot-id": 10, + "manifest-list": manifest_list + } + ], + "refs": { + "main": { + "snapshot-id": 10, + "type": "branch" + } + } + })) + .unwrap(), + ) + .await; + + let report = store + .plan_table_metadata_maintenance(bucket, "sales", "orders", 0) + .await + .expect("metadata maintenance dry-run should succeed"); + + assert_eq!(report.reachability_graph.status, TableMaintenanceReachabilityGraphStatus::Complete); + assert_eq!(report.reachability_graph.manifest_list_count, 1); + assert_eq!(report.reachability_graph.manifest_file_count, 1); + assert_eq!(report.reachability_graph.data_file_count, 1); + assert_eq!(report.reachability_graph.delete_file_count, 1); + assert_eq!(report.reachability_graph.manual_review_count, 0); + assert!( + !report + .reachability_graph + .reasons + .contains(&TableMaintenanceReachabilityGraphReason::ManifestAvroReaderUnavailable) + ); + for (location, kind) in [ + (&manifest_list, TableMetadataMaintenanceObjectKind::ManifestList), + (&manifest, TableMetadataMaintenanceObjectKind::ManifestFile), + (&data_file, TableMetadataMaintenanceObjectKind::DataFile), + (&delete_file, TableMetadataMaintenanceObjectKind::DeleteFile), + ] { + let referenced = report + .referenced_object_reports + .iter() + .find(|object| object.object_location == *location) + .expect("referenced object should be reported"); + assert_eq!(referenced.object_kind, kind); + assert_eq!(referenced.state, TableMetadataMaintenanceObjectState::Retained); + } + assert!(report.cleanup_object_candidate_locations.is_empty()); + assert!(report.deletable_object_locations.is_empty()); + } + + #[tokio::test] + async fn maintenance_reachability_treats_v1_snapshot_manifests_as_reachable() { + 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, "00002.metadata.json"); + let metadata_dir = default_table_metadata_dir_path(&namespace, &table); + let data_dir = default_table_data_dir_path(&namespace, &table); + let manifest = format!("{metadata_dir}/manifest-10.avro"); + let data_file = format!("{data_dir}/part-00001.parquet"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend + .seed_object(bucket, &manifest, manifest_avro_bytes(&[(&data_file, 0)])) + .await; + backend.seed_object(bucket, &data_file, b"data".to_vec()).await; + backend + .seed_object( + bucket, + ¤t, + serde_json::to_vec(&serde_json::json!({ + "metadata-log": [], + "snapshots": [ + { + "snapshot-id": 10, + "manifests": [manifest] + } + ], + "refs": { + "main": { + "snapshot-id": 10, + "type": "branch" + } + } + })) + .unwrap(), + ) + .await; + + let report = store + .plan_table_metadata_maintenance(bucket, "sales", "orders", 0) + .await + .expect("metadata maintenance dry-run should succeed"); + + assert_eq!(report.reachability_graph.status, TableMaintenanceReachabilityGraphStatus::Complete); + assert_eq!(report.reachability_graph.manifest_file_count, 1); + assert_eq!(report.reachability_graph.data_file_count, 1); + assert!(report.cleanup_object_candidate_locations.is_empty()); + assert!(report.deletable_object_locations.is_empty()); + for location in [&manifest, &data_file] { + let referenced = report + .referenced_object_reports + .iter() + .find(|object| object.object_location == *location) + .expect("v1 manifest reference should be retained"); + assert_eq!(referenced.state, TableMetadataMaintenanceObjectState::Retained); + } + } + + #[tokio::test] + async fn maintenance_reachability_uses_table_warehouse_object_paths() { + 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, "00002.metadata.json"); + let manifest_list = "tables/table-id/metadata/snap-10.avro".to_string(); + let manifest = "tables/table-id/metadata/manifest-10.avro".to_string(); + let data_file = "tables/table-id/data/part-00001.parquet".to_string(); + let orphan_data = "tables/table-id/data/orphan.parquet".to_string(); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend + .seed_object(bucket, &manifest_list, manifest_list_avro_bytes(&[&manifest])) + .await; + backend + .seed_object(bucket, &manifest, manifest_avro_bytes(&[(&data_file, 0)])) + .await; + backend.seed_object(bucket, &data_file, b"data".to_vec()).await; + backend.seed_object(bucket, &orphan_data, b"orphan".to_vec()).await; + backend + .seed_object( + bucket, + ¤t, + serde_json::to_vec(&serde_json::json!({ + "metadata-log": [], + "snapshots": [ + { + "snapshot-id": 10, + "manifest-list": format!("s3://{bucket}/{manifest_list}") + } + ], + "refs": { + "main": { + "snapshot-id": 10, + "type": "branch" + } + } + })) + .unwrap(), + ) + .await; + + let report = store + .plan_table_metadata_maintenance(bucket, "sales", "orders", 0) + .await + .expect("metadata maintenance dry-run should succeed"); + + assert_eq!(report.reachability_graph.status, TableMaintenanceReachabilityGraphStatus::Complete); + assert!(report.referenced_object_reports.iter().any( + |object| object.object_location == manifest_list && object.state == TableMetadataMaintenanceObjectState::Retained + )); + assert!( + report.referenced_object_reports.iter().any( + |object| object.object_location == data_file && object.state == TableMetadataMaintenanceObjectState::Retained + ) + ); + assert_eq!(report.cleanup_object_candidate_locations, vec![orphan_data.clone()]); + assert_eq!(report.deletable_object_locations, vec![orphan_data]); + } + + #[tokio::test] + async fn maintenance_reachability_fails_closed_when_retained_metadata_is_unreadable() { + 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 old = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let current = default_table_metadata_file_path(&namespace, &table, "00002.metadata.json"); + let orphan_data = format!("{}/orphan.parquet", default_table_data_dir_path(&namespace, &table)); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend.seed_object(bucket, &old, b"not-json".to_vec()).await; + backend.seed_object(bucket, &orphan_data, b"orphan".to_vec()).await; + backend + .seed_object( + bucket, + ¤t, + serde_json::to_vec(&serde_json::json!({ + "metadata-log": [ + { + "timestamp-ms": 1, + "metadata-file": old + } + ], + "snapshots": [] + })) + .unwrap(), + ) + .await; + + let report = store + .plan_table_metadata_maintenance(bucket, "sales", "orders", 0) + .await + .expect("metadata maintenance dry-run should succeed"); + + assert_eq!( + report.reachability_graph.status, + TableMaintenanceReachabilityGraphStatus::ManualReviewRequired + ); + assert!(report.cleanup_object_candidate_locations.is_empty()); + assert!(report.deletable_object_locations.is_empty()); + let retained_metadata = report + .referenced_object_reports + .iter() + .find(|object| object.object_location == old) + .expect("unreadable retained metadata should be reported"); + assert_eq!(retained_metadata.object_kind, TableMetadataMaintenanceObjectKind::MetadataFile); + assert_eq!(retained_metadata.state, TableMetadataMaintenanceObjectState::ManualReviewRequired); + assert!( + retained_metadata + .reasons + .contains(&TableMetadataMaintenanceReason::UnreadableMetadata) + ); + } + + #[tokio::test] + async fn maintenance_dry_run_reports_unreachable_manifest_data_and_delete_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, "00002.metadata.json"); + let metadata_dir = default_table_metadata_dir_path(&namespace, &table); + let table_root = format!("{}{}/", default_table_root_prefix(&namespace), table.as_str()); + let manifest_list = format!("{metadata_dir}/snap-10.avro"); + let manifest = format!("{metadata_dir}/manifest-10.avro"); + let data_file = format!("{table_root}data/part-00001.parquet"); + let orphan_manifest = format!("{metadata_dir}/manifest-orphan.avro"); + let orphan_data = format!("{table_root}data/orphan.parquet"); + let orphan_delete = format!("{table_root}delete/orphan-delete.parquet"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend + .seed_object(bucket, &manifest_list, manifest_list_avro_bytes(&[&manifest])) + .await; + backend + .seed_object(bucket, &manifest, manifest_avro_bytes(&[(&data_file, 0)])) + .await; + backend.seed_object(bucket, &data_file, b"data".to_vec()).await; + backend.seed_object(bucket, &orphan_manifest, manifest_avro_bytes(&[])).await; + backend.seed_object(bucket, &orphan_data, b"orphan-data".to_vec()).await; + backend.seed_object(bucket, &orphan_delete, b"orphan-delete".to_vec()).await; + backend + .seed_object( + bucket, + ¤t, + serde_json::to_vec(&serde_json::json!({ + "metadata-log": [], + "snapshots": [ + { + "snapshot-id": 10, + "manifest-list": manifest_list + } + ], + "refs": { + "main": { + "snapshot-id": 10, + "type": "branch" + } + } + })) + .unwrap(), + ) + .await; + + let report = store + .plan_table_metadata_maintenance(bucket, "sales", "orders", 0) + .await + .expect("metadata maintenance dry-run should succeed"); + let candidates = report + .cleanup_object_candidate_locations + .iter() + .cloned() + .collect::>(); + let deletable = report.deletable_object_locations.iter().cloned().collect::>(); + let expected = BTreeSet::from([orphan_data.clone(), orphan_delete.clone(), orphan_manifest.clone()]); + + assert_eq!(candidates, expected); + assert_eq!(deletable, expected); + assert_eq!( + object_cleanup_report(&report, &orphan_manifest).object_kind, + TableMetadataMaintenanceObjectKind::ManifestFile + ); + assert_eq!( + object_cleanup_report(&report, &orphan_data).object_kind, + TableMetadataMaintenanceObjectKind::DataFile + ); + assert_eq!( + object_cleanup_report(&report, &orphan_delete).object_kind, + TableMetadataMaintenanceObjectKind::DeleteFile + ); + } + + #[tokio::test] + async fn maintenance_delete_removes_only_planned_unreachable_table_objects() { + 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, "00002.metadata.json"); + let metadata_dir = default_table_metadata_dir_path(&namespace, &table); + let table_root = format!("{}{}/", default_table_root_prefix(&namespace), table.as_str()); + let manifest_list = format!("{metadata_dir}/snap-10.avro"); + let manifest = format!("{metadata_dir}/manifest-10.avro"); + let data_file = format!("{table_root}data/part-00001.parquet"); + let orphan_manifest = format!("{metadata_dir}/manifest-orphan.avro"); + let orphan_data = format!("{table_root}data/orphan.parquet"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await; + backend + .seed_object(bucket, &manifest_list, manifest_list_avro_bytes(&[&manifest])) + .await; + backend + .seed_object(bucket, &manifest, manifest_avro_bytes(&[(&data_file, 0)])) + .await; + backend.seed_object(bucket, &data_file, b"data".to_vec()).await; + backend.seed_object(bucket, &orphan_manifest, manifest_avro_bytes(&[])).await; + backend.seed_object(bucket, &orphan_data, b"orphan-data".to_vec()).await; + backend + .seed_object( + bucket, + ¤t, + serde_json::to_vec(&serde_json::json!({ + "metadata-log": [], + "snapshots": [ + { + "snapshot-id": 10, + "manifest-list": manifest_list + } + ], + "refs": { + "main": { + "snapshot-id": 10, + "type": "branch" + } + } + })) + .unwrap(), + ) + .await; + + let deleted = store + .delete_table_metadata_maintenance_candidates(bucket, "sales", "orders", 0) + .await + .expect("maintenance delete should succeed"); + + assert_eq!( + deleted.deletable_object_locations.iter().cloned().collect::>(), + BTreeSet::from([orphan_data.clone(), orphan_manifest.clone()]) + ); + assert!(backend.object_exists(bucket, &manifest_list).await.unwrap()); + assert!(backend.object_exists(bucket, &manifest).await.unwrap()); + assert!(backend.object_exists(bucket, &data_file).await.unwrap()); + assert!(!backend.object_exists(bucket, &orphan_manifest).await.unwrap()); + assert!(!backend.object_exists(bucket, &orphan_data).await.unwrap()); + } + #[tokio::test] async fn maintenance_dry_run_keeps_metadata_log_references() { let backend = TestCatalogObjectBackend::default(); @@ -6724,12 +7988,18 @@ mod tests { "timestamp-ms": 1, "metadata-file": retained } + ], + "snapshots": [ + { + "snapshot-id": 1, + "manifest-list": manifest + } ] })) .unwrap(), ) .await; - backend.seed_object(bucket, &manifest, b"manifest".to_vec()).await; + backend.seed_object(bucket, &manifest, manifest_list_avro_bytes(&[])).await; let report = store .delete_table_metadata_maintenance_candidates(bucket, "sales", "orders", 0) @@ -6737,6 +8007,7 @@ mod tests { .unwrap(); assert_eq!(report.cleanup_candidate_locations, vec![old.clone()]); + assert!(report.cleanup_object_candidate_locations.is_empty()); assert!(!backend.object_exists(bucket, &old).await.unwrap()); assert!(backend.object_exists(bucket, &retained).await.unwrap()); assert!(backend.object_exists(bucket, ¤t).await.unwrap()); diff --git a/scripts/table-catalog/README.md b/scripts/table-catalog/README.md index a3211609c..a22dd09e2 100644 --- a/scripts/table-catalog/README.md +++ b/scripts/table-catalog/README.md @@ -150,10 +150,10 @@ 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 and heartbeat endpoints are registered; continuous in-process scheduling is not claimed -- manifest/data reachability cleanup: fail-closed reachability graph reporting only; metadata cleanup must not delete manifest, data, or delete files +- 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; built-in periodic scheduling is not claimed -- compaction rewrite: unsupported; planning reports fail closed until manifest Avro reading and rewrite support are implemented +- compaction rewrite: unsupported; planning reports fail closed until data file rewrite support is implemented - Iceberg views: stable unsupported routes are registered and return explicit unsupported JSON - external catalog bridges: metadata import/register is supported, but Polaris/Glue/DLF/Hive synchronization is 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 837891a02..40e855fae 100755 --- a/scripts/table-catalog/pyiceberg_smoke.py +++ b/scripts/table-catalog/pyiceberg_smoke.py @@ -140,9 +140,9 @@ UNSUPPORTED_INVENTORY: list[dict[str, str]] = [ }, { "capability": "manifest-data-reachability-cleanup", - "status": "fail-closed-reporting", + "status": "conservative-cleanup-supported", "roadmap_area": "reachability-cleanup", - "expected_behavior": "metadata maintenance reports expose reachability graph status and referenced manifest-list objects, but cleanup must not delete manifest, data, or delete files", + "expected_behavior": "metadata maintenance reads manifest-list and manifest Avro references, reports retained manifest/data/delete objects, and deletes only unreferenced table objects that pass the safety window", }, { "capability": "snapshot-expiration", @@ -154,7 +154,7 @@ UNSUPPORTED_INVENTORY: list[dict[str, str]] = [ "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", + "expected_behavior": "compaction planning fails closed until data file rewrite support is implemented; no data file rewrite or automatic compaction is claimed", }, { "capability": "iceberg-views",