From e4031940e266a1ce878bc4d371650135729669da Mon Sep 17 00:00:00 2001 From: GatewayJ <835269233@qq.com> Date: Mon, 14 Sep 2026 09:58:52 +0800 Subject: [PATCH] fix(tables): bound warehouse index recovery (#7677) Co-authored-by: Chris --- docs/architecture/s3-tables-support-matrix.md | 1 + .../admin/handlers/table_catalog/config.rs | 20 ++ .../src/admin/handlers/table_catalog/mod.rs | 1 + .../admin/handlers/table_catalog/routes.rs | 5 + .../src/admin/handlers/table_catalog/tests.rs | 2 + rustfs/src/admin/route_policy.rs | 24 ++- rustfs/src/admin/route_registration_test.rs | 10 + rustfs/src/table_catalog/mod.rs | 1 + rustfs/src/table_catalog/store/mod.rs | 2 + rustfs/src/table_catalog/store/object.rs | 191 ++++++++++++++++-- rustfs/src/table_catalog/tests.rs | 144 ++++++++++++- 11 files changed, 378 insertions(+), 23 deletions(-) diff --git a/docs/architecture/s3-tables-support-matrix.md b/docs/architecture/s3-tables-support-matrix.md index 7d476d5f9..f2bee2d7d 100644 --- a/docs/architecture/s3-tables-support-matrix.md +++ b/docs/architecture/s3-tables-support-matrix.md @@ -25,6 +25,7 @@ RustFS S3 Tables is an Iceberg REST Catalog and table-bucket implementation on t | `/iceberg/v1` | Supported | Canonical REST Catalog prefix; default REST signing name `s3`. | | `/_iceberg/v1` | Supported compatibility alias | MinIO AIStor-style alias; smoke profile defaults to signing name `s3tables`. | | S3 object data plane | Supported | Data, metadata, manifest, and delete files are ordinary S3 objects with table-aware policy checks on warehouse paths. | +| Object-backed warehouse-index recovery | Supported with bounded request fallback | Data-plane requests repair or verify a missing warehouse-prefix index against bounded pages totaling at most 4,096 catalog objects. Larger catalogs fail closed with a retryable service error and require `POST /iceberg/v1/{warehouse}/catalog/warehouse-index/backfill` with `admin:MigrateTableCatalog` instead of allowing one S3 request to trigger an unbounded scan. | | Table bucket enablement | Supported | A regular bucket is enabled for catalog use and addressed as the REST catalog warehouse. | | Catalog-vended table credentials | Automated when enabled | Disabled by default. LoadTable vends credentials only when `X-Iceberg-Access-Delegation` contains the exact `vended-credentials` token; the dedicated credentials endpoint uses the same issuer path. | | AWS S3 Tables endpoint shape | Profile generator | Generates the AWS catalog URI and warehouse ARN shape for migration docs. API parity not claimed. | diff --git a/rustfs/src/admin/handlers/table_catalog/config.rs b/rustfs/src/admin/handlers/table_catalog/config.rs index 46959a717..610beb037 100644 --- a/rustfs/src/admin/handlers/table_catalog/config.rs +++ b/rustfs/src/admin/handlers/table_catalog/config.rs @@ -119,6 +119,26 @@ impl Operation for CancelTableCatalogMigrationHandler { } } +pub struct BackfillTableWarehouseIndexHandler {} + +#[async_trait::async_trait] +impl Operation for BackfillTableWarehouseIndexHandler { + 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_from_extensions(&req.extensions, &warehouse).await?; + let store = table_catalog_object_store_from_extensions(&req.extensions)?; + let started = Instant::now(); + let result = store + .backfill_table_warehouse_index(&warehouse) + .await + .map_err(catalog_store_error); + record_table_catalog_admin_operation_result("warehouse-index-backfill", &warehouse, "", "", started, &result); + result?; + Ok(empty_response(StatusCode::NO_CONTENT)) + } +} + pub struct ExternalCatalogBridgeHandler {} #[async_trait::async_trait] diff --git a/rustfs/src/admin/handlers/table_catalog/mod.rs b/rustfs/src/admin/handlers/table_catalog/mod.rs index 62c757eea..1c4127a63 100644 --- a/rustfs/src/admin/handlers/table_catalog/mod.rs +++ b/rustfs/src/admin/handlers/table_catalog/mod.rs @@ -186,6 +186,7 @@ static GET_TABLE_CATALOG_MIGRATION_HANDLER: GetTableCatalogMigrationHandler = Ge static MATERIALIZE_TABLE_CATALOG_MIGRATION_HANDLER: MaterializeTableCatalogMigrationHandler = MaterializeTableCatalogMigrationHandler {}; static CANCEL_TABLE_CATALOG_MIGRATION_HANDLER: CancelTableCatalogMigrationHandler = CancelTableCatalogMigrationHandler {}; +static BACKFILL_TABLE_WAREHOUSE_INDEX_HANDLER: BackfillTableWarehouseIndexHandler = BackfillTableWarehouseIndexHandler {}; static LIST_NAMESPACES_HANDLER: RestListNamespacesHandler = RestListNamespacesHandler {}; static CREATE_NAMESPACE_HANDLER: RestCreateNamespaceHandler = RestCreateNamespaceHandler {}; static GET_NAMESPACE_HANDLER: RestGetNamespaceHandler = RestGetNamespaceHandler {}; diff --git a/rustfs/src/admin/handlers/table_catalog/routes.rs b/rustfs/src/admin/handlers/table_catalog/routes.rs index fcfe781f1..2e8177f7d 100644 --- a/rustfs/src/admin/handlers/table_catalog/routes.rs +++ b/rustfs/src/admin/handlers/table_catalog/routes.rs @@ -49,6 +49,11 @@ fn register_table_catalog_prefix_routes(r: &mut S3Router, prefix format!("{prefix}/{{warehouse}}/catalog/migration").as_str(), AdminOperation(&CANCEL_TABLE_CATALOG_MIGRATION_HANDLER), )?; + r.insert( + Method::POST, + format!("{prefix}/{{warehouse}}/catalog/warehouse-index/backfill").as_str(), + AdminOperation(&BACKFILL_TABLE_WAREHOUSE_INDEX_HANDLER), + )?; r.insert( Method::GET, format!("{prefix}/{{warehouse}}/namespaces").as_str(), diff --git a/rustfs/src/admin/handlers/table_catalog/tests.rs b/rustfs/src/admin/handlers/table_catalog/tests.rs index 240523434..4792b44fe 100644 --- a/rustfs/src/admin/handlers/table_catalog/tests.rs +++ b/rustfs/src/admin/handlers/table_catalog/tests.rs @@ -696,6 +696,7 @@ fn table_catalog_handlers_require_table_admin_actions() { for handler in [ "MaterializeTableCatalogMigrationHandler", "CancelTableCatalogMigrationHandler", + "BackfillTableWarehouseIndexHandler", ] { let block = operation_block(&src, handler); assert!( @@ -848,6 +849,7 @@ fn table_catalog_handlers_require_enabled_table_bucket_marker_before_catalog_sta "GetTableCatalogMigrationHandler", "MaterializeTableCatalogMigrationHandler", "CancelTableCatalogMigrationHandler", + "BackfillTableWarehouseIndexHandler", "RestListNamespacesHandler", "RestCreateNamespaceHandler", "RestGetNamespaceHandler", diff --git a/rustfs/src/admin/route_policy.rs b/rustfs/src/admin/route_policy.rs index f53366d6d..18b70d757 100644 --- a/rustfs/src/admin/route_policy.rs +++ b/rustfs/src/admin/route_policy.rs @@ -1013,6 +1013,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ MIGRATE_TABLE_CATALOG, RouteRiskLevel::High, ), + admin( + HttpMethod::Post, + "/iceberg/v1/{warehouse}/catalog/warehouse-index/backfill", + MIGRATE_TABLE_CATALOG, + RouteRiskLevel::High, + ), admin( HttpMethod::Get, "/iceberg/v1/{warehouse}/namespaces", @@ -1297,6 +1303,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ MIGRATE_TABLE_CATALOG, RouteRiskLevel::High, ), + admin( + HttpMethod::Post, + "/_iceberg/v1/{warehouse}/catalog/warehouse-index/backfill", + MIGRATE_TABLE_CATALOG, + RouteRiskLevel::High, + ), admin( HttpMethod::Get, "/_iceberg/v1/{warehouse}/namespaces", @@ -1845,7 +1857,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(), 98); + assert_eq!(table_specs.count(), 100); 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); @@ -2028,6 +2040,16 @@ mod tests { 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}/catalog/warehouse-index/backfill", + MIGRATE_TABLE_CATALOG, + ); + assert_action( + HttpMethod::Post, + "/_iceberg/v1/{warehouse}/catalog/warehouse-index/backfill", + 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 260f098b1..a0a178a9a 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -439,6 +439,11 @@ fn expected_admin_route_matrix() -> Vec { 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::POST, + "/{warehouse}/catalog/warehouse-index/backfill", + "/analytics/catalog/warehouse-index/backfill", + ), 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"), @@ -636,6 +641,11 @@ fn expected_admin_route_matrix() -> Vec { 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::POST, + "/{warehouse}/catalog/warehouse-index/backfill", + "/analytics/catalog/warehouse-index/backfill", + ), 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/mod.rs b/rustfs/src/table_catalog/mod.rs index 4772511cb..b023dad9a 100644 --- a/rustfs/src/table_catalog/mod.rs +++ b/rustfs/src/table_catalog/mod.rs @@ -169,6 +169,7 @@ const WAREHOUSE_INDEX_ROOT: &str = "warehouse-index"; const WAREHOUSE_INDEX_STATE_FILE: &str = "state.json"; const TABLE_RENAME_ROOT: &str = "renames"; const WAREHOUSE_INDEX_MAX_PREFIX_DEPTH: usize = 64; +const TABLE_DATA_PLANE_INDEX_MISS_SCAN_MAX_CATALOG_OBJECTS: usize = 4096; const EXTERNAL_CATALOG_ROOT: &str = "external-catalog"; const EXTERNAL_CATALOG_BRIDGE_FILE: &str = "bridge.json"; const MAINTENANCE_ROOT: &str = "maintenance"; diff --git a/rustfs/src/table_catalog/store/mod.rs b/rustfs/src/table_catalog/store/mod.rs index b79c29322..dd39449c6 100644 --- a/rustfs/src/table_catalog/store/mod.rs +++ b/rustfs/src/table_catalog/store/mod.rs @@ -21,6 +21,8 @@ mod strong; use migration::table_catalog_backing_manifest; pub(crate) use object::ObjectTableCatalogStore; #[cfg(test)] +pub(super) use object::bounded_table_entry_objects_for_data_plane_scan; +#[cfg(test)] pub(super) use strong::{ STRONG_TABLE_CATALOG_RELOAD_MAX_ATTEMPTS, STRONG_TABLE_CATALOG_SNAPSHOT_MAX_SIZE, StrongCommitSnapshotRecord, StrongTableCatalogBucketSnapshot, StrongTableCatalogSnapshot, strong_snapshot_write_version, diff --git a/rustfs/src/table_catalog/store/object.rs b/rustfs/src/table_catalog/store/object.rs index f80db59d4..a0501ab5a 100644 --- a/rustfs/src/table_catalog/store/object.rs +++ b/rustfs/src/table_catalog/store/object.rs @@ -245,6 +245,27 @@ enum TableWarehouseIndexResolution { Missing, } +enum TableWarehouseIndexBackfillMode { + InitializeIfNeeded, + ScanAllTables, +} + +pub(in crate::table_catalog) fn bounded_table_entry_objects_for_data_plane_scan( + objects: Vec, + is_truncated: bool, + max_catalog_objects: usize, +) -> TableCatalogStoreResult> { + if is_truncated || objects.len() > max_catalog_objects { + return Err(TableCatalogStoreError::Unavailable(format!( + "table data-plane warehouse-index miss scan exceeds the {max_catalog_objects}-catalog-object safety limit" + ))); + } + Ok(objects + .into_iter() + .filter(|object| object.ends_with(TABLE_ENTRY_FILE)) + .collect()) +} + #[derive(Clone)] pub(crate) struct ObjectTableCatalogStore { pub(in crate::table_catalog) backend: B, @@ -1717,12 +1738,13 @@ where .await? { TableWarehouseIndexResolution::Active(resource) => Ok(Some(resource)), - TableWarehouseIndexResolution::Deleted(tombstone) => { - Ok(scan_table_data_plane_resource_for_object(self, table_bucket, object) - .await? - .or(Some(tombstone))) + TableWarehouseIndexResolution::Deleted(tombstone) => Ok(self + .scan_table_data_plane_resource_for_index_miss(table_bucket, object) + .await? + .or(Some(tombstone))), + TableWarehouseIndexResolution::Missing => { + self.scan_table_data_plane_resource_for_index_miss(table_bucket, object).await } - TableWarehouseIndexResolution::Missing => scan_table_data_plane_resource_for_object(self, table_bucket, object).await, } } @@ -1771,11 +1793,61 @@ where Ok(matched) } + async fn scan_table_data_plane_resource_for_index_miss( + &self, + table_bucket: &str, + object: &str, + ) -> TableCatalogStoreResult> { + let mut matched: Option = None; + for table in self + .list_all_table_entries_with_limit(table_bucket, Some(TABLE_DATA_PLANE_INDEX_MISS_SCAN_MAX_CATALOG_OBJECTS)) + .await? + { + if table.state != TableCatalogEntryState::Active { + continue; + } + let warehouse_object_prefix = table_warehouse_object_prefix(&table)?; + if !object.starts_with(&warehouse_object_prefix) { + continue; + } + if let Some(current) = matched.as_ref() { + return Err(TableCatalogStoreError::Invalid(format!( + "object {object} matches overlapping active table warehouse prefixes {} and {warehouse_object_prefix}", + current.warehouse_object_prefix + ))); + } + matched = Some(table_data_plane_resource_from_entry(table, warehouse_object_prefix)); + } + if let Some(resource) = matched.as_ref() { + let _migration_guard = self.acquire_object_backed_catalog_write_permit(table_bucket).await?; + self.backfill_active_table_warehouse_index_with_prefix_check( + &resource.table_bucket, + &resource.namespace, + &resource.table, + true, + ) + .await?; + } + Ok(matched) + } + + #[cfg(test)] pub(in crate::table_catalog) async fn backfill_active_table_warehouse_index( &self, table_bucket: &str, namespace: &str, table: &str, + ) -> TableCatalogStoreResult<()> { + self.backfill_active_table_warehouse_index_with_prefix_check(table_bucket, namespace, table, false) + .await + } + + async fn backfill_active_table_warehouse_index_with_prefix_check( + &self, + table_bucket: &str, + namespace: &str, + table: &str, + prefix_already_checked: bool, ) -> TableCatalogStoreResult<()> { let namespace = parse_namespace_for_store(namespace)?; let table = parse_table_for_store(table)?; @@ -1787,21 +1859,40 @@ where if current.state != TableCatalogEntryState::Active { return Ok(()); } - self.reserve_table_warehouse_index(¤t, false).await.map(|_| ()) + self.reserve_table_warehouse_index(¤t, prefix_already_checked) + .await + .map(|_| ()) } - pub(in crate::table_catalog) async fn backfill_table_warehouse_index( + pub(crate) async fn backfill_table_warehouse_index(&self, table_bucket: &str) -> TableCatalogStoreResult<()> { + self.backfill_table_warehouse_index_with_limit(table_bucket, None, TableWarehouseIndexBackfillMode::ScanAllTables) + .await + } + + async fn backfill_table_warehouse_index_for_data_plane(&self, table_bucket: &str) -> TableCatalogStoreResult<()> { + self.backfill_table_warehouse_index_with_limit( + table_bucket, + Some(TABLE_DATA_PLANE_INDEX_MISS_SCAN_MAX_CATALOG_OBJECTS), + TableWarehouseIndexBackfillMode::InitializeIfNeeded, + ) + .await + } + + async fn backfill_table_warehouse_index_with_limit( &self, table_bucket: &str, + max_catalog_objects: Option, + mode: TableWarehouseIndexBackfillMode, ) -> TableCatalogStoreResult<()> { let _migration_guard = self.acquire_object_backed_catalog_write_permit(table_bucket).await?; let state_object = self.paths.warehouse_index_state_path(table_bucket); let _guard = self.backend.acquire_write_lock(self.catalog_bucket(), &state_object).await?; - if self.read_warehouse_index_state_unlocked(table_bucket).await? { + let index_ready = self.read_warehouse_index_state_unlocked(table_bucket).await?; + if matches!(mode, TableWarehouseIndexBackfillMode::InitializeIfNeeded) && index_ready { return Ok(()); } let tables = self - .list_all_table_entries(table_bucket) + .list_all_table_entries_with_limit(table_bucket, max_catalog_objects) .await? .into_iter() .filter(|table| table.state == TableCatalogEntryState::Active) @@ -1828,8 +1919,13 @@ where ))); } for table in tables { - self.backfill_active_table_warehouse_index(&table.table_bucket, &table.namespace, &table.table) - .await?; + self.backfill_active_table_warehouse_index_with_prefix_check( + &table.table_bucket, + &table.namespace, + &table.table, + true, + ) + .await?; } self.write_warehouse_index_state_unlocked(table_bucket).await } @@ -1876,15 +1972,70 @@ where } async fn list_all_table_entries(&self, table_bucket: &str) -> TableCatalogStoreResult> { - let mut entries = Vec::new(); - for object in self - .backend - .list_objects(self.catalog_bucket(), &self.paths.namespace_entries_prefix(table_bucket)) - .await? - { - if !object.ends_with(TABLE_ENTRY_FILE) { - continue; + self.list_all_table_entries_with_limit(table_bucket, None).await + } + + async fn list_table_entry_objects_for_data_plane_scan( + &self, + table_bucket: &str, + max_catalog_objects: usize, + ) -> TableCatalogStoreResult> { + let mut objects = Vec::new(); + let mut cursor = None; + loop { + let remaining = max_catalog_objects.saturating_sub(objects.len()); + let page_size = remaining.min(TABLE_CATALOG_LIST_MAX_KEYS); + let limit = NonZeroUsize::new(page_size).ok_or_else(|| { + TableCatalogStoreError::Unavailable(format!( + "table data-plane warehouse-index miss scan exceeds the {max_catalog_objects}-catalog-object safety limit" + )) + })?; + let page = self + .backend + .list_objects_page( + self.catalog_bucket(), + &self.paths.namespace_entries_prefix(table_bucket), + cursor.as_deref(), + limit, + ) + .await?; + let next_cursor = page.objects.last().cloned(); + objects.extend(page.objects); + if !page.is_truncated { + return bounded_table_entry_objects_for_data_plane_scan(objects, false, max_catalog_objects); } + if objects.len() >= max_catalog_objects { + return bounded_table_entry_objects_for_data_plane_scan(objects, true, max_catalog_objects); + } + if next_cursor.is_none() || next_cursor == cursor { + return Err(TableCatalogStoreError::Internal( + "catalog object pagination did not advance during table data-plane index-miss scan".to_string(), + )); + } + cursor = next_cursor; + } + } + + async fn list_all_table_entries_with_limit( + &self, + table_bucket: &str, + max_catalog_objects: Option, + ) -> TableCatalogStoreResult> { + let mut entries = Vec::new(); + let table_objects = match max_catalog_objects { + Some(max_catalog_objects) => { + self.list_table_entry_objects_for_data_plane_scan(table_bucket, max_catalog_objects) + .await? + } + None => self + .backend + .list_objects(self.catalog_bucket(), &self.paths.namespace_entries_prefix(table_bucket)) + .await? + .into_iter() + .filter(|object| object.ends_with(TABLE_ENTRY_FILE)) + .collect(), + }; + for object in table_objects { let Some((entry, _)) = self.read_entry::(self.catalog_bucket(), &object).await? else { continue; }; @@ -5269,7 +5420,7 @@ where let resource = if self.warehouse_index_ready(table_bucket).await? { self.resolve_table_data_plane_resource_with_scan(table_bucket, object).await } else { - match self.backfill_table_warehouse_index(table_bucket).await { + match self.backfill_table_warehouse_index_for_data_plane(table_bucket).await { Ok(()) => self.resolve_table_data_plane_resource_with_scan(table_bucket, object).await, Err(err @ TableCatalogStoreError::Internal(_)) => { tracing::warn!( diff --git a/rustfs/src/table_catalog/tests.rs b/rustfs/src/table_catalog/tests.rs index f5767c15d..0d7c082a3 100644 --- a/rustfs/src/table_catalog/tests.rs +++ b/rustfs/src/table_catalog/tests.rs @@ -20,6 +20,31 @@ use std::sync::Arc; const TABLE_CATALOG_TEST_TIMEOUT: StdDuration = StdDuration::from_secs(30); +#[test] +fn table_data_plane_index_miss_scan_rejects_catalogs_above_its_object_limit() { + let objects = vec![ + "catalog/table-1/table-entry.json".to_string(), + "catalog/namespace.json".to_string(), + "catalog/table-2/table-entry.json".to_string(), + ]; + + assert_eq!( + bounded_table_entry_objects_for_data_plane_scan(objects.clone(), false, 3).unwrap(), + vec![ + "catalog/table-1/table-entry.json".to_string(), + "catalog/table-2/table-entry.json".to_string(), + ] + ); + assert_matches!( + bounded_table_entry_objects_for_data_plane_scan(objects, false, 2), + Err(TableCatalogStoreError::Unavailable(message)) if message.contains("2-catalog-object safety limit") + ); + assert_matches!( + bounded_table_entry_objects_for_data_plane_scan(vec!["catalog/table-entry.json".to_string()], true, 1), + Err(TableCatalogStoreError::Unavailable(message)) if message.contains("1-catalog-object safety limit") + ); +} + #[test] fn catalog_lock_authority_failures_are_typed_as_unavailable() { for error in [ @@ -4481,6 +4506,8 @@ async fn table_data_plane_resource_scans_when_a_ready_index_entry_is_missing() { .delete_object(RUSTFS_META_BUCKET, &store.paths.warehouse_index_entry_path(bucket, "tables/table-id/")) .await .expect("warehouse index entry should be removed"); + let migration_lock = store.paths.backing_migration_fence_lock_path(bucket); + let migration_permits_before = backend.read_lock_acquisition_count(RUSTFS_META_BUCKET, &migration_lock).await; backend.reset_call_counts().await; let resource = table_data_plane_resource_for_object(&store, bucket, object) @@ -4489,7 +4516,120 @@ async fn table_data_plane_resource_scans_when_a_ready_index_entry_is_missing() { .expect("the table scan must retain table-aware protection"); assert_eq!(resource.table, "orders"); - assert!(backend.list_call_count().await > 0); + assert_eq!(backend.list_call_count().await, 1); + assert_eq!( + backend.read_lock_acquisition_count(RUSTFS_META_BUCKET, &migration_lock).await, + migration_permits_before + 1, + "repairing a missing warehouse index must hold the object-backed migration permit" + ); + assert!( + store + .read_entry::( + RUSTFS_META_BUCKET, + &store.paths.warehouse_index_entry_path(bucket, "tables/table-id/"), + ) + .await + .expect("repaired warehouse index lookup should succeed") + .is_some(), + "the bounded scan should repair the missing index" + ); + + backend.reset_call_counts().await; + let indexed = table_data_plane_resource_for_object(&store, bucket, object) + .await + .expect("repaired warehouse index lookup should succeed") + .expect("repaired warehouse index should retain table-aware protection"); + assert_eq!(indexed.table, "orders"); + assert_eq!(backend.list_call_count().await, 0); + assert_eq!( + backend.read_lock_acquisition_count(RUSTFS_META_BUCKET, &migration_lock).await, + migration_permits_before + 1, + "an indexed lookup must not acquire a write permit" + ); +} + +#[tokio::test] +async fn operator_backfill_reconciles_a_missing_ready_index_above_the_data_plane_limit() { + let backend = TestCatalogObjectBackend::default(); + let store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").expect("namespace should parse"); + let table = IdentifierSegment::parse("orders").expect("table should parse"); + let current = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let object = "tables/table-id/data/part-00001.parquet"; + let index_path = store.paths.warehouse_index_entry_path(bucket, "tables/table-id/"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current).await; + backend + .delete_object(RUSTFS_META_BUCKET, &index_path) + .await + .expect("warehouse index entry should be removed"); + let catalog_prefix = store.paths.namespace_entries_prefix(bucket); + for index in 0..TABLE_DATA_PLANE_INDEX_MISS_SCAN_MAX_CATALOG_OBJECTS { + backend + .seed_object(RUSTFS_META_BUCKET, &format!("{catalog_prefix}noise-{index}.json"), b"{}".to_vec()) + .await; + } + + assert_matches!( + table_data_plane_resource_for_object(&store, bucket, object).await, + Err(TableCatalogStoreError::Unavailable(message)) if message.contains("4096-catalog-object safety limit") + ); + store + .backfill_table_warehouse_index(bucket) + .await + .expect("operator backfill should reconcile a ready index without the data-plane limit"); + assert!( + store + .read_entry::(RUSTFS_META_BUCKET, &index_path) + .await + .expect("reconciled warehouse index lookup should succeed") + .is_some(), + "operator reconciliation must recreate the missing ready-state index" + ); + + backend.reset_call_counts().await; + let resource = table_data_plane_resource_for_object(&store, bucket, object) + .await + .expect("reconciled index lookup should succeed") + .expect("reconciled index should retain table-aware protection"); + assert_eq!(resource.table, "orders"); + assert_eq!(backend.list_call_count().await, 0); +} + +#[tokio::test] +async fn table_data_plane_index_repair_honors_an_active_migration_fence() { + let backend = TestCatalogObjectBackend::default(); + let store = ObjectTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").expect("namespace should parse"); + let table = IdentifierSegment::parse("orders").expect("table should parse"); + let current = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let object = "tables/table-id/data/part-00001.parquet"; + let index_path = store.paths.warehouse_index_entry_path(bucket, "tables/table-id/"); + + seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current).await; + store + .materialize_durable_strong_backing_migration(bucket) + .await + .expect("durable strong migration should install its object-backed write fence"); + backend + .delete_object(RUSTFS_META_BUCKET, &index_path) + .await + .expect("warehouse index entry should be removed after migration fencing"); + + assert_matches!( + table_data_plane_resource_for_object(&store, bucket, object).await, + Err(TableCatalogStoreError::Conflict(message)) if message.contains("writes are fenced") + ); + assert!( + store + .read_entry::(RUSTFS_META_BUCKET, &index_path) + .await + .expect("warehouse index lookup should succeed") + .is_none(), + "a fenced object-backed catalog must not repair the missing warehouse index" + ); } #[tokio::test] @@ -5161,7 +5301,7 @@ async fn table_data_plane_resource_falls_back_to_scan_without_index_state() { .expect("legacy table entry should resolve"); assert_eq!(resource.table, "orders"); - assert!(backend.list_call_count().await > 0); + assert_eq!(backend.list_call_count().await, 1); assert!(store.warehouse_index_ready(bucket).await.unwrap()); backend.reset_call_counts().await;