mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-03 20:20:29 +00:00
fix(tables): bound warehouse index recovery (#7677)
Co-authored-by: Chris <anzhengchao@gmail.com>
This commit is contained in:
@@ -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. |
|
||||
|
||||
@@ -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<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||
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]
|
||||
|
||||
@@ -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 {};
|
||||
|
||||
@@ -49,6 +49,11 @@ fn register_table_catalog_prefix_routes(r: &mut S3Router<AdminOperation>, 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(),
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -439,6 +439,11 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
|
||||
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<RouteMatrixEntry> {
|
||||
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"),
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<String>,
|
||||
is_truncated: bool,
|
||||
max_catalog_objects: usize,
|
||||
) -> TableCatalogStoreResult<Vec<String>> {
|
||||
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<B> {
|
||||
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<Option<TableDataPlaneResource>> {
|
||||
let mut matched: Option<TableDataPlaneResource> = 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<usize>,
|
||||
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<Vec<TableEntry>> {
|
||||
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<Vec<String>> {
|
||||
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<usize>,
|
||||
) -> TableCatalogStoreResult<Vec<TableEntry>> {
|
||||
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::<TableEntry>(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!(
|
||||
|
||||
@@ -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::<TableWarehouseIndexEntry>(
|
||||
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::<TableWarehouseIndexEntry>(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::<TableWarehouseIndexEntry>(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;
|
||||
|
||||
Reference in New Issue
Block a user