diff --git a/crates/policy/src/policy/action.rs b/crates/policy/src/policy/action.rs index 402269fb5..45db55f75 100644 --- a/crates/policy/src/policy/action.rs +++ b/crates/policy/src/policy/action.rs @@ -508,6 +508,8 @@ pub enum AdminAction { ExportBucketMetadataAction, #[strum(serialize = "admin:GetTableCatalog")] GetTableCatalogAction, + #[strum(serialize = "admin:MigrateTableCatalog")] + MigrateTableCatalogAction, #[strum(serialize = "admin:GetTableBucket")] GetTableBucketAction, #[strum(serialize = "admin:SetTableBucket")] @@ -658,6 +660,7 @@ impl AdminAction { | AdminAction::ImportBucketMetadataAction | AdminAction::ExportBucketMetadataAction | AdminAction::GetTableCatalogAction + | AdminAction::MigrateTableCatalogAction | AdminAction::GetTableBucketAction | AdminAction::SetTableBucketAction | AdminAction::GetTableNamespaceAction @@ -832,9 +835,12 @@ mod tests { #[test] fn test_table_catalog_admin_action_is_valid() { let get_action = AdminAction::try_from("admin:GetTableCatalog").expect("Should parse GetTableCatalog action"); + let migrate_action = AdminAction::try_from("admin:MigrateTableCatalog").expect("Should parse MigrateTableCatalog action"); assert_eq!(get_action, AdminAction::GetTableCatalogAction); + assert_eq!(migrate_action, AdminAction::MigrateTableCatalogAction); assert!(get_action.is_valid()); + assert!(migrate_action.is_valid()); } #[test] diff --git a/docs/architecture/s3-tables-support-matrix.md b/docs/architecture/s3-tables-support-matrix.md index 5f112c748..bf77c5815 100644 --- a/docs/architecture/s3-tables-support-matrix.md +++ b/docs/architecture/s3-tables-support-matrix.md @@ -114,11 +114,12 @@ catalog extension. | Idempotent retry | Supported | Repeated commit IDs can return the already finalized result or surface recoverable finalization gaps. | | Post-CAS finalization recovery | Supported | Diagnostics and recovery can repair stale or missing idempotency indexes without changing the current table pointer. | | Catalog export | Supported | Exposes table state, commit recovery state, and backing migration information for operator inspection. | -| Strong backing migration contract | Supported as a contract | Diagnostics publish object-backed manifest state, recoverable commit-log WAL state, target backing type, replay requirements, and blockers. | -| Durable backing migration dry-run | Supported | `GET /iceberg/v1/{warehouse}/catalog/migration` and the `/_iceberg/v1` alias inspect object-backed catalog inventory, commit recovery blockers, idempotency indexes, warehouse prefix index readiness, recommended actions, and rollback configuration without mutating catalog state. | +| Strong backing state transfer | Supported | Object-backed table bucket, namespace, table, view, commit-log, and idempotency state can be materialized into the durable strong snapshot. The transfer is deterministic, ETag-CAS protected, idempotent after an interrupted finalization, and fails closed when a table or view has no owning namespace entry. | +| Durable backing migration preflight | Supported | `GET /iceberg/v1/{warehouse}/catalog/migration` and the `/_iceberg/v1` alias inspect object-backed catalog inventory, recovery blockers, warehouse prefix index readiness, persistent write-fence state, target snapshot agreement, and whether every table bucket is ready for cutover. | +| Durable backing migration execution | Preview / controlled | `POST /iceberg/v1/{warehouse}/catalog/migration` fences table-bucket registry changes, acquires a persistent per-bucket write fence, drains in-flight catalog mutations, materializes the target snapshot, and reports `ready_to_enable_durable_strong`. `DELETE` safely releases the bucket fence only while its target state has not advanced, and releases the registry fence after the last bucket is cancelled. Both mutations require `admin:MigrateTableCatalog`. | | Disaster recovery rehearsal | Manual/live harness | `failure_coverage.py --print-disaster-recovery-rehearsal` generates an operator runbook covering catalog export, diagnostics, safe recovery repair, rollback/import, durable backing migration dry-run, post-recovery loadTable, and table data-plane policy probes. | | Scale and fault rehearsal | Manual/live harness | `failure_coverage.py --print-scale-fault-rehearsal` generates an opt-in runbook for concurrent writer stress, maintenance scheduler lease recovery, durable backing cutover preflight, recovery/rollback/import under load, and post-run evidence capture. | -| Strong KV/WAL backing cutover | Preview / controlled | Operators can explicitly select durable strong backing with `RUSTFS_TABLE_CATALOG_BACKING=durable-strong`. Run the migration dry-run first and keep the object-backed catalog available for rollback. Object-only advanced operations fail closed in durable strong mode. | +| Strong KV/WAL backing cutover | Preview / controlled | Operators can select durable strong backing with `RUSTFS_TABLE_CATALOG_BACKING=durable-strong` only after every table bucket reports `SNAPSHOT_MATERIALIZED` and `ready_to_enable_durable_strong: true`. Object-only advanced operations fail closed in durable strong mode. | | Single active writer region | Supported policy | Diagnostics publish single-active-writer semantics and read-only replica limits. | | Active-active multi-region writes | Not claimed | A table must not accept independent concurrent writers in multiple active regions. | @@ -127,20 +128,29 @@ catalog extension. Use the migration dry-run before changing the table catalog backing for a warehouse: -1. Run `GET /iceberg/v1/{warehouse}/catalog/migration` with a principal that has - `GetTableCatalogAction` on the table bucket. -2. Treat any `blockers` entry as fail-closed. Run catalog recovery first when - commit recovery or idempotency repair is required, and backfill the warehouse - prefix index before switching backing modes. -3. Take an object-backed catalog snapshot or backup before setting - `RUSTFS_TABLE_CATALOG_BACKING=durable-strong`. -4. Restart with durable strong backing enabled, then verify catalog config, - table load, commit recovery diagnostics, and table data-plane policy - resolution for representative tables. -5. To roll back, restore `RUSTFS_TABLE_CATALOG_BACKING=object-backed` and restart - while preserving the object-backed catalog objects. Do not delete object-backed - catalog state until durable strong backing has passed the operator's retention - window. +1. Take an object-backed catalog backup and record the current metadata pointer + and version token for representative tables. +2. Run `GET /iceberg/v1/{warehouse}/catalog/migration` with a principal that has + `GetTableCatalogAction` on each table bucket. Treat every `blockers` entry as + fail-closed; repair commit recovery state and backfill the warehouse prefix + index before continuing. +3. Run `POST /iceberg/v1/{warehouse}/catalog/migration` with + `admin:MigrateTableCatalog`. This persists the source write fence before it + drains in-flight mutations and copies the catalog state. +4. Repeat the preflight and materialization for every table bucket. Do not set + `RUSTFS_TABLE_CATALOG_BACKING=durable-strong` until the preflight reports + `SNAPSHOT_MATERIALIZED`, no blockers, and + `ready_to_enable_durable_strong: true`. +5. Restart with durable strong backing enabled, then verify catalog config, + table and view loads, commit idempotency, and table data-plane policy + resolution before admitting writers. +6. Before restarting into durable-strong mode, `DELETE` on the migration + endpoint can remove a migration-created target bucket snapshot and release + the source fence. After the durable-strong state advances, cancellation + fails closed; recovery requires an operator-selected restore or reverse + migration instead of restarting against the stale object-backed pointer. +7. Preserve the object-backed catalog backup until durable strong backing has + passed the operator's retention window. ## Production Failure Coverage diff --git a/rustfs/src/admin/handlers/table_catalog.rs b/rustfs/src/admin/handlers/table_catalog.rs index e7264ec66..24d0b4e2b 100644 --- a/rustfs/src/admin/handlers/table_catalog.rs +++ b/rustfs/src/admin/handlers/table_catalog.rs @@ -118,6 +118,8 @@ const TABLE_CATALOG_ENDPOINTS: &[&str] = &[ "PUT /buckets/{warehouse}", "GET /buckets/{warehouse}", "GET /{warehouse}/catalog/migration", + "POST /{warehouse}/catalog/migration", + "DELETE /{warehouse}/catalog/migration", "GET /{warehouse}/namespaces", "POST /{warehouse}/namespaces", "GET /{warehouse}/namespaces/{namespace}", @@ -165,6 +167,9 @@ static GET_CONFIG_HANDLER: GetCatalogConfigHandler = GetCatalogConfigHandler {}; static ENABLE_TABLE_BUCKET_HANDLER: EnableTableBucketHandler = EnableTableBucketHandler {}; static GET_TABLE_BUCKET_HANDLER: GetTableBucketHandler = GetTableBucketHandler {}; static GET_TABLE_CATALOG_MIGRATION_HANDLER: GetTableCatalogMigrationHandler = GetTableCatalogMigrationHandler {}; +static MATERIALIZE_TABLE_CATALOG_MIGRATION_HANDLER: MaterializeTableCatalogMigrationHandler = + MaterializeTableCatalogMigrationHandler {}; +static CANCEL_TABLE_CATALOG_MIGRATION_HANDLER: CancelTableCatalogMigrationHandler = CancelTableCatalogMigrationHandler {}; static LIST_NAMESPACES_HANDLER: RestListNamespacesHandler = RestListNamespacesHandler {}; static CREATE_NAMESPACE_HANDLER: RestCreateNamespaceHandler = RestCreateNamespaceHandler {}; static GET_NAMESPACE_HANDLER: RestGetNamespaceHandler = RestGetNamespaceHandler {}; @@ -832,6 +837,16 @@ fn register_table_catalog_prefix_routes(r: &mut S3Router, prefix format!("{prefix}/{{warehouse}}/catalog/migration").as_str(), AdminOperation(&GET_TABLE_CATALOG_MIGRATION_HANDLER), )?; + r.insert( + Method::POST, + format!("{prefix}/{{warehouse}}/catalog/migration").as_str(), + AdminOperation(&MATERIALIZE_TABLE_CATALOG_MIGRATION_HANDLER), + )?; + r.insert( + Method::DELETE, + format!("{prefix}/{{warehouse}}/catalog/migration").as_str(), + AdminOperation(&CANCEL_TABLE_CATALOG_MIGRATION_HANDLER), + )?; r.insert( Method::GET, format!("{prefix}/{{warehouse}}/namespaces").as_str(), @@ -5012,6 +5027,46 @@ impl Operation for GetTableCatalogMigrationHandler { } } +pub struct MaterializeTableCatalogMigrationHandler {} + +#[async_trait::async_trait] +impl Operation for MaterializeTableCatalogMigrationHandler { + async fn call(&self, req: S3Request, params: Params<'_, '_>) -> S3Result> { + let warehouse = warehouse_from_params(¶ms)?; + authorize_table_catalog_request(&req, AdminAction::MigrateTableCatalogAction).await?; + ensure_table_bucket_enabled(&warehouse).await?; + let store = table_catalog_object_store()?; + let started = Instant::now(); + let result = store + .materialize_durable_strong_backing_migration(&warehouse) + .await + .map_err(catalog_store_error); + record_table_catalog_admin_operation_result("migration-materialize", &warehouse, "", "", started, &result); + let response = result?; + build_json_response(StatusCode::OK, &response) + } +} + +pub struct CancelTableCatalogMigrationHandler {} + +#[async_trait::async_trait] +impl Operation for CancelTableCatalogMigrationHandler { + async fn call(&self, req: S3Request, params: Params<'_, '_>) -> S3Result> { + let warehouse = warehouse_from_params(¶ms)?; + authorize_table_catalog_request(&req, AdminAction::MigrateTableCatalogAction).await?; + ensure_table_bucket_enabled(&warehouse).await?; + let store = table_catalog_object_store()?; + let started = Instant::now(); + let result = store + .cancel_durable_strong_backing_migration(&warehouse) + .await + .map_err(catalog_store_error); + record_table_catalog_admin_operation_result("migration-cancel", &warehouse, "", "", started, &result); + let response = result?; + build_json_response(StatusCode::OK, &response) + } +} + pub struct RestListNamespacesHandler {} #[async_trait::async_trait] @@ -5871,6 +5926,8 @@ mod tests { assert_eq!(response.admin_discovery.extensions_catalog, "/rustfs/admin/v4/extensions/catalog"); assert!(response.endpoints.contains(&"GET /v1/{prefix}/namespaces")); assert!(response.endpoints.contains(&"GET /{warehouse}/catalog/migration")); + assert!(response.endpoints.contains(&"POST /{warehouse}/catalog/migration")); + assert!(response.endpoints.contains(&"DELETE /{warehouse}/catalog/migration")); assert!(response.endpoints.contains(&"HEAD /v1/{prefix}/namespaces/{namespace}")); assert!( response @@ -6122,6 +6179,20 @@ mod tests { migration_block.contains("TableCatalogResource::warehouse(&warehouse)"), "catalog migration dry-run should authorize against the warehouse resource" ); + for handler in [ + "MaterializeTableCatalogMigrationHandler", + "CancelTableCatalogMigrationHandler", + ] { + let block = operation_block(src, handler); + assert!( + block.contains("authorize_table_catalog_request(&req, AdminAction::MigrateTableCatalogAction).await?;"), + "{handler} should require the global catalog migration action" + ); + assert!( + !block.contains("authorize_table_catalog_resource_request("), + "{handler} must not imply that a global backing cutover is warehouse-scoped" + ); + } for (handler, action) in [ ("RestLoadTableHandler", "AdminAction::GetTableMetadataAction"), @@ -6166,6 +6237,8 @@ mod tests { for handler in [ "GetTableCatalogMigrationHandler", + "MaterializeTableCatalogMigrationHandler", + "CancelTableCatalogMigrationHandler", "RestListNamespacesHandler", "RestCreateNamespaceHandler", "RestGetNamespaceHandler", @@ -6280,6 +6353,8 @@ mod tests { let _: &EnableTableBucketHandler = &ENABLE_TABLE_BUCKET_HANDLER; let _: &GetTableBucketHandler = &GET_TABLE_BUCKET_HANDLER; let _: &GetTableCatalogMigrationHandler = &GET_TABLE_CATALOG_MIGRATION_HANDLER; + let _: &MaterializeTableCatalogMigrationHandler = &MATERIALIZE_TABLE_CATALOG_MIGRATION_HANDLER; + let _: &CancelTableCatalogMigrationHandler = &CANCEL_TABLE_CATALOG_MIGRATION_HANDLER; let _: &RestListNamespacesHandler = &LIST_NAMESPACES_HANDLER; let _: &RestCreateNamespaceHandler = &CREATE_NAMESPACE_HANDLER; let _: &RestGetNamespaceHandler = &GET_NAMESPACE_HANDLER; @@ -6320,6 +6395,8 @@ mod tests { assert_operation::(); assert_operation::(); assert_operation::(); + assert_operation::(); + assert_operation::(); assert_operation::(); assert_operation::(); assert_operation::(); diff --git a/rustfs/src/admin/route_policy.rs b/rustfs/src/admin/route_policy.rs index b413742b6..af765cc70 100644 --- a/rustfs/src/admin/route_policy.rs +++ b/rustfs/src/admin/route_policy.rs @@ -65,6 +65,7 @@ const LIST_TEMPORARY_ACCOUNTS: AdminActionRef = AdminActionRef::new("ListTempora const LIST_TIER: AdminActionRef = AdminActionRef::new("ListTierAction"); const LIST_USER_POLICIES: AdminActionRef = AdminActionRef::new("ListUserPoliciesAdminAction"); const LIST_USERS: AdminActionRef = AdminActionRef::new("ListUsersAdminAction"); +const MIGRATE_TABLE_CATALOG: AdminActionRef = AdminActionRef::new("MigrateTableCatalogAction"); const PROFILING: AdminActionRef = AdminActionRef::new("ProfilingAdminAction"); const REBALANCE: AdminActionRef = AdminActionRef::new("RebalanceAdminAction"); const REGISTER_TABLE: AdminActionRef = AdminActionRef::new("RegisterTableAction"); @@ -789,6 +790,18 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ GET_TABLE_CATALOG, RouteRiskLevel::Sensitive, ), + admin( + HttpMethod::Post, + "/iceberg/v1/{warehouse}/catalog/migration", + MIGRATE_TABLE_CATALOG, + RouteRiskLevel::High, + ), + admin( + HttpMethod::Delete, + "/iceberg/v1/{warehouse}/catalog/migration", + MIGRATE_TABLE_CATALOG, + RouteRiskLevel::High, + ), admin( HttpMethod::Get, "/iceberg/v1/{warehouse}/namespaces", @@ -1054,6 +1067,18 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ GET_TABLE_CATALOG, RouteRiskLevel::Sensitive, ), + admin( + HttpMethod::Post, + "/_iceberg/v1/{warehouse}/catalog/migration", + MIGRATE_TABLE_CATALOG, + RouteRiskLevel::High, + ), + admin( + HttpMethod::Delete, + "/_iceberg/v1/{warehouse}/catalog/migration", + MIGRATE_TABLE_CATALOG, + RouteRiskLevel::High, + ), admin( HttpMethod::Get, "/_iceberg/v1/{warehouse}/namespaces", @@ -1531,7 +1556,7 @@ mod tests { let table_specs = ADMIN_ROUTE_POLICY_SPECS .iter() .filter(|spec| spec.path().starts_with("/iceberg/v1") || spec.path().starts_with("/_iceberg/v1")); - assert_eq!(table_specs.count(), 90); + assert_eq!(table_specs.count(), 94); assert_action(HttpMethod::Put, "/iceberg/v1/buckets/{warehouse}", SET_TABLE_BUCKET); assert_action(HttpMethod::Get, "/_iceberg/v1/buckets/{warehouse}", GET_TABLE_BUCKET); assert_action(HttpMethod::Get, "/iceberg/v1/{warehouse}/namespaces", GET_TABLE_NAMESPACE); @@ -1698,6 +1723,10 @@ mod tests { ); assert_action(HttpMethod::Get, "/iceberg/v1/{warehouse}/catalog/migration", GET_TABLE_CATALOG); assert_action(HttpMethod::Get, "/_iceberg/v1/{warehouse}/catalog/migration", GET_TABLE_CATALOG); + assert_action(HttpMethod::Post, "/iceberg/v1/{warehouse}/catalog/migration", MIGRATE_TABLE_CATALOG); + assert_action(HttpMethod::Post, "/_iceberg/v1/{warehouse}/catalog/migration", MIGRATE_TABLE_CATALOG); + assert_action(HttpMethod::Delete, "/iceberg/v1/{warehouse}/catalog/migration", MIGRATE_TABLE_CATALOG); + assert_action(HttpMethod::Delete, "/_iceberg/v1/{warehouse}/catalog/migration", MIGRATE_TABLE_CATALOG); assert_action( HttpMethod::Post, "/iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/catalog/import", diff --git a/rustfs/src/admin/route_registration_test.rs b/rustfs/src/admin/route_registration_test.rs index 18f3ffd36..e633e8fc9 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -342,6 +342,8 @@ fn expected_admin_route_matrix() -> Vec { table_route_sample(Method::PUT, "/buckets/{warehouse}", "/buckets/analytics"), table_route_sample(Method::GET, "/buckets/{warehouse}", "/buckets/analytics"), table_route_sample(Method::GET, "/{warehouse}/catalog/migration", "/analytics/catalog/migration"), + table_route_sample(Method::POST, "/{warehouse}/catalog/migration", "/analytics/catalog/migration"), + table_route_sample(Method::DELETE, "/{warehouse}/catalog/migration", "/analytics/catalog/migration"), table_route_sample(Method::GET, "/{warehouse}/namespaces", "/analytics/namespaces"), table_route_sample(Method::POST, "/{warehouse}/namespaces", "/analytics/namespaces"), table_route_sample(Method::GET, "/{warehouse}/namespaces/{namespace}", "/analytics/namespaces/sales"), @@ -531,6 +533,8 @@ fn expected_admin_route_matrix() -> Vec { compat_table_route_sample(Method::PUT, "/buckets/{warehouse}", "/buckets/analytics"), compat_table_route_sample(Method::GET, "/buckets/{warehouse}", "/buckets/analytics"), compat_table_route_sample(Method::GET, "/{warehouse}/catalog/migration", "/analytics/catalog/migration"), + compat_table_route_sample(Method::POST, "/{warehouse}/catalog/migration", "/analytics/catalog/migration"), + compat_table_route_sample(Method::DELETE, "/{warehouse}/catalog/migration", "/analytics/catalog/migration"), compat_table_route_sample(Method::GET, "/{warehouse}/namespaces", "/analytics/namespaces"), compat_table_route_sample(Method::POST, "/{warehouse}/namespaces", "/analytics/namespaces"), compat_table_route_sample(Method::GET, "/{warehouse}/namespaces/{namespace}", "/analytics/namespaces/sales"), diff --git a/rustfs/src/table_catalog.rs b/rustfs/src/table_catalog.rs index e21b4ae77..7a517ace3 100644 --- a/rustfs/src/table_catalog.rs +++ b/rustfs/src/table_catalog.rs @@ -49,6 +49,7 @@ use http::HeaderMap; use metrics::{counter, histogram}; use rustfs_filemeta::FileInfo; use serde::{Deserialize, Serialize, de::DeserializeOwned}; +use sha2::{Digest, Sha256}; use time::{Duration, OffsetDateTime}; use tokio::io::AsyncReadExt; use uuid::Uuid; @@ -124,6 +125,12 @@ const ICEBERG_REF_MAX_REF_AGE_MS_FIELD: &str = "max-ref-age-ms"; const STRONG_TABLE_CATALOG_SNAPSHOT_VERSION: u16 = 1; const STRONG_TABLE_CATALOG_BACKING_ROOT: &str = "strong-backing"; const STRONG_TABLE_CATALOG_SNAPSHOT_FILE: &str = "snapshot.json"; +const TABLE_CATALOG_MIGRATION_VERSION: u16 = 1; +const TABLE_CATALOG_MIGRATION_ROOT: &str = "backing-migration"; +const TABLE_CATALOG_MIGRATION_FENCE_FILE: &str = "durable-strong-fence.json"; +const TABLE_CATALOG_MIGRATION_FENCE_LOCK: &str = "durable-strong-fence.lock"; +const TABLE_CATALOG_MIGRATION_GLOBAL_FENCE_FILE: &str = "durable-strong-global-fence.json"; +const TABLE_CATALOG_MIGRATION_GLOBAL_FENCE_LOCK: &str = "durable-strong-global-fence.lock"; type CatalogListObjectsV2Info = StorageListObjectsV2Info; type CatalogListObjectVersionsInfo = StorageListObjectVersionsInfo; @@ -1241,15 +1248,56 @@ pub(crate) struct TableCatalogBackingMigrationDryRunReport { pub idempotency_index_count: usize, pub warehouse_prefix_count: usize, pub warehouse_index_ready: bool, + pub object_backed_writes_fenced: bool, + pub ready_to_enable_durable_strong: bool, pub blockers: Vec, pub recommended_actions: Vec, pub rollback: TableCatalogBackingRollbackPlan, } +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub(crate) struct TableCatalogBackingMigrationExecutionReport { + pub table_bucket: String, + pub source_kind: TableCatalogBackingKind, + pub target_kind: TableCatalogBackingKind, + pub status: TableCatalogBackingMigrationExecutionStatus, + pub namespace_count: usize, + pub table_count: usize, + pub view_count: usize, + pub commit_log_count: usize, + pub idempotency_index_count: usize, + pub source_fingerprint: String, + pub target_snapshot_etag: String, + pub object_backed_writes_fenced: bool, + pub ready_to_enable_durable_strong: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[serde(rename_all = "SCREAMING_SNAKE_CASE")] +pub(crate) enum TableCatalogBackingMigrationExecutionStatus { + SnapshotMaterialized, + SnapshotAlreadyMaterialized, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub(crate) struct TableCatalogBackingMigrationCancelReport { + pub table_bucket: String, + pub status: TableCatalogBackingMigrationCancelStatus, + pub object_backed_writes_fenced: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[serde(rename_all = "SCREAMING_SNAKE_CASE")] +pub(crate) enum TableCatalogBackingMigrationCancelStatus { + FenceReleased, + NoMigrationFence, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize)] #[serde(rename_all = "SCREAMING_SNAKE_CASE")] pub(crate) enum TableCatalogBackingMigrationStatus { ReadyToSnapshot, + SnapshotMaterialized, RecoveryRequired, ManualReviewRequired, } @@ -1271,6 +1319,7 @@ pub(crate) enum TableCatalogBackingMigrationBlocker { CommitManualReviewRequired, WarehouseIndexBackfillRequired, DuplicateWarehousePrefix, + DurableStrongSnapshotChanged, } #[derive(Debug, Clone, PartialEq, Eq, Serialize)] @@ -1282,6 +1331,8 @@ pub(crate) enum TableCatalogBackingMigrationAction { SnapshotObjectBackedCatalog, EnableDurableStrongBacking, VerifyDurableStrongSnapshot, + SnapshotRemainingTableBuckets, + ReviewDurableStrongSnapshot, KeepObjectBackedRollbackConfig, } @@ -1842,6 +1893,10 @@ pub(crate) trait TableCatalogObjectBackend: Clone + Send + Sync + 'static { async fn list_objects(&self, bucket: &str, prefix: &str) -> TableCatalogStoreResult>; + async fn acquire_read_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult> { + self.acquire_write_lock(bucket, object).await + } + async fn acquire_write_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult>; } @@ -1859,6 +1914,10 @@ impl Default for TableCatalogObjectPaths { } impl TableCatalogObjectPaths { + pub fn table_bucket_entries_prefix(&self) -> String { + format!("{}/{}/", self.catalog_root, TABLE_BUCKET_ROOT) + } + pub fn table_bucket_entry_path(&self, table_bucket: &str) -> String { format!("{}{}", self.table_bucket_root_prefix(table_bucket), TABLE_BUCKET_ENTRY_FILE) } @@ -2055,6 +2114,38 @@ impl TableCatalogObjectPaths { ) } + pub fn backing_migration_fence_path(&self, table_bucket: &str) -> String { + format!( + "{}{}/{}", + self.table_bucket_root_prefix(table_bucket), + TABLE_CATALOG_MIGRATION_ROOT, + TABLE_CATALOG_MIGRATION_FENCE_FILE + ) + } + + pub fn backing_migration_fence_lock_path(&self, table_bucket: &str) -> String { + format!( + "{}{}/{}", + self.table_bucket_root_prefix(table_bucket), + TABLE_CATALOG_MIGRATION_ROOT, + TABLE_CATALOG_MIGRATION_FENCE_LOCK + ) + } + + pub fn backing_migration_global_fence_path(&self) -> String { + format!( + "{}/{}/{}", + self.catalog_root, TABLE_CATALOG_MIGRATION_ROOT, TABLE_CATALOG_MIGRATION_GLOBAL_FENCE_FILE + ) + } + + pub fn backing_migration_global_fence_lock_path(&self) -> String { + format!( + "{}/{}/{}", + self.catalog_root, TABLE_CATALOG_MIGRATION_ROOT, TABLE_CATALOG_MIGRATION_GLOBAL_FENCE_LOCK + ) + } + fn table_bucket_root_prefix(&self, table_bucket: &str) -> String { format!("{}/{}/{}/", self.catalog_root, TABLE_BUCKET_ROOT, table_catalog_path_hash(table_bucket)) } @@ -2177,7 +2268,7 @@ struct StrongTableCatalogState { warehouse_index: StrongWarehouseIndex, } -#[derive(Debug, Clone, Serialize, Deserialize)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] struct StrongCommitSnapshotRecord { table_bucket: String, table_id: String, @@ -2185,7 +2276,7 @@ struct StrongCommitSnapshotRecord { commit: CommitLogEntry, } -#[derive(Debug, Clone, Serialize, Deserialize)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] struct StrongTableCatalogSnapshot { version: u16, table_buckets: Vec, @@ -2196,6 +2287,48 @@ struct StrongTableCatalogSnapshot { idempotency: Vec, } +#[derive(Debug, Clone, PartialEq, Serialize)] +struct StrongTableCatalogBucketSnapshot { + table_bucket: TableBucketEntry, + namespaces: Vec, + tables: Vec, + views: Vec, + commits: Vec, + idempotency: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "SCREAMING_SNAKE_CASE")] +enum TableCatalogBackingMigrationFenceStatus { + Preparing, + Materialized, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct TableCatalogBackingMigrationGlobalFence { + version: u16, + migration_id: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct TableCatalogBackingMigrationFence { + version: u16, + table_bucket: String, + migration_id: String, + status: TableCatalogBackingMigrationFenceStatus, + target_bucket_existed: bool, + source_fingerprint: Option, + target_snapshot_etag: Option, +} + +fn table_catalog_bucket_snapshot_fingerprint(snapshot: &StrongTableCatalogBucketSnapshot) -> TableCatalogStoreResult { + let data = serde_json::to_vec(snapshot) + .map_err(|err| TableCatalogStoreError::Internal(format!("failed to encode catalog migration snapshot: {err}")))?; + Ok(hex_simd::encode_to_string(Sha256::digest(data), hex_simd::AsciiCase::Lower)) +} + #[derive(Clone)] pub(crate) struct StrongTableCatalogStore { object_backend: B, @@ -2273,6 +2406,104 @@ where } } + fn bucket_snapshot_from_state_locked( + state: &StrongTableCatalogState, + table_bucket: &str, + ) -> Option { + let table_bucket_entry = state.table_buckets.get(table_bucket)?.clone(); + Some(StrongTableCatalogBucketSnapshot { + table_bucket: table_bucket_entry, + namespaces: state + .namespaces + .iter() + .filter(|((entry_bucket, _), _)| entry_bucket == table_bucket) + .map(|(_, entry)| entry.clone()) + .collect(), + tables: state + .tables + .iter() + .filter(|((entry_bucket, _, _), _)| entry_bucket == table_bucket) + .map(|(_, entry)| entry.clone()) + .collect(), + views: state + .views + .iter() + .filter(|((entry_bucket, _, _), _)| entry_bucket == table_bucket) + .map(|(_, entry)| entry.clone()) + .collect(), + commits: state + .commits + .iter() + .filter(|((entry_bucket, _, _), _)| entry_bucket == table_bucket) + .map(|((entry_bucket, table_id, lookup_key), commit)| StrongCommitSnapshotRecord { + table_bucket: entry_bucket.clone(), + table_id: table_id.clone(), + lookup_key: lookup_key.clone(), + commit: commit.clone(), + }) + .collect(), + idempotency: state + .idempotency + .iter() + .filter(|((entry_bucket, _, _), _)| entry_bucket == table_bucket) + .map(|((entry_bucket, table_id, lookup_key), commit)| StrongCommitSnapshotRecord { + table_bucket: entry_bucket.clone(), + table_id: table_id.clone(), + lookup_key: lookup_key.clone(), + commit: commit.clone(), + }) + .collect(), + }) + } + + fn remove_bucket_from_state_locked(state: &mut StrongTableCatalogState, table_bucket: &str) { + state.table_buckets.remove(table_bucket); + state.namespaces.retain(|(entry_bucket, _), _| entry_bucket != table_bucket); + state.tables.retain(|(entry_bucket, _, _), _| entry_bucket != table_bucket); + state.views.retain(|(entry_bucket, _, _), _| entry_bucket != table_bucket); + state.commits.retain(|(entry_bucket, _, _), _| entry_bucket != table_bucket); + state + .idempotency + .retain(|(entry_bucket, _, _), _| entry_bucket != table_bucket); + state.warehouse_index.remove(table_bucket); + } + + fn insert_bucket_snapshot_locked( + state: &mut StrongTableCatalogState, + snapshot: StrongTableCatalogBucketSnapshot, + ) -> TableCatalogStoreResult<()> { + let table_bucket = snapshot.table_bucket.table_bucket.clone(); + Self::remove_bucket_from_state_locked(state, &table_bucket); + state.table_buckets.insert(table_bucket.clone(), snapshot.table_bucket); + for entry in snapshot.namespaces { + let namespace = parse_namespace_for_store(&entry.namespace)?; + state.namespaces.insert(Self::namespace_key(&table_bucket, &namespace), entry); + } + for entry in snapshot.tables { + let namespace = parse_namespace_for_store(&entry.namespace)?; + let table = parse_table_for_store(&entry.table)?; + state.tables.insert(Self::table_key(&table_bucket, &namespace, &table), entry); + } + for entry in snapshot.views { + let namespace = parse_namespace_for_store(&entry.namespace)?; + let view = parse_table_for_store(&entry.view)?; + state.views.insert(Self::table_key(&table_bucket, &namespace, &view), entry); + } + for record in snapshot.commits { + state.commits.insert( + Self::commit_key(&record.table_bucket, &record.table_id, &record.lookup_key), + record.commit, + ); + } + for record in snapshot.idempotency { + state.idempotency.insert( + Self::idempotency_key(&record.table_bucket, &record.table_id, &record.lookup_key), + record.commit, + ); + } + Self::rebuild_warehouse_index_locked(state) + } + fn rebuild_warehouse_index_locked(state: &mut StrongTableCatalogState) -> TableCatalogStoreResult<()> { let mut warehouse_index: StrongWarehouseIndex = BTreeMap::new(); for ((table_bucket, namespace, table), entry) in &state.tables { @@ -2455,6 +2686,73 @@ where } } + async fn materialize_bucket_snapshot( + &self, + source: StrongTableCatalogBucketSnapshot, + ) -> TableCatalogStoreResult<(String, bool)> { + let _write_guard = self.write_lock.lock().await; + self.hydrate_state().await?; + let table_bucket = source.table_bucket.table_bucket.clone(); + let (snapshot, precondition) = { + let state = self.state.lock().await; + if let Some(current) = Self::bucket_snapshot_from_state_locked(&state, &table_bucket) { + if current != source { + return Err(TableCatalogStoreError::Conflict(format!( + "durable strong catalog already contains different state for table bucket {table_bucket}" + ))); + } + let snapshot_etag = state + .snapshot_etag + .clone() + .ok_or_else(|| TableCatalogStoreError::Internal("durable strong catalog snapshot has no etag".to_string()))?; + return Ok((snapshot_etag, false)); + } + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + Self::insert_bucket_snapshot_locked(&mut draft_state, source)?; + (Self::snapshot_from_mutated_state_locked(&mut draft_state)?, precondition) + }; + self.finalize_snapshot_write(snapshot, precondition).await?; + let state = self.state.lock().await; + let snapshot_etag = state + .snapshot_etag + .clone() + .ok_or_else(|| TableCatalogStoreError::Internal("durable strong catalog snapshot has no etag".to_string()))?; + Ok((snapshot_etag, true)) + } + + async fn remove_bucket_snapshot_if_unchanged( + &self, + table_bucket: &str, + expected_fingerprint: &str, + ) -> TableCatalogStoreResult<()> { + let _write_guard = self.write_lock.lock().await; + self.hydrate_state().await?; + let (snapshot, precondition) = { + let state = self.state.lock().await; + let Some(current) = Self::bucket_snapshot_from_state_locked(&state, table_bucket) else { + return Ok(()); + }; + if table_catalog_bucket_snapshot_fingerprint(¤t)? != expected_fingerprint { + return Err(TableCatalogStoreError::Conflict(format!( + "durable strong catalog state changed after materializing table bucket {table_bucket}" + ))); + } + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + Self::remove_bucket_from_state_locked(&mut draft_state, table_bucket); + (Self::snapshot_from_mutated_state_locked(&mut draft_state)?, precondition) + }; + self.finalize_snapshot_write(snapshot, precondition).await + } + + async fn bucket_snapshot_fingerprint(&self, table_bucket: &str) -> TableCatalogStoreResult> { + self.hydrate_state().await?; + let state = self.state.lock().await; + Self::bucket_snapshot_from_state_locked(&state, table_bucket) + .as_ref() + .map(table_catalog_bucket_snapshot_fingerprint) + .transpose() + } + fn require_table_bucket_in_state(state: &StrongTableCatalogState, table_bucket: &str) -> TableCatalogStoreResult<()> { if !state.table_buckets.contains_key(table_bucket) { return Err(TableCatalogStoreError::NotFound(format!("table bucket {table_bucket}"))); @@ -3260,6 +3558,400 @@ where RUSTFS_META_BUCKET } + async fn read_backing_migration_fence( + &self, + table_bucket: &str, + ) -> TableCatalogStoreResult)>> { + self.read_entry(self.catalog_bucket(), &self.paths.backing_migration_fence_path(table_bucket)) + .await + } + + async fn acquire_table_bucket_registry_write_permit(&self) -> TableCatalogStoreResult> { + let fence_path = self.paths.backing_migration_global_fence_path(); + let lock_path = self.paths.backing_migration_global_fence_lock_path(); + let guard = self.backend.acquire_read_lock(self.catalog_bucket(), &lock_path).await?; + if self + .read_entry::(self.catalog_bucket(), &fence_path) + .await? + .is_some() + { + return Err(TableCatalogStoreError::Conflict( + "table bucket registry writes are fenced while durable strong migration is in progress".to_string(), + )); + } + Ok(guard) + } + + async fn acquire_object_backed_catalog_write_permit(&self, table_bucket: &str) -> TableCatalogStoreResult> { + let lock_path = self.paths.backing_migration_fence_lock_path(table_bucket); + let guard = self.backend.acquire_read_lock(self.catalog_bucket(), &lock_path).await?; + if self.read_backing_migration_fence(table_bucket).await?.is_some() { + return Err(TableCatalogStoreError::Conflict(format!( + "object-backed catalog writes are fenced while table bucket {table_bucket} is prepared for durable strong cutover" + ))); + } + Ok(guard) + } + + async fn ensure_global_backing_migration_fence( + &self, + fence_path: &str, + ) -> TableCatalogStoreResult { + if let Some((fence, _)) = self + .read_entry::(self.catalog_bucket(), fence_path) + .await? + { + if fence.version != TABLE_CATALOG_MIGRATION_VERSION { + return Err(TableCatalogStoreError::Invalid( + "invalid durable strong global migration fence".to_string(), + )); + } + return Ok(fence); + } + let fence = TableCatalogBackingMigrationGlobalFence { + version: TABLE_CATALOG_MIGRATION_VERSION, + migration_id: Uuid::new_v4().to_string(), + }; + self.write_entry(self.catalog_bucket(), fence_path, &fence, TableCatalogPutPrecondition::IfAbsent) + .await?; + Ok(fence) + } + + async fn clear_global_backing_migration_fence_if_unused(&self, fence_path: &str) -> TableCatalogStoreResult<()> { + let bucket_objects = self + .backend + .list_objects(self.catalog_bucket(), &self.paths.table_bucket_entries_prefix()) + .await?; + if bucket_objects + .iter() + .any(|object| object.ends_with(TABLE_CATALOG_MIGRATION_FENCE_FILE)) + { + return Ok(()); + } + self.backend.delete_object(self.catalog_bucket(), fence_path).await + } + + async fn ensure_object_backed_writes_allowed(&self, table_bucket: &str) -> TableCatalogStoreResult<()> { + if self + .backend + .object_exists(self.catalog_bucket(), &self.paths.backing_migration_fence_path(table_bucket)) + .await? + { + return Err(TableCatalogStoreError::Conflict(format!( + "object-backed catalog writes are fenced while table bucket {table_bucket} is prepared for durable strong cutover" + ))); + } + Ok(()) + } + + async fn collect_bucket_snapshot_with_locks( + &self, + table_bucket: &str, + guards: &mut Vec>, + ) -> TableCatalogStoreResult { + let bucket_path = self.paths.table_bucket_entry_path(table_bucket); + guards.push(self.backend.acquire_write_lock(self.catalog_bucket(), &bucket_path).await?); + let Some((table_bucket_entry, _)) = self + .read_entry_unlocked::(self.catalog_bucket(), &bucket_path) + .await? + else { + return Err(TableCatalogStoreError::NotFound(format!("table bucket {table_bucket}"))); + }; + if table_bucket_entry.table_bucket != table_bucket { + return Err(TableCatalogStoreError::Invalid(format!( + "table bucket entry does not match migration target {table_bucket}" + ))); + } + + let mut namespaces = Vec::new(); + let mut tables = Vec::new(); + let mut views = Vec::new(); + let mut commits = Vec::new(); + let mut idempotency = Vec::new(); + let namespace_objects = self + .backend + .list_objects(self.catalog_bucket(), &self.paths.namespace_entries_prefix(table_bucket)) + .await?; + let mut unmatched_table_objects = namespace_objects + .iter() + .filter(|object| object.ends_with(TABLE_ENTRY_FILE)) + .cloned() + .collect::>(); + let mut unmatched_view_objects = namespace_objects + .iter() + .filter(|object| object.ends_with(VIEW_ENTRY_FILE)) + .cloned() + .collect::>(); + for namespace_object in namespace_objects + .iter() + .filter(|object| object.ends_with(NAMESPACE_ENTRY_FILE)) + { + guards.push( + self.backend + .acquire_write_lock(self.catalog_bucket(), namespace_object) + .await?, + ); + let Some((namespace_entry, _)) = self + .read_entry_unlocked::(self.catalog_bucket(), namespace_object) + .await? + else { + return Err(TableCatalogStoreError::Conflict(format!( + "namespace changed while preparing durable strong snapshot: {namespace_object}" + ))); + }; + if namespace_entry.table_bucket != table_bucket { + return Err(TableCatalogStoreError::Invalid(format!( + "namespace {} belongs to a different table bucket", + namespace_entry.namespace + ))); + } + let namespace = parse_namespace_for_store(&namespace_entry.namespace)?; + + let table_objects = self + .backend + .list_objects(self.catalog_bucket(), &self.paths.table_entries_prefix(table_bucket, &namespace)) + .await?; + for table_object in table_objects.iter().filter(|object| object.ends_with(TABLE_ENTRY_FILE)) { + unmatched_table_objects.remove(table_object); + guards.push(self.backend.acquire_write_lock(self.catalog_bucket(), table_object).await?); + let Some((table_entry, _)) = self + .read_entry_unlocked::(self.catalog_bucket(), table_object) + .await? + else { + return Err(TableCatalogStoreError::Conflict(format!( + "table changed while preparing durable strong snapshot: {table_object}" + ))); + }; + if table_entry.table_bucket != table_bucket || table_entry.namespace != namespace_entry.namespace { + return Err(TableCatalogStoreError::Invalid(format!( + "table {} does not match its catalog namespace", + table_entry.table + ))); + } + + for commit_object in self + .backend + .list_objects( + self.catalog_bucket(), + &self.paths.commit_log_entries_prefix(table_bucket, &table_entry.table_id), + ) + .await? + .into_iter() + .filter(|object| object.ends_with(".json")) + { + let Some((commit, _)) = self + .read_entry_unlocked::(self.catalog_bucket(), &commit_object) + .await? + else { + return Err(TableCatalogStoreError::Conflict(format!( + "commit log changed while preparing durable strong snapshot: {commit_object}" + ))); + }; + commits.push(StrongCommitSnapshotRecord { + table_bucket: table_bucket.to_string(), + table_id: table_entry.table_id.clone(), + lookup_key: commit.commit_id.clone(), + commit, + }); + } + for idempotency_object in self + .backend + .list_objects( + self.catalog_bucket(), + &self + .paths + .commit_idempotency_entries_prefix(table_bucket, &table_entry.table_id), + ) + .await? + .into_iter() + .filter(|object| object.ends_with(".json")) + { + let Some((commit, _)) = self + .read_entry_unlocked::(self.catalog_bucket(), &idempotency_object) + .await? + else { + return Err(TableCatalogStoreError::Conflict(format!( + "idempotency index changed while preparing durable strong snapshot: {idempotency_object}" + ))); + }; + let lookup_key = commit.idempotency_key.clone().ok_or_else(|| { + TableCatalogStoreError::Invalid(format!("idempotency index {idempotency_object} has no idempotency key")) + })?; + idempotency.push(StrongCommitSnapshotRecord { + table_bucket: table_bucket.to_string(), + table_id: table_entry.table_id.clone(), + lookup_key, + commit, + }); + } + tables.push(table_entry); + } + + let view_objects = self + .backend + .list_objects(self.catalog_bucket(), &self.paths.view_entries_prefix(table_bucket, &namespace)) + .await?; + for view_object in view_objects.iter().filter(|object| object.ends_with(VIEW_ENTRY_FILE)) { + unmatched_view_objects.remove(view_object); + guards.push(self.backend.acquire_write_lock(self.catalog_bucket(), view_object).await?); + let Some((view_entry, _)) = self + .read_entry_unlocked::(self.catalog_bucket(), view_object) + .await? + else { + return Err(TableCatalogStoreError::Conflict(format!( + "view changed while preparing durable strong snapshot: {view_object}" + ))); + }; + if view_entry.table_bucket != table_bucket || view_entry.namespace != namespace_entry.namespace { + return Err(TableCatalogStoreError::Invalid(format!( + "view {} does not match its catalog namespace", + view_entry.view + ))); + } + views.push(view_entry); + } + namespaces.push(namespace_entry); + } + if let Some(object) = unmatched_table_objects.first() { + return Err(TableCatalogStoreError::Invalid(format!( + "table entry has no namespace entry during durable strong migration: {object}" + ))); + } + if let Some(object) = unmatched_view_objects.first() { + return Err(TableCatalogStoreError::Invalid(format!( + "view entry has no namespace entry during durable strong migration: {object}" + ))); + } + + namespaces.sort_by(|left, right| left.namespace.cmp(&right.namespace)); + tables.sort_by(|left, right| (&left.namespace, &left.table).cmp(&(&right.namespace, &right.table))); + views.sort_by(|left, right| (&left.namespace, &left.view).cmp(&(&right.namespace, &right.view))); + commits.sort_by(|left, right| (&left.table_id, &left.lookup_key).cmp(&(&right.table_id, &right.lookup_key))); + idempotency.sort_by(|left, right| (&left.table_id, &left.lookup_key).cmp(&(&right.table_id, &right.lookup_key))); + + let snapshot = StrongTableCatalogBucketSnapshot { + table_bucket: table_bucket_entry, + namespaces, + tables, + views, + commits, + idempotency, + }; + self.validate_bucket_snapshot_for_migration(&snapshot)?; + Ok(snapshot) + } + + fn validate_bucket_snapshot_for_migration(&self, snapshot: &StrongTableCatalogBucketSnapshot) -> TableCatalogStoreResult<()> { + let table_bucket = &snapshot.table_bucket.table_bucket; + let tables_by_id = snapshot + .tables + .iter() + .map(|table| (table.table_id.as_str(), table)) + .collect::>(); + if tables_by_id.len() != snapshot.tables.len() { + return Err(TableCatalogStoreError::Invalid( + "migration snapshot contains duplicate table ids".to_string(), + )); + } + let commits_by_key = snapshot + .commits + .iter() + .map(|record| ((record.table_id.as_str(), record.lookup_key.as_str()), &record.commit)) + .collect::>(); + if commits_by_key.len() != snapshot.commits.len() { + return Err(TableCatalogStoreError::Invalid( + "migration snapshot contains duplicate commit lookup keys".to_string(), + )); + } + let idempotency_by_key = snapshot + .idempotency + .iter() + .map(|record| ((record.table_id.as_str(), record.lookup_key.as_str()), &record.commit)) + .collect::>(); + if idempotency_by_key.len() != snapshot.idempotency.len() { + return Err(TableCatalogStoreError::Invalid( + "migration snapshot contains duplicate idempotency lookup keys".to_string(), + )); + } + for record in &snapshot.commits { + let table = tables_by_id.get(record.table_id.as_str()).ok_or_else(|| { + TableCatalogStoreError::Invalid(format!("commit {} has no table in migration snapshot", record.commit.commit_id)) + })?; + if record.table_bucket != *table_bucket + || record.commit.table_id != record.table_id + || record.lookup_key != record.commit.commit_id + { + return Err(TableCatalogStoreError::Invalid(format!( + "commit {} does not match its migration snapshot owner", + record.commit.commit_id + ))); + } + let indexed = record + .commit + .idempotency_key + .as_deref() + .and_then(|idempotency_key| idempotency_by_key.get(&(record.table_id.as_str(), idempotency_key)).copied()); + let recovery = table_commit_recovery_entry(table, &record.commit, indexed); + if recovery.recovery_state != TableCommitRecoveryState::Committed { + return Err(TableCatalogStoreError::Conflict(format!( + "commit {} requires catalog recovery before durable strong migration", + record.commit.commit_id + ))); + } + } + for record in &snapshot.idempotency { + let _table = tables_by_id.get(record.table_id.as_str()).ok_or_else(|| { + TableCatalogStoreError::Invalid(format!( + "idempotency index {} has no table in migration snapshot", + record.lookup_key + )) + })?; + if record.table_bucket != *table_bucket || record.commit.table_id != record.table_id { + return Err(TableCatalogStoreError::Invalid(format!( + "idempotency index {} does not match its migration snapshot owner", + record.lookup_key + ))); + } + if record.commit.idempotency_key.as_deref() != Some(record.lookup_key.as_str()) { + return Err(TableCatalogStoreError::Invalid(format!( + "idempotency index {} does not match its commit payload", + record.lookup_key + ))); + } + let committed = commits_by_key + .get(&(record.table_id.as_str(), record.commit.commit_id.as_str())) + .ok_or_else(|| { + TableCatalogStoreError::Invalid(format!( + "idempotency index {} has no commit record in migration snapshot", + record.lookup_key + )) + })?; + if *committed != &record.commit { + return Err(TableCatalogStoreError::Conflict(format!( + "idempotency index {} requires catalog recovery before durable strong migration", + record.lookup_key + ))); + } + } + + let mut state = StrongTableCatalogState { + hydrated: true, + ..StrongTableCatalogState::default() + }; + StrongTableCatalogStore::::insert_bucket_snapshot_locked(&mut state, snapshot.clone())?; + if state.namespaces.len() != snapshot.namespaces.len() + || state.tables.len() != snapshot.tables.len() + || state.views.len() != snapshot.views.len() + || state.commits.len() != snapshot.commits.len() + || state.idempotency.len() != snapshot.idempotency.len() + { + return Err(TableCatalogStoreError::Invalid( + "migration snapshot contains duplicate catalog identities".to_string(), + )); + } + Ok(()) + } + async fn read_entry(&self, bucket: &str, object: &str) -> TableCatalogStoreResult)>> where T: DeserializeOwned, @@ -3784,16 +4476,27 @@ where let namespace = parse_namespace_for_store(&entry.namespace)?; let table = parse_table_for_store(&entry.table)?; validate_table_warehouse_location(&entry.table_bucket, &entry.warehouse_location)?; - if self.get_namespace(&entry.table_bucket, &entry.namespace).await?.is_none() { + let _migration_guard = self.acquire_object_backed_catalog_write_permit(&entry.table_bucket).await?; + let namespace_path = self.paths.namespace_entry_path(&entry.table_bucket, &namespace); + let _namespace_guard = self + .backend + .acquire_write_lock(self.catalog_bucket(), &namespace_path) + .await?; + if self + .read_entry_unlocked::(self.catalog_bucket(), &namespace_path) + .await? + .is_none() + { return Err(TableCatalogStoreError::NotFound(format!( "namespace {}/{}", entry.table_bucket, entry.namespace ))); } - let reservation = self.reserve_table_warehouse_index(&entry).await?; let table_path = self.paths.table_entry_path(&entry.table_bucket, &namespace, &table); + let _table_guard = self.backend.acquire_write_lock(self.catalog_bucket(), &table_path).await?; + let reservation = self.reserve_table_warehouse_index(&entry).await?; let result = self - .write_entry(self.catalog_bucket(), &table_path, &entry, precondition) + .write_entry_unlocked(self.catalog_bucket(), &table_path, &entry, precondition) .await; if result.is_err() { self.delete_created_table_warehouse_index(&entry, reservation, "table entry write failed") @@ -3808,14 +4511,25 @@ where let namespace = parse_namespace_for_store(&entry.namespace)?; let view = parse_table_for_store(&entry.view)?; validate_view_warehouse_location(&entry.table_bucket, &entry.warehouse_location)?; - if self.get_namespace(&entry.table_bucket, &entry.namespace).await?.is_none() { + let _migration_guard = self.acquire_object_backed_catalog_write_permit(&entry.table_bucket).await?; + let namespace_path = self.paths.namespace_entry_path(&entry.table_bucket, &namespace); + let _namespace_guard = self + .backend + .acquire_write_lock(self.catalog_bucket(), &namespace_path) + .await?; + if self + .read_entry_unlocked::(self.catalog_bucket(), &namespace_path) + .await? + .is_none() + { return Err(TableCatalogStoreError::NotFound(format!( "namespace {}/{}", entry.table_bucket, entry.namespace ))); } let view_path = self.paths.view_entry_path(&entry.table_bucket, &namespace, &view); - self.write_entry(self.catalog_bucket(), &view_path, &entry, precondition) + let _view_guard = self.backend.acquire_write_lock(self.catalog_bucket(), &view_path).await?; + self.write_entry_unlocked(self.catalog_bucket(), &view_path, &entry, precondition) .await } @@ -3847,6 +4561,7 @@ where ) -> TableCatalogStoreResult { validate_catalog_entry_version("external catalog bridge", entry.version)?; self.require_table_bucket(&entry.table_bucket).await?; + self.ensure_object_backed_writes_allowed(&entry.table_bucket).await?; let namespace = parse_namespace_for_store(&entry.namespace)?; let table = parse_table_for_store(&entry.table)?; if self.get_namespace(&entry.table_bucket, &entry.namespace).await?.is_none() { @@ -3973,6 +4688,7 @@ where ) -> TableCatalogStoreResult { let namespace = parse_namespace_for_store(namespace)?; let table = parse_table_for_store(table)?; + let _migration_guard = self.acquire_object_backed_catalog_write_permit(table_bucket).await?; let table_path = self.paths.table_entry_path(table_bucket, &namespace, &table); let _guard = self.backend.acquire_write_lock(self.catalog_bucket(), &table_path).await?; let Some((entry, _)) = self.read_table_with_etag_unlocked(table_bucket, &namespace, &table).await? else { @@ -4054,6 +4770,7 @@ where ) -> TableCatalogStoreResult { validate_table_maintenance_config(&config)?; self.require_table_bucket(table_bucket).await?; + self.ensure_object_backed_writes_allowed(table_bucket).await?; let config_path = self.paths.table_bucket_maintenance_config_path(table_bucket); self.write_entry(self.catalog_bucket(), &config_path, &config, TableCatalogPutPrecondition::Any) .await?; @@ -4129,6 +4846,7 @@ where config: TableMaintenanceConfig, ) -> TableCatalogStoreResult { validate_table_maintenance_config(&config)?; + self.ensure_object_backed_writes_allowed(table_bucket).await?; 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); @@ -4154,6 +4872,7 @@ where report: &TableMetadataMaintenanceReport, ) -> TableCatalogStoreResult<()> { let report = table_maintenance_report_with_recommended_actions(report.clone()); + self.ensure_object_backed_writes_allowed(&report.job.table_bucket).await?; let namespace = parse_namespace_for_store(&report.job.namespace)?; let table = parse_table_for_store(&report.job.table)?; let job_path = self.paths.table_maintenance_job_path( @@ -5634,7 +6353,7 @@ where recommended_actions.push(TableCatalogBackingMigrationAction::ReviewDuplicateWarehousePrefixes); } - let status = if manual_review_count > 0 || duplicate_warehouse_prefix_count > 0 { + let mut status = if manual_review_count > 0 || duplicate_warehouse_prefix_count > 0 { TableCatalogBackingMigrationStatus::ManualReviewRequired } else if recovery_required_count > 0 || !warehouse_index_ready { TableCatalogBackingMigrationStatus::RecoveryRequired @@ -5642,13 +6361,38 @@ where TableCatalogBackingMigrationStatus::ReadyToSnapshot }; + let strong_store = StrongTableCatalogStore::new(self.backend.clone()); + let migration_fence = self.read_backing_migration_fence(table_bucket).await?.map(|(fence, _)| fence); + let object_backed_writes_fenced = migration_fence.is_some(); + if status == TableCatalogBackingMigrationStatus::ReadyToSnapshot + && let Some(fence) = migration_fence.as_ref() + && fence.status == TableCatalogBackingMigrationFenceStatus::Materialized + && let Some(source_fingerprint) = fence.source_fingerprint.as_deref() + { + if strong_store.bucket_snapshot_fingerprint(table_bucket).await?.as_deref() == Some(source_fingerprint) { + status = TableCatalogBackingMigrationStatus::SnapshotMaterialized; + } else { + status = TableCatalogBackingMigrationStatus::ManualReviewRequired; + blockers.push(TableCatalogBackingMigrationBlocker::DurableStrongSnapshotChanged); + recommended_actions.push(TableCatalogBackingMigrationAction::ReviewDurableStrongSnapshot); + } + } + + let ready_to_enable_durable_strong = status == TableCatalogBackingMigrationStatus::SnapshotMaterialized + && self.all_table_buckets_materialized(&strong_store).await?; if status == TableCatalogBackingMigrationStatus::ReadyToSnapshot { recommended_actions.extend([ TableCatalogBackingMigrationAction::SnapshotObjectBackedCatalog, - TableCatalogBackingMigrationAction::EnableDurableStrongBacking, - TableCatalogBackingMigrationAction::VerifyDurableStrongSnapshot, TableCatalogBackingMigrationAction::KeepObjectBackedRollbackConfig, ]); + } else if status == TableCatalogBackingMigrationStatus::SnapshotMaterialized { + recommended_actions.push(TableCatalogBackingMigrationAction::VerifyDurableStrongSnapshot); + if ready_to_enable_durable_strong { + recommended_actions.push(TableCatalogBackingMigrationAction::EnableDurableStrongBacking); + } else { + recommended_actions.push(TableCatalogBackingMigrationAction::SnapshotRemainingTableBuckets); + } + recommended_actions.push(TableCatalogBackingMigrationAction::KeepObjectBackedRollbackConfig); } Ok(TableCatalogBackingMigrationDryRunReport { @@ -5663,6 +6407,8 @@ where idempotency_index_count, warehouse_prefix_count: warehouse_prefix_owners.len(), warehouse_index_ready, + object_backed_writes_fenced, + ready_to_enable_durable_strong, blockers, recommended_actions, rollback: TableCatalogBackingRollbackPlan { @@ -5675,6 +6421,231 @@ where }) } + pub(crate) async fn materialize_durable_strong_backing_migration( + &self, + table_bucket: &str, + ) -> TableCatalogStoreResult { + let fence_path = self.paths.backing_migration_fence_path(table_bucket); + let fence_lock_path = self.paths.backing_migration_fence_lock_path(table_bucket); + let global_fence_path = self.paths.backing_migration_global_fence_path(); + let global_fence_lock_path = self.paths.backing_migration_global_fence_lock_path(); + if self.read_backing_migration_fence(table_bucket).await?.is_none() { + let preflight = self.plan_durable_strong_backing_migration(table_bucket).await?; + if preflight.status != TableCatalogBackingMigrationStatus::ReadyToSnapshot { + return Err(TableCatalogStoreError::Conflict(format!( + "table bucket {table_bucket} is not ready for durable strong snapshot materialization" + ))); + } + } + + let _global_fence_guard = self + .backend + .acquire_write_lock(self.catalog_bucket(), &global_fence_lock_path) + .await?; + let _fence_guard = self + .backend + .acquire_write_lock(self.catalog_bucket(), &fence_lock_path) + .await?; + let existing_fence = self.read_backing_migration_fence(table_bucket).await?; + let strong_store = StrongTableCatalogStore::new(self.backend.clone()); + if let Some((fence, _)) = existing_fence.as_ref() + && (fence.version != TABLE_CATALOG_MIGRATION_VERSION || fence.table_bucket != table_bucket) + { + return Err(TableCatalogStoreError::Invalid(format!( + "invalid durable strong migration fence for table bucket {table_bucket}" + ))); + } + + if !self.warehouse_index_ready(table_bucket).await? { + return Err(TableCatalogStoreError::Conflict(format!( + "table bucket {table_bucket} warehouse index must be backfilled before durable strong migration" + ))); + } + + let mut source_guards = Vec::new(); + let source = self + .collect_bucket_snapshot_with_locks(table_bucket, &mut source_guards) + .await?; + let source_fingerprint = table_catalog_bucket_snapshot_fingerprint(&source)?; + if let Some((fence, _)) = existing_fence.as_ref() + && fence.status == TableCatalogBackingMigrationFenceStatus::Materialized + && fence.source_fingerprint.as_deref() != Some(source_fingerprint.as_str()) + { + return Err(TableCatalogStoreError::Conflict(format!( + "object-backed catalog state no longer matches the materialized snapshot for table bucket {table_bucket}" + ))); + } + + self.ensure_global_backing_migration_fence(&global_fence_path).await?; + let (migration_id, target_bucket_existed) = if let Some((fence, _)) = existing_fence.as_ref() { + (fence.migration_id.clone(), fence.target_bucket_existed) + } else { + let target_bucket_existed = strong_store.bucket_snapshot_fingerprint(table_bucket).await?.is_some(); + let fence = TableCatalogBackingMigrationFence { + version: TABLE_CATALOG_MIGRATION_VERSION, + table_bucket: table_bucket.to_string(), + migration_id: Uuid::new_v4().to_string(), + status: TableCatalogBackingMigrationFenceStatus::Preparing, + target_bucket_existed, + source_fingerprint: None, + target_snapshot_etag: None, + }; + self.write_entry(self.catalog_bucket(), &fence_path, &fence, TableCatalogPutPrecondition::IfAbsent) + .await?; + (fence.migration_id, target_bucket_existed) + }; + + let (target_snapshot_etag, created) = strong_store.materialize_bucket_snapshot(source.clone()).await?; + let completed_fence = TableCatalogBackingMigrationFence { + version: TABLE_CATALOG_MIGRATION_VERSION, + table_bucket: table_bucket.to_string(), + migration_id, + status: TableCatalogBackingMigrationFenceStatus::Materialized, + target_bucket_existed, + source_fingerprint: Some(source_fingerprint.clone()), + target_snapshot_etag: Some(target_snapshot_etag.clone()), + }; + self.write_entry(self.catalog_bucket(), &fence_path, &completed_fence, TableCatalogPutPrecondition::Any) + .await?; + + drop(source_guards); + drop(_fence_guard); + let ready_to_enable_durable_strong = self.all_table_buckets_materialized(&strong_store).await?; + Ok(TableCatalogBackingMigrationExecutionReport { + table_bucket: table_bucket.to_string(), + source_kind: TableCatalogBackingKind::ObjectBacked, + target_kind: TableCatalogBackingKind::StrongKvWal, + status: if created { + TableCatalogBackingMigrationExecutionStatus::SnapshotMaterialized + } else { + TableCatalogBackingMigrationExecutionStatus::SnapshotAlreadyMaterialized + }, + namespace_count: source.namespaces.len(), + table_count: source.tables.len(), + view_count: source.views.len(), + commit_log_count: source.commits.len(), + idempotency_index_count: source.idempotency.len(), + source_fingerprint, + target_snapshot_etag, + object_backed_writes_fenced: true, + ready_to_enable_durable_strong, + }) + } + + pub(crate) async fn cancel_durable_strong_backing_migration( + &self, + table_bucket: &str, + ) -> TableCatalogStoreResult { + let fence_path = self.paths.backing_migration_fence_path(table_bucket); + let fence_lock_path = self.paths.backing_migration_fence_lock_path(table_bucket); + let global_fence_path = self.paths.backing_migration_global_fence_path(); + let global_fence_lock_path = self.paths.backing_migration_global_fence_lock_path(); + let _global_fence_guard = self + .backend + .acquire_write_lock(self.catalog_bucket(), &global_fence_lock_path) + .await?; + let _fence_guard = self + .backend + .acquire_write_lock(self.catalog_bucket(), &fence_lock_path) + .await?; + let Some((fence, _)) = self.read_backing_migration_fence(table_bucket).await? else { + self.clear_global_backing_migration_fence_if_unused(&global_fence_path) + .await?; + return Ok(TableCatalogBackingMigrationCancelReport { + table_bucket: table_bucket.to_string(), + status: TableCatalogBackingMigrationCancelStatus::NoMigrationFence, + object_backed_writes_fenced: false, + }); + }; + self.ensure_global_backing_migration_fence(&global_fence_path).await?; + if fence.version != TABLE_CATALOG_MIGRATION_VERSION || fence.table_bucket != table_bucket { + return Err(TableCatalogStoreError::Invalid(format!( + "invalid durable strong migration fence for table bucket {table_bucket}" + ))); + } + + let mut source_guards = Vec::new(); + let source = self + .collect_bucket_snapshot_with_locks(table_bucket, &mut source_guards) + .await?; + let source_fingerprint = table_catalog_bucket_snapshot_fingerprint(&source)?; + if fence.status == TableCatalogBackingMigrationFenceStatus::Materialized + && fence.source_fingerprint.as_deref() != Some(source_fingerprint.as_str()) + { + return Err(TableCatalogStoreError::Conflict(format!( + "object-backed catalog state changed after materializing table bucket {table_bucket}" + ))); + } + + let strong_store = StrongTableCatalogStore::new(self.backend.clone()); + if fence.status == TableCatalogBackingMigrationFenceStatus::Materialized + && strong_store.bucket_snapshot_fingerprint(table_bucket).await?.as_deref() != Some(&source_fingerprint) + { + return Err(TableCatalogStoreError::Conflict(format!( + "durable strong catalog state changed after materializing table bucket {table_bucket}" + ))); + } + if !fence.target_bucket_existed { + strong_store + .remove_bucket_snapshot_if_unchanged(table_bucket, &source_fingerprint) + .await?; + } + self.backend.delete_object(self.catalog_bucket(), &fence_path).await?; + self.clear_global_backing_migration_fence_if_unused(&global_fence_path) + .await?; + Ok(TableCatalogBackingMigrationCancelReport { + table_bucket: table_bucket.to_string(), + status: TableCatalogBackingMigrationCancelStatus::FenceReleased, + object_backed_writes_fenced: false, + }) + } + + async fn all_table_buckets_materialized(&self, strong_store: &StrongTableCatalogStore) -> TableCatalogStoreResult { + if self + .read_entry::( + self.catalog_bucket(), + &self.paths.backing_migration_global_fence_path(), + ) + .await? + .is_none() + { + return Ok(false); + } + let table_bucket_objects = self + .backend + .list_objects(self.catalog_bucket(), &self.paths.table_bucket_entries_prefix()) + .await?; + for table_bucket_object in table_bucket_objects + .iter() + .filter(|object| object.ends_with(TABLE_BUCKET_ENTRY_FILE)) + { + let Some((entry, _)) = self + .read_entry::(self.catalog_bucket(), table_bucket_object) + .await? + else { + return Ok(false); + }; + let Some((fence, _)) = self.read_backing_migration_fence(&entry.table_bucket).await? else { + return Ok(false); + }; + if fence.status != TableCatalogBackingMigrationFenceStatus::Materialized { + return Ok(false); + } + let Some(source_fingerprint) = fence.source_fingerprint.as_deref() else { + return Ok(false); + }; + if strong_store + .bucket_snapshot_fingerprint(&entry.table_bucket) + .await? + .as_deref() + != Some(source_fingerprint) + { + return Ok(false); + } + } + Ok(true) + } + pub(crate) async fn diagnose_table_catalog( &self, table_bucket: &str, @@ -6391,20 +7362,21 @@ where if entry.catalog_type != TABLE_BUCKET_CATALOG_TYPE { return Err(TableCatalogStoreError::Invalid("unsupported table bucket catalog type".to_string())); } - - self.write_entry( - self.catalog_bucket(), - &self.paths.table_bucket_entry_path(&entry.table_bucket), - &entry, - TableCatalogPutPrecondition::Any, - ) - .await + let _registry_guard = self.acquire_table_bucket_registry_write_permit().await?; + let _migration_guard = self.acquire_object_backed_catalog_write_permit(&entry.table_bucket).await?; + let object = self.paths.table_bucket_entry_path(&entry.table_bucket); + let _guard = self.backend.acquire_write_lock(self.catalog_bucket(), &object).await?; + self.write_entry_unlocked(self.catalog_bucket(), &object, &entry, TableCatalogPutPrecondition::Any) + .await } async fn create_namespace(&self, entry: NamespaceEntry) -> TableCatalogStoreResult<()> { validate_catalog_entry_version("namespace", entry.version)?; self.require_table_bucket(&entry.table_bucket).await?; let namespace = parse_namespace_for_store(&entry.namespace)?; + let _migration_guard = self.acquire_object_backed_catalog_write_permit(&entry.table_bucket).await?; + let bucket_path = self.paths.table_bucket_entry_path(&entry.table_bucket); + let _bucket_guard = self.backend.acquire_write_lock(self.catalog_bucket(), &bucket_path).await?; let object = self.paths.namespace_entry_path(&entry.table_bucket, &namespace); self.write_entry(self.catalog_bucket(), &object, &entry, TableCatalogPutPrecondition::IfAbsent) .await @@ -6437,7 +7409,20 @@ where async fn drop_namespace(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult<()> { let namespace = parse_namespace_for_store(namespace)?; - if self.get_namespace(table_bucket, &namespace.public_name()).await?.is_none() { + let _migration_guard = self.acquire_object_backed_catalog_write_permit(table_bucket).await?; + let bucket_path = self.paths.table_bucket_entry_path(table_bucket); + let _bucket_guard = self.backend.acquire_write_lock(self.catalog_bucket(), &bucket_path).await?; + let namespace_path = self.paths.namespace_entry_path(table_bucket, &namespace); + // Match create_namespace and migration lock order while draining table/view creation. + let _namespace_guard = self + .backend + .acquire_write_lock(self.catalog_bucket(), &namespace_path) + .await?; + if self + .read_entry_unlocked::(self.catalog_bucket(), &namespace_path) + .await? + .is_none() + { return Err(TableCatalogStoreError::NotFound(format!( "namespace {}/{}", table_bucket, @@ -6459,7 +7444,7 @@ where ))); } self.backend - .delete_object(self.catalog_bucket(), &self.paths.namespace_entry_path(table_bucket, &namespace)) + .delete_object_unlocked(self.catalog_bucket(), &namespace_path) .await } @@ -6535,6 +7520,7 @@ where record_table_commit_attempt(&request.operation); let namespace = parse_namespace_for_store(&request.namespace)?; let table = parse_table_for_store(&request.table)?; + let _migration_guard = self.acquire_object_backed_catalog_write_permit(&request.table_bucket).await?; let table_path = self.paths.table_entry_path(&request.table_bucket, &namespace, &table); let _guard = self.backend.acquire_write_lock(self.catalog_bucket(), &table_path).await?; @@ -6817,6 +7803,12 @@ where async fn drop_table(&self, table_bucket: &str, namespace: &str, table: &str) -> TableCatalogStoreResult<()> { let namespace = parse_namespace_for_store(namespace)?; let table = parse_table_for_store(table)?; + let _migration_guard = self.acquire_object_backed_catalog_write_permit(table_bucket).await?; + let namespace_path = self.paths.namespace_entry_path(table_bucket, &namespace); + let _namespace_guard = self + .backend + .acquire_write_lock(self.catalog_bucket(), &namespace_path) + .await?; let object = self.paths.table_entry_path(table_bucket, &namespace, &table); let _guard = self.backend.acquire_write_lock(self.catalog_bucket(), &object).await?; let Some((entry, _)) = self.read_table_with_etag_unlocked(table_bucket, &namespace, &table).await? else { @@ -6889,6 +7881,12 @@ where async fn replace_view(&self, request: ViewCommitRequest) -> TableCatalogStoreResult { let namespace = parse_namespace_for_store(&request.namespace)?; let view = parse_table_for_store(&request.view)?; + let _migration_guard = self.acquire_object_backed_catalog_write_permit(&request.table_bucket).await?; + let namespace_path = self.paths.namespace_entry_path(&request.table_bucket, &namespace); + let _namespace_guard = self + .backend + .acquire_write_lock(self.catalog_bucket(), &namespace_path) + .await?; let view_path = self.paths.view_entry_path(&request.table_bucket, &namespace, &view); let _guard = self.backend.acquire_write_lock(self.catalog_bucket(), &view_path).await?; let Some((current, current_etag)) = self @@ -6948,9 +7946,16 @@ where async fn drop_view(&self, table_bucket: &str, namespace: &str, view: &str) -> TableCatalogStoreResult<()> { let namespace = parse_namespace_for_store(namespace)?; let view = parse_table_for_store(view)?; + let _migration_guard = self.acquire_object_backed_catalog_write_permit(table_bucket).await?; + let namespace_path = self.paths.namespace_entry_path(table_bucket, &namespace); + let _namespace_guard = self + .backend + .acquire_write_lock(self.catalog_bucket(), &namespace_path) + .await?; let object = self.paths.view_entry_path(table_bucket, &namespace, &view); + let _view_guard = self.backend.acquire_write_lock(self.catalog_bucket(), &object).await?; if self - .load_view(table_bucket, &namespace.public_name(), view.as_str()) + .read_entry_unlocked::(self.catalog_bucket(), &object) .await? .is_none() { @@ -6961,7 +7966,7 @@ where view.as_str() ))); } - self.backend.delete_object(self.catalog_bucket(), &object).await + self.backend.delete_object_unlocked(self.catalog_bucket(), &object).await } async fn get_commit_by_id( @@ -7507,6 +8512,19 @@ where .map_err(|err| TableCatalogStoreError::Internal(format!("failed to acquire catalog table lock: {err}")))?; Ok(Box::new(guard)) } + + async fn acquire_read_lock(&self, bucket: &str, object: &str) -> TableCatalogStoreResult> { + let lock = self + .store + .new_ns_lock(bucket, object) + .await + .map_err(|err| storage_error_to_catalog("create catalog migration lock", err))?; + let guard = lock + .get_read_lock(get_lock_acquire_timeout()) + .await + .map_err(|err| TableCatalogStoreError::Internal(format!("failed to acquire catalog migration lock: {err}")))?; + Ok(Box::new(guard)) + } } impl EcStoreTableCatalogObjectBackend @@ -16778,6 +17796,52 @@ mod tests { assert!(!export.backing_manifest.ha.active_active_supported); } + #[tokio::test] + async fn object_table_catalog_store_serializes_namespace_drop_with_table_creation() { + let backend = TestCatalogObjectBackend::default(); + let store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").expect("namespace should parse"); + let table = IdentifierSegment::parse("orders").expect("table should parse"); + let metadata_location = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + + store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap(); + store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .unwrap(); + + let table_path = store.paths.table_entry_path(bucket, &namespace, &table); + let pause = backend.pause_next_put(RUSTFS_META_BUCKET, &table_path).await; + let create_store = store.clone(); + let table_entry = test_table_entry(bucket, &namespace, &table, metadata_location); + let create = tokio::spawn(async move { create_store.create_table(table_entry).await }); + pause.wait_started().await; + + let drop_store = store.clone(); + let namespace_name = namespace.public_name(); + let mut drop_namespace = tokio::spawn(async move { drop_store.drop_namespace(bucket, &namespace_name).await }); + assert!( + tokio::time::timeout(std::time::Duration::from_millis(50), &mut drop_namespace) + .await + .is_err() + ); + + pause.release(); + create.await.unwrap().unwrap(); + assert_matches!( + drop_namespace.await.unwrap(), + Err(TableCatalogStoreError::Conflict(message)) if message.contains("is not empty") + ); + assert!( + store + .load_table(bucket, &namespace.public_name(), table.as_str()) + .await + .unwrap() + .is_some() + ); + } + #[tokio::test] async fn durable_strong_migration_dry_run_reports_ready_catalog_inventory() { let backend = TestCatalogObjectBackend::default(); @@ -16825,6 +17889,8 @@ mod tests { assert_eq!(report.commit_log_count, 1); assert_eq!(report.idempotency_index_count, 1); assert_eq!(report.warehouse_prefix_count, 1); + assert!(!report.object_backed_writes_fenced); + assert!(!report.ready_to_enable_durable_strong); assert!(report.blockers.is_empty()); assert!( report @@ -16832,7 +17898,7 @@ mod tests { .contains(&TableCatalogBackingMigrationAction::SnapshotObjectBackedCatalog) ); assert!( - report + !report .recommended_actions .contains(&TableCatalogBackingMigrationAction::EnableDurableStrongBacking) ); @@ -16840,6 +17906,479 @@ mod tests { assert_eq!(report.rollback.rollback_backing_value, TABLE_CATALOG_BACKING_OBJECT); } + #[tokio::test] + async fn durable_strong_migration_rejects_orphan_table_without_namespace() { + let backend = TestCatalogObjectBackend::default(); + let object_store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").expect("namespace should parse"); + let table = IdentifierSegment::parse("orders").expect("table should parse"); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + seed_table_for_metadata_maintenance(&object_store, bucket, &namespace, &table, current_metadata).await; + + let namespace_path = object_store.paths.namespace_entry_path(bucket, &namespace); + backend + .delete_object(RUSTFS_META_BUCKET, &namespace_path) + .await + .expect("namespace marker should be removed"); + assert!( + object_store + .load_table(bucket, &namespace.public_name(), table.as_str()) + .await + .unwrap() + .is_some() + ); + + let error = object_store + .materialize_durable_strong_backing_migration(bucket) + .await + .unwrap_err(); + assert_matches!( + error, + TableCatalogStoreError::Invalid(message) if message.contains("table entry has no namespace entry") + ); + object_store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .expect("failed migration must leave source catalog writable"); + object_store.put_table_bucket(test_bucket_entry("research")).await.unwrap(); + } + + #[tokio::test] + async fn durable_strong_migration_rejects_orphan_view_without_namespace() { + let backend = TestCatalogObjectBackend::default(); + let object_store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").expect("namespace should parse"); + let view = IdentifierSegment::parse("recent_orders").expect("view should parse"); + let metadata_location = default_view_metadata_file_path(&namespace, &view, "00001.metadata.json"); + + object_store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap(); + object_store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .unwrap(); + object_store + .create_view(test_view_entry(bucket, &namespace, &view, metadata_location)) + .await + .unwrap(); + object_store.backfill_table_warehouse_index(bucket).await.unwrap(); + + let namespace_path = object_store.paths.namespace_entry_path(bucket, &namespace); + backend + .delete_object(RUSTFS_META_BUCKET, &namespace_path) + .await + .expect("namespace marker should be removed"); + assert!( + object_store + .load_view(bucket, &namespace.public_name(), view.as_str()) + .await + .unwrap() + .is_some() + ); + + let error = object_store + .materialize_durable_strong_backing_migration(bucket) + .await + .unwrap_err(); + assert_matches!( + error, + TableCatalogStoreError::Invalid(message) if message.contains("view entry has no namespace entry") + ); + object_store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .expect("failed migration must leave source catalog writable"); + } + + #[tokio::test] + async fn durable_strong_migration_rejects_orphan_idempotency_without_fencing_source() { + let backend = TestCatalogObjectBackend::default(); + let object_store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + let table = IdentifierSegment::parse("orders").unwrap(); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + seed_table_for_metadata_maintenance(&object_store, bucket, &namespace, &table, current_metadata.clone()).await; + + let orphan = CommitLogEntry { + version: TABLE_CATALOG_ENTRY_VERSION, + commit_id: "orphan-commit".to_string(), + idempotency_key: Some("orphan-request".to_string()), + table_id: "table-id".to_string(), + operation: "append".to_string(), + expected_version_token: "token-v1".to_string(), + new_version_token: "token-v2".to_string(), + previous_metadata_location: current_metadata.clone(), + new_metadata_location: current_metadata, + requirements: Vec::new(), + status: CommitLogStatus::Committed, + writer: None, + created_at: None, + updated_at: None, + }; + let idempotency_path = object_store + .paths + .commit_idempotency_entry_path(bucket, "table-id", "orphan-request"); + backend + .seed_object(RUSTFS_META_BUCKET, &idempotency_path, serde_json::to_vec(&orphan).unwrap()) + .await; + + let error = object_store + .materialize_durable_strong_backing_migration(bucket) + .await + .unwrap_err(); + assert_matches!( + error, + TableCatalogStoreError::Invalid(message) if message.contains("has no commit record") + ); + + object_store + .create_namespace(test_namespace_entry(bucket, &Namespace::parse("still_writable").unwrap())) + .await + .unwrap(); + object_store.put_table_bucket(test_bucket_entry("research")).await.unwrap(); + } + + #[tokio::test] + async fn durable_strong_migration_materializes_catalog_state_and_fences_object_writes() { + let backend = TestCatalogObjectBackend::default(); + let object_store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + let table = IdentifierSegment::parse("orders").unwrap(); + let view = IdentifierSegment::parse("recent_orders").unwrap(); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let next_metadata = default_table_metadata_file_path(&namespace, &table, "00002.metadata.json"); + let strong_metadata = default_table_metadata_file_path(&namespace, &table, "00003.metadata.json"); + let view_metadata = default_view_metadata_file_path(&namespace, &view, "00001.metadata.json"); + + seed_table_for_metadata_maintenance(&object_store, bucket, &namespace, &table, current_metadata.clone()).await; + object_store + .create_view(test_view_entry(bucket, &namespace, &view, view_metadata)) + .await + .unwrap(); + backend.seed_object(bucket, &next_metadata, b"{}".to_vec()).await; + let object_commit = object_store + .commit_table(TableCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + commit_id: "object-commit".to_string(), + idempotency_key: Some("object-request".to_string()), + operation: "append".to_string(), + expected_version_token: "token-v1".to_string(), + expected_metadata_location: current_metadata, + new_metadata_location: next_metadata.clone(), + requirements: Vec::new(), + writer: Some("pyiceberg/test".to_string()), + }) + .await + .unwrap(); + + let materialized = object_store + .materialize_durable_strong_backing_migration(bucket) + .await + .unwrap(); + assert_eq!(materialized.status, TableCatalogBackingMigrationExecutionStatus::SnapshotMaterialized); + assert_eq!(materialized.namespace_count, 1); + assert_eq!(materialized.table_count, 1); + assert_eq!(materialized.view_count, 1); + assert_eq!(materialized.commit_log_count, 1); + assert_eq!(materialized.idempotency_index_count, 1); + assert!(materialized.object_backed_writes_fenced); + assert!(materialized.ready_to_enable_durable_strong); + + let retry = object_store + .materialize_durable_strong_backing_migration(bucket) + .await + .unwrap(); + assert_eq!(retry.status, TableCatalogBackingMigrationExecutionStatus::SnapshotAlreadyMaterialized); + assert_eq!(retry.source_fingerprint, materialized.source_fingerprint); + + let migration = object_store.plan_durable_strong_backing_migration(bucket).await.unwrap(); + assert_eq!(migration.status, TableCatalogBackingMigrationStatus::SnapshotMaterialized); + assert!(migration.object_backed_writes_fenced); + assert!(migration.ready_to_enable_durable_strong); + assert!( + migration + .recommended_actions + .contains(&TableCatalogBackingMigrationAction::EnableDurableStrongBacking) + ); + + let strong_store = StrongTableCatalogStore::new(backend.clone()); + let migrated_table = strong_store.load_table(bucket, "sales", "orders").await.unwrap().unwrap(); + assert_eq!(migrated_table.metadata_location, next_metadata); + assert_eq!(migrated_table.version_token, object_commit.table.version_token); + assert_eq!(migrated_table.generation, 2); + assert!( + strong_store + .load_view(bucket, "sales", "recent_orders") + .await + .unwrap() + .is_some() + ); + assert!( + strong_store + .get_commit_by_id(bucket, "table-id", "object-commit") + .await + .unwrap() + .is_some() + ); + assert!( + strong_store + .get_commit_by_idempotency_key(bucket, "table-id", "object-request") + .await + .unwrap() + .is_some() + ); + + let object_err = object_store + .commit_table(TableCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + commit_id: "stale-object-commit".to_string(), + idempotency_key: None, + operation: "append".to_string(), + expected_version_token: migrated_table.version_token.clone(), + expected_metadata_location: migrated_table.metadata_location.clone(), + new_metadata_location: strong_metadata.clone(), + requirements: Vec::new(), + writer: None, + }) + .await + .unwrap_err(); + assert_matches!(object_err, TableCatalogStoreError::Conflict(message) if message.contains("writes are fenced")); + + backend.seed_object(bucket, &strong_metadata, b"{}".to_vec()).await; + let strong_commit = strong_store + .commit_table(TableCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + commit_id: "strong-commit".to_string(), + idempotency_key: Some("strong-request".to_string()), + operation: "append".to_string(), + expected_version_token: migrated_table.version_token, + expected_metadata_location: migrated_table.metadata_location, + new_metadata_location: strong_metadata, + requirements: Vec::new(), + writer: None, + }) + .await + .unwrap(); + assert_eq!(strong_commit.table.generation, 3); + + let cancel_err = object_store + .cancel_durable_strong_backing_migration(bucket) + .await + .unwrap_err(); + assert_matches!(cancel_err, TableCatalogStoreError::Conflict(message) if message.contains("changed after materializing")); + } + + #[tokio::test] + async fn durable_strong_migration_cancel_restores_object_backed_writes_before_target_changes() { + let backend = TestCatalogObjectBackend::default(); + let object_store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + let table = IdentifierSegment::parse("orders").unwrap(); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let next_metadata = default_table_metadata_file_path(&namespace, &table, "00002.metadata.json"); + + seed_table_for_metadata_maintenance(&object_store, bucket, &namespace, &table, current_metadata.clone()).await; + object_store + .materialize_durable_strong_backing_migration(bucket) + .await + .unwrap(); + + let cancelled = object_store.cancel_durable_strong_backing_migration(bucket).await.unwrap(); + assert_eq!(cancelled.status, TableCatalogBackingMigrationCancelStatus::FenceReleased); + assert!(!cancelled.object_backed_writes_fenced); + assert!( + StrongTableCatalogStore::new(backend.clone()) + .load_table(bucket, "sales", "orders") + .await + .unwrap() + .is_none() + ); + object_store.put_table_bucket(test_bucket_entry("research")).await.unwrap(); + + backend.seed_object(bucket, &next_metadata, b"{}".to_vec()).await; + let result = object_store + .commit_table(TableCommitRequest { + table_bucket: bucket.to_string(), + namespace: namespace.public_name(), + table: table.as_str().to_string(), + commit_id: "rollback-object-commit".to_string(), + idempotency_key: None, + operation: "append".to_string(), + expected_version_token: "token-v1".to_string(), + expected_metadata_location: current_metadata, + new_metadata_location: next_metadata, + requirements: Vec::new(), + writer: None, + }) + .await + .unwrap(); + assert_eq!(result.table.generation, 2); + } + + #[tokio::test] + async fn durable_strong_migration_waits_for_in_flight_catalog_writes() { + let backend = TestCatalogObjectBackend::default(); + let object_store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + object_store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap(); + object_store.backfill_table_warehouse_index(bucket).await.unwrap(); + + let namespace_path = object_store.paths.namespace_entry_path(bucket, &namespace); + let pause = backend.pause_next_put(RUSTFS_META_BUCKET, &namespace_path).await; + let namespace_store = object_store.clone(); + let namespace_write = tokio::spawn(async move { + namespace_store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + }); + pause.wait_started().await; + + let migration_store = object_store.clone(); + let mut migration = + tokio::spawn(async move { migration_store.materialize_durable_strong_backing_migration(bucket).await }); + assert!( + tokio::time::timeout(std::time::Duration::from_millis(50), &mut migration) + .await + .is_err() + ); + + pause.release(); + namespace_write.await.unwrap().unwrap(); + let report = migration.await.unwrap().unwrap(); + assert_eq!(report.namespace_count, 1); + assert!( + StrongTableCatalogStore::new(backend) + .get_namespace(bucket, "sales") + .await + .unwrap() + .is_some() + ); + } + + #[tokio::test] + async fn durable_strong_migration_retry_recovers_after_fence_finalization_failure() { + let backend = TestCatalogObjectBackend::default(); + let object_store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + let table = IdentifierSegment::parse("orders").unwrap(); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + seed_table_for_metadata_maintenance(&object_store, bucket, &namespace, &table, current_metadata).await; + + let fence_path = object_store.paths.backing_migration_fence_path(bucket); + backend.fail_put_attempt(RUSTFS_META_BUCKET, &fence_path, 2).await; + let first = object_store + .materialize_durable_strong_backing_migration(bucket) + .await + .unwrap_err(); + assert_matches!(first, TableCatalogStoreError::Internal(message) if message.contains("injected put failure")); + assert_matches!( + object_store + .create_namespace(test_namespace_entry(bucket, &Namespace::parse("blocked").unwrap())) + .await + .unwrap_err(), + TableCatalogStoreError::Conflict(message) if message.contains("writes are fenced") + ); + + let retry = object_store + .materialize_durable_strong_backing_migration(bucket) + .await + .unwrap(); + assert_eq!(retry.status, TableCatalogBackingMigrationExecutionStatus::SnapshotAlreadyMaterialized); + assert!(retry.ready_to_enable_durable_strong); + } + + #[tokio::test] + async fn durable_strong_migration_cancel_preserves_preexisting_target_state() { + let backend = TestCatalogObjectBackend::default(); + let object_store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").unwrap(); + let table = IdentifierSegment::parse("orders").unwrap(); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + seed_table_for_metadata_maintenance(&object_store, bucket, &namespace, &table, current_metadata).await; + + let strong_store = StrongTableCatalogStore::new(backend.clone()); + strong_store + .put_table_bucket(object_store.get_table_bucket(bucket).await.unwrap().unwrap()) + .await + .unwrap(); + strong_store + .create_namespace(object_store.get_namespace(bucket, "sales").await.unwrap().unwrap()) + .await + .unwrap(); + strong_store + .create_table(object_store.load_table(bucket, "sales", "orders").await.unwrap().unwrap()) + .await + .unwrap(); + + let materialized = object_store + .materialize_durable_strong_backing_migration(bucket) + .await + .unwrap(); + assert_eq!( + materialized.status, + TableCatalogBackingMigrationExecutionStatus::SnapshotAlreadyMaterialized + ); + object_store.cancel_durable_strong_backing_migration(bucket).await.unwrap(); + assert!(strong_store.load_table(bucket, "sales", "orders").await.unwrap().is_some()); + } + + #[tokio::test] + async fn durable_strong_migration_requires_every_table_bucket_before_cutover() { + let backend = TestCatalogObjectBackend::default(); + let object_store = ObjectTableCatalogStore::new(backend); + for bucket in ["analytics", "research"] { + object_store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap(); + object_store.backfill_table_warehouse_index(bucket).await.unwrap(); + } + + let first = object_store + .materialize_durable_strong_backing_migration("analytics") + .await + .unwrap(); + assert!(!first.ready_to_enable_durable_strong); + assert_matches!( + object_store + .put_table_bucket(test_bucket_entry("new-bucket")) + .await + .unwrap_err(), + TableCatalogStoreError::Conflict(message) if message.contains("registry writes are fenced") + ); + assert!( + object_store + .plan_durable_strong_backing_migration("analytics") + .await + .unwrap() + .recommended_actions + .contains(&TableCatalogBackingMigrationAction::SnapshotRemainingTableBuckets) + ); + + let second = object_store + .materialize_durable_strong_backing_migration("research") + .await + .unwrap(); + assert!(second.ready_to_enable_durable_strong); + assert!( + object_store + .plan_durable_strong_backing_migration("analytics") + .await + .unwrap() + .ready_to_enable_durable_strong + ); + } + #[tokio::test] async fn durable_strong_migration_dry_run_reports_recovery_blockers() { let backend = TestCatalogObjectBackend::default(); diff --git a/scripts/table-catalog/failure_coverage.py b/scripts/table-catalog/failure_coverage.py index bed9915b9..aae4f1e93 100644 --- a/scripts/table-catalog/failure_coverage.py +++ b/scripts/table-catalog/failure_coverage.py @@ -335,7 +335,14 @@ def disaster_recovery_rehearsal_plan( "GET", warehouse_path(warehouse, "/catalog/migration", rest_path), "200", - "migration blockers must be empty before cutover", + "migration blockers must be empty and status must be READY_TO_SNAPSHOT before state transfer", + ), + probe_step( + "durable-backing-migration-materialize", + "POST", + warehouse_path(warehouse, "/catalog/migration", rest_path), + "200", + "source writes are fenced and the durable strong snapshot is materialized before cutover", ), ], ), @@ -504,6 +511,13 @@ def scale_fault_rehearsal_plan( "200", "migration dry-run reports inventory, replay blockers, idempotency blockers, and recommended actions", ), + probe_step( + "materialize-durable-backing-before-cutover", + "POST", + warehouse_path(warehouse, "/catalog/migration", rest_path), + "200", + "materialization drains in-flight writes and reports SNAPSHOT_MATERIALIZED before backing restart", + ), ], ), rehearsal_phase( diff --git a/scripts/table-catalog/pyiceberg_smoke.py b/scripts/table-catalog/pyiceberg_smoke.py index b1f786231..d743747cc 100755 --- a/scripts/table-catalog/pyiceberg_smoke.py +++ b/scripts/table-catalog/pyiceberg_smoke.py @@ -239,8 +239,8 @@ UNSUPPORTED_INVENTORY: list[dict[str, str]] = [ PRODUCTION_READINESS_INVENTORY: list[dict[str, str]] = [ { "capability": "strong-catalog-backing", - "status": "migration-contract-supported", - "validation": "catalog export and diagnostics expose an object-backed manifest, recoverable commit-log WAL state, strong-kv-wal migration target, required migration steps, and recovery blockers before cutover", + "status": "state-transfer-supported", + "validation": "preflight exposes recovery blockers and target agreement; migration execution fences object-backed writes, drains in-flight catalog mutations, and idempotently materializes table bucket, namespace, table, view, commit-log, and idempotency state before cutover", }, { "capability": "single-active-writer-ha", diff --git a/scripts/table-catalog/test_failure_coverage.py b/scripts/table-catalog/test_failure_coverage.py index 8e015939c..164e8f492 100644 --- a/scripts/table-catalog/test_failure_coverage.py +++ b/scripts/table-catalog/test_failure_coverage.py @@ -132,7 +132,11 @@ class FailureCoverageTest(unittest.TestCase): migration_steps = {step["name"]: step for step in phases["migration-preflight"]["steps"]} self.assertEqual(migration_steps["durable-backing-migration-dry-run"]["path"], "/iceberg/v1/lake/catalog/migration") - self.assertEqual(migration_steps["durable-backing-migration-dry-run"]["expected_behavior"], "migration blockers must be empty before cutover") + self.assertEqual(migration_steps["durable-backing-migration-materialize"]["method"], "POST") + self.assertIn( + "READY_TO_SNAPSHOT", + migration_steps["durable-backing-migration-dry-run"]["expected_behavior"], + ) validation_steps = {step["name"]: step for step in phases["post-recovery-validation"]["steps"]} self.assertEqual(validation_steps["load-table-after-recovery"]["path"], "/iceberg/v1/lake/namespaces/sales/tables/orders") @@ -202,6 +206,7 @@ class FailureCoverageTest(unittest.TestCase): cutover_steps = {step["name"]: step for step in phases["durable-backing-cutover-under-load"]["steps"]} self.assertEqual(cutover_steps["migration-dry-run-before-cutover"]["path"], "/iceberg/v1/lake/catalog/migration") + self.assertEqual(cutover_steps["materialize-durable-backing-before-cutover"]["method"], "POST") self.assertIn("catalog-backup", cutover_steps["capture-cutover-backup"]["expected_behavior"]) recovery_steps = {step["name"]: step for step in phases["recovery-rollback-import-under-load"]["steps"]} diff --git a/scripts/table-catalog/test_pyiceberg_smoke.py b/scripts/table-catalog/test_pyiceberg_smoke.py index 48432504e..573fa3aaf 100644 --- a/scripts/table-catalog/test_pyiceberg_smoke.py +++ b/scripts/table-catalog/test_pyiceberg_smoke.py @@ -864,7 +864,7 @@ class PyIcebergSmokeConfigTest(unittest.TestCase): self.assertIn("single-active-writer-ha", capabilities) self.assertIn("scale-validation-matrix", capabilities) strong_backing = next(entry for entry in inventory if entry["capability"] == "strong-catalog-backing") - self.assertEqual(strong_backing["status"], "migration-contract-supported") + self.assertEqual(strong_backing["status"], "state-transfer-supported") for entry in inventory: self.assertIn("status", entry) self.assertIn("validation", entry)