feat(table-catalog): expose backing migration contract (#3574)

* feat(table-catalog): expose backing migration contract

* fix(table-catalog): require register auth for external sync

---------

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
Henry Guo
2026-06-18 21:14:40 +08:00
committed by GitHub
parent 7b0cb9e725
commit 5d9fee5c0c
5 changed files with 437 additions and 1 deletions
+19 -1
View File
@@ -5193,9 +5193,17 @@ impl Operation for SyncExternalCatalogBridgeHandler {
let table = table_name_from_params(&params)?;
let resource = TableCatalogResource::table(&warehouse, &namespace, &table);
authorize_table_catalog_resource_request(&req, &resource, AdminAction::SetTableMetadataLocationAction).await?;
let request = read_json_body::<ExternalCatalogBridgeSyncRequest>(req.input).await?;
let metadata_backend = table_catalog_backend()?;
let store = crate::table_catalog::ObjectTableCatalogStore::new(metadata_backend.clone());
if store
.load_table(&warehouse, &namespace.public_name(), &table)
.await
.map_err(catalog_store_error)?
.is_none()
{
authorize_table_catalog_resource_request(&req, &resource, AdminAction::RegisterTableAction).await?;
}
let request = read_json_body::<ExternalCatalogBridgeSyncRequest>(req.input).await?;
let table_bucket_enabled = table_bucket_enabled_from_metadata(&warehouse).await?;
let response = sync_external_catalog_bridge_response(
&store,
@@ -5474,6 +5482,16 @@ mod tests {
);
}
let sync_bridge_block = operation_block(src, "SyncExternalCatalogBridgeHandler");
assert!(
sync_bridge_block.contains("AdminAction::RegisterTableAction"),
"external catalog sync should require register authorization before creating a missing table"
);
assert!(
sync_bridge_block.contains(".load_table(&warehouse, &namespace.public_name(), &table)"),
"external catalog sync should branch authorization on current table existence"
);
for (handler, action) in [
("RestLoadTableHandler", "AdminAction::GetTableMetadataAction"),
("RestTableExistsHandler", "AdminAction::GetTableAction"),
+370
View File
@@ -72,6 +72,7 @@ pub(crate) const TABLE_METADATA_POINTER_VERSION: u16 = 1;
pub(crate) const TABLE_CATALOG_ENTRY_VERSION: u16 = 1;
pub(crate) const TABLE_MAINTENANCE_CONFIG_VERSION: u16 = 1;
pub(crate) const TABLE_EXTERNAL_CATALOG_BRIDGE_VERSION: u16 = 1;
pub(crate) const TABLE_CATALOG_BACKING_MANIFEST_VERSION: u16 = 1;
pub(crate) const TABLE_METADATA_FILE_NAME_MAX_LEN: usize = 128;
pub const TABLE_RESERVED_PREFIX: &str = BUCKET_TABLE_RESERVED_PREFIX;
const WAREHOUSE_ROOT: &str = "warehouses";
@@ -843,6 +844,168 @@ pub(crate) struct TableCatalogExport {
pub table_bucket: TableBucketEntry,
pub namespace: NamespaceEntry,
pub table: TableEntry,
pub backing_manifest: TableCatalogBackingManifest,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub(crate) struct TableCatalogBackingManifest {
pub version: u16,
pub current: TableCatalogBackingProfile,
pub migration: TableCatalogBackingMigrationPlan,
pub ha: TableCatalogHaPolicy,
pub scale_validation: TableCatalogScaleValidation,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub(crate) struct TableCatalogBackingProfile {
pub kind: TableCatalogBackingKind,
pub authority: TableCatalogAuthority,
pub consistency: TableCatalogConsistencyMode,
pub durability: TableCatalogDurabilityMode,
pub current_pointer_path: String,
pub wal: TableCatalogWalState,
pub snapshot: TableCatalogSnapshotState,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogBackingKind {
ObjectBacked,
StrongKvWal,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogAuthority {
RustfsSysObject,
LinearizableMetadataKv,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogConsistencyMode {
ConditionalObjectCas,
LinearizableCas,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogDurabilityMode {
StagedCommitLogBeforePointerUpdate,
WalBeforeStateMachineApply,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub(crate) struct TableCatalogWalState {
pub status: TableCatalogWalStatus,
pub commit_log_prefix: String,
pub idempotency_index_prefix: String,
pub committed_generation: u64,
pub staged_before_table_update_count: usize,
pub finalization_required_count: usize,
pub idempotency_repair_required_count: usize,
pub manual_review_count: usize,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogWalStatus {
Recoverable,
RecoveryRequired,
ManualReviewRequired,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub(crate) struct TableCatalogSnapshotState {
pub export_api: String,
pub includes_table_bucket: bool,
pub includes_namespace: bool,
pub includes_table_pointer: bool,
pub includes_backing_manifest: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub(crate) struct TableCatalogBackingMigrationPlan {
pub source_kind: TableCatalogBackingKind,
pub target_kind: TableCatalogBackingKind,
pub status: TableCatalogBackingMigrationStatus,
pub required_steps: Vec<TableCatalogBackingMigrationStep>,
pub blockers: Vec<TableCatalogBackingMigrationBlocker>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogBackingMigrationStatus {
ReadyToSnapshot,
RecoveryRequired,
ManualReviewRequired,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogBackingMigrationStep {
SnapshotCatalogExport,
ReplayCommitLog,
VerifyCurrentPointer,
EnableSingleWriterFencing,
CutOverLinearizableReads,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogBackingMigrationBlocker {
CommitRecoveryRequired,
CommitManualReviewRequired,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub(crate) struct TableCatalogHaPolicy {
pub writer_region_model: TableCatalogHaWriterModel,
pub read_replica_strategy: TableCatalogReadReplicaStrategy,
pub commit_read_requirement: TableCatalogCommitReadRequirement,
pub active_active_supported: bool,
pub failover_requires_operator_promotion: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogHaWriterModel {
SingleActiveWriterRegion,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogReadReplicaStrategy {
ReadOnlyReplicasForListAndLoad,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogCommitReadRequirement {
LinearizableLeaderRead,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub(crate) struct TableCatalogScaleValidation {
pub status: TableCatalogScaleValidationStatus,
pub benchmark_required: bool,
pub required_scenarios: Vec<TableCatalogScaleValidationScenario>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogScaleValidationStatus {
MatrixPublished,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub(crate) enum TableCatalogScaleValidationScenario {
ConcurrentCommitCas,
CommitLogRecoveryReplay,
MigrationSnapshotReplay,
ReadReplicaStaleReadGuard,
ClientConformanceMatrix,
}
#[derive(Debug, Clone, PartialEq, Eq)]
@@ -899,6 +1062,7 @@ pub(crate) struct TableCatalogDiagnosticsReport {
pub recovery_status: TableCatalogRecoveryStatus,
pub recommended_actions: Vec<TableCatalogRecoveryAction>,
pub commit_recovery: TableCommitRecoveryReport,
pub backing_manifest: TableCatalogBackingManifest,
pub orphan_metadata_candidate_locations: Vec<String>,
}
@@ -1366,6 +1530,15 @@ impl TableCatalogObjectPaths {
)
}
pub fn commit_idempotency_entries_prefix(&self, table_bucket: &str, table_id: &str) -> String {
format!(
"{}{}/{}/",
self.table_bucket_root_prefix(table_bucket),
COMMIT_IDEMPOTENCY_ROOT,
table_catalog_path_hash(table_id)
)
}
fn table_bucket_root_prefix(&self, table_bucket: &str) -> String {
format!("{}/{}/{}/", self.catalog_root, TABLE_BUCKET_ROOT, table_catalog_path_hash(table_bucket))
}
@@ -1377,6 +1550,99 @@ pub(crate) struct ObjectTableCatalogStore<B> {
paths: TableCatalogObjectPaths,
}
fn table_catalog_backing_manifest(
paths: &TableCatalogObjectPaths,
namespace: &Namespace,
table: &IdentifierSegment,
entry: &TableEntry,
commit_recovery: &TableCommitRecoveryReport,
) -> TableCatalogBackingManifest {
let recovery_required = commit_recovery.staged_before_table_update_count > 0
|| commit_recovery.finalization_required_count > 0
|| commit_recovery.idempotency_repair_required_count > 0;
let manual_review_required = commit_recovery.manual_review_count > 0;
let wal_status = if manual_review_required {
TableCatalogWalStatus::ManualReviewRequired
} else if recovery_required {
TableCatalogWalStatus::RecoveryRequired
} else {
TableCatalogWalStatus::Recoverable
};
let migration_status = if manual_review_required {
TableCatalogBackingMigrationStatus::ManualReviewRequired
} else if recovery_required {
TableCatalogBackingMigrationStatus::RecoveryRequired
} else {
TableCatalogBackingMigrationStatus::ReadyToSnapshot
};
let mut blockers = Vec::new();
if recovery_required {
blockers.push(TableCatalogBackingMigrationBlocker::CommitRecoveryRequired);
}
if manual_review_required {
blockers.push(TableCatalogBackingMigrationBlocker::CommitManualReviewRequired);
}
TableCatalogBackingManifest {
version: TABLE_CATALOG_BACKING_MANIFEST_VERSION,
current: TableCatalogBackingProfile {
kind: TableCatalogBackingKind::ObjectBacked,
authority: TableCatalogAuthority::RustfsSysObject,
consistency: TableCatalogConsistencyMode::ConditionalObjectCas,
durability: TableCatalogDurabilityMode::StagedCommitLogBeforePointerUpdate,
current_pointer_path: paths.table_entry_path(&entry.table_bucket, namespace, table),
wal: TableCatalogWalState {
status: wal_status,
commit_log_prefix: paths.commit_log_entries_prefix(&entry.table_bucket, &entry.table_id),
idempotency_index_prefix: paths.commit_idempotency_entries_prefix(&entry.table_bucket, &entry.table_id),
committed_generation: entry.generation,
staged_before_table_update_count: commit_recovery.staged_before_table_update_count,
finalization_required_count: commit_recovery.finalization_required_count,
idempotency_repair_required_count: commit_recovery.idempotency_repair_required_count,
manual_review_count: commit_recovery.manual_review_count,
},
snapshot: TableCatalogSnapshotState {
export_api: "GET /iceberg/v1/{warehouse}/namespaces/{namespace}/tables/{table}/catalog/export".to_string(),
includes_table_bucket: true,
includes_namespace: true,
includes_table_pointer: true,
includes_backing_manifest: true,
},
},
migration: TableCatalogBackingMigrationPlan {
source_kind: TableCatalogBackingKind::ObjectBacked,
target_kind: TableCatalogBackingKind::StrongKvWal,
status: migration_status,
required_steps: vec![
TableCatalogBackingMigrationStep::SnapshotCatalogExport,
TableCatalogBackingMigrationStep::ReplayCommitLog,
TableCatalogBackingMigrationStep::VerifyCurrentPointer,
TableCatalogBackingMigrationStep::EnableSingleWriterFencing,
TableCatalogBackingMigrationStep::CutOverLinearizableReads,
],
blockers,
},
ha: TableCatalogHaPolicy {
writer_region_model: TableCatalogHaWriterModel::SingleActiveWriterRegion,
read_replica_strategy: TableCatalogReadReplicaStrategy::ReadOnlyReplicasForListAndLoad,
commit_read_requirement: TableCatalogCommitReadRequirement::LinearizableLeaderRead,
active_active_supported: false,
failover_requires_operator_promotion: true,
},
scale_validation: TableCatalogScaleValidation {
status: TableCatalogScaleValidationStatus::MatrixPublished,
benchmark_required: true,
required_scenarios: vec![
TableCatalogScaleValidationScenario::ConcurrentCommitCas,
TableCatalogScaleValidationScenario::CommitLogRecoveryReplay,
TableCatalogScaleValidationScenario::MigrationSnapshotReplay,
TableCatalogScaleValidationScenario::ReadReplicaStaleReadGuard,
TableCatalogScaleValidationScenario::ClientConformanceMatrix,
],
},
}
}
impl<B> ObjectTableCatalogStore<B>
where
B: TableCatalogObjectBackend,
@@ -2473,11 +2739,14 @@ where
table.as_str()
)));
};
let commit_recovery = self.table_commit_recovery_report_for_entry(&table_entry, 0).await?;
let backing_manifest = table_catalog_backing_manifest(&self.paths, &namespace, &table, &table_entry, &commit_recovery);
Ok(TableCatalogExport {
table_bucket: table_bucket_entry,
namespace: namespace_entry,
table: table_entry,
backing_manifest,
})
}
@@ -2550,6 +2819,8 @@ where
let commit_recovery = self.plan_table_commit_recovery(table_bucket, namespace, table).await?;
let (recovery_status, recommended_actions) = table_catalog_recovery_summary(&current_metadata_status, &commit_recovery);
let backing_manifest =
table_catalog_backing_manifest(&self.paths, &parsed_namespace, &parsed_table, &catalog.table, &commit_recovery);
Ok(TableCatalogDiagnosticsReport {
catalog,
@@ -2557,6 +2828,7 @@ where
recovery_status,
recommended_actions,
commit_recovery,
backing_manifest,
orphan_metadata_candidate_locations,
})
}
@@ -10148,6 +10420,104 @@ mod tests {
assert_eq!(export.table.generation, 1);
}
#[tokio::test]
async fn export_catalog_entry_includes_backing_migration_manifest() {
let backend = TestCatalogObjectBackend::default();
let store = ObjectTableCatalogStore::new(backend);
let bucket = "analytics";
let namespace = Namespace::parse("sales").unwrap();
let table = IdentifierSegment::parse("orders").unwrap();
let current = default_table_metadata_file_path(&namespace, &table, "00002.metadata.json");
seed_table_for_metadata_maintenance(&store, bucket, &namespace, &table, current.clone()).await;
let export = store.export_table_catalog_entry(bucket, "sales", "orders").await.unwrap();
assert_eq!(export.backing_manifest.version, TABLE_CATALOG_BACKING_MANIFEST_VERSION);
assert_eq!(export.backing_manifest.current.kind, TableCatalogBackingKind::ObjectBacked);
assert_eq!(export.backing_manifest.current.authority, TableCatalogAuthority::RustfsSysObject);
assert_eq!(
export.backing_manifest.current.consistency,
TableCatalogConsistencyMode::ConditionalObjectCas
);
assert_eq!(export.backing_manifest.current.wal.finalization_required_count, 0);
assert_eq!(export.backing_manifest.migration.target_kind, TableCatalogBackingKind::StrongKvWal);
assert_eq!(
export.backing_manifest.migration.status,
TableCatalogBackingMigrationStatus::ReadyToSnapshot
);
assert!(
export
.backing_manifest
.migration
.required_steps
.contains(&TableCatalogBackingMigrationStep::ReplayCommitLog)
);
assert_eq!(
export.backing_manifest.ha.writer_region_model,
TableCatalogHaWriterModel::SingleActiveWriterRegion
);
assert!(!export.backing_manifest.ha.active_active_supported);
}
#[tokio::test]
async fn diagnostics_backing_manifest_requires_recovery_before_migration() {
let backend = TestCatalogObjectBackend::default();
let store = ObjectTableCatalogStore::new(backend.clone());
let bucket = "analytics";
let namespace = Namespace::parse("sales").unwrap();
let table = IdentifierSegment::parse("orders").unwrap();
let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json");
let new_metadata = default_table_metadata_file_path(&namespace, &table, "00002.metadata.json");
let commit_path = TableCatalogObjectPaths::default().commit_log_entry_path(bucket, "table-id", "commit-1");
store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap();
store
.create_namespace(test_namespace_entry(bucket, &namespace))
.await
.unwrap();
store
.create_table(test_table_entry(bucket, &namespace, &table, current_metadata.clone()))
.await
.unwrap();
backend.seed_object(bucket, &new_metadata, b"{}".to_vec()).await;
backend
.fail_put_attempt(rustfs_ecstore::disk::RUSTFS_META_BUCKET, &commit_path, 2)
.await;
store
.commit_table(TableCommitRequest {
table_bucket: bucket.to_string(),
namespace: namespace.public_name(),
table: table.as_str().to_string(),
commit_id: "commit-1".to_string(),
idempotency_key: None,
operation: "append".to_string(),
expected_version_token: "token-v1".to_string(),
expected_metadata_location: current_metadata,
new_metadata_location: new_metadata,
requirements: Vec::new(),
writer: None,
})
.await
.unwrap();
let diagnostics = store.diagnose_table_catalog(bucket, "sales", "orders", 0).await.unwrap();
assert_eq!(diagnostics.backing_manifest.current.wal.finalization_required_count, 1);
assert_eq!(
diagnostics.backing_manifest.migration.status,
TableCatalogBackingMigrationStatus::RecoveryRequired
);
assert!(
diagnostics
.backing_manifest
.migration
.blockers
.contains(&TableCatalogBackingMigrationBlocker::CommitRecoveryRequired)
);
}
#[tokio::test]
async fn consistency_check_reports_missing_metadata_object() {
let backend = TestCatalogObjectBackend::default();
+4
View File
@@ -125,6 +125,7 @@ PyIceberg, PyArrow, or boto3:
python3 scripts/table-catalog/pyiceberg_smoke.py --print-client-matrix
python3 scripts/table-catalog/pyiceberg_smoke.py --print-vendor-profiles
python3 scripts/table-catalog/pyiceberg_smoke.py --print-unsupported-inventory
python3 scripts/table-catalog/pyiceberg_smoke.py --print-production-readiness
```
Use these outputs when updating release notes, PR descriptions, or follow-up
@@ -143,6 +144,9 @@ The smoke test also probes catalog-backed advanced Iceberg surfaces:
execution checks
- catalog diagnostics exposes the table recovery and consistency state used by
operators
- catalog export and diagnostics expose the current catalog backing manifest,
recoverable commit-log WAL state, strong backing migration target, single
active writer HA policy, and scale validation matrix
## Client Matrix
+30
View File
@@ -178,6 +178,24 @@ 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",
},
{
"capability": "single-active-writer-ha",
"status": "policy-published",
"validation": "diagnostics publish single-active-writer-region semantics, linearizable leader-read requirements for commits, read-only replica limits, and explicit no active-active support",
},
{
"capability": "scale-validation-matrix",
"status": "matrix-published",
"validation": "required validation scenarios are machine-readable: concurrent commit CAS, commit-log recovery replay, migration snapshot replay, stale replica read guards, and long-running client conformance",
},
]
@dataclass(frozen=True)
class RuntimeDeps:
@@ -225,6 +243,10 @@ def unsupported_inventory() -> list[dict[str, str]]:
return [entry.copy() for entry in UNSUPPORTED_INVENTORY]
def production_readiness_inventory() -> list[dict[str, str]]:
return [entry.copy() for entry in PRODUCTION_READINESS_INVENTORY]
def normalized_rest_path(rest_path: str) -> str:
stripped = rest_path.strip()
if not stripped:
@@ -278,6 +300,11 @@ def parse_args() -> argparse.Namespace:
parser.add_argument("--print-client-matrix", action="store_true", help="Print the current client conformance matrix as JSON and exit.")
parser.add_argument("--print-vendor-profiles", action="store_true", help="Print vendor connection profile references as JSON and exit.")
parser.add_argument("--print-unsupported-inventory", action="store_true", help="Print unsupported capability inventory as JSON and exit.")
parser.add_argument(
"--print-production-readiness",
action="store_true",
help="Print catalog backing, HA, and scale-readiness inventory as JSON and exit.",
)
return apply_profile_defaults(parser.parse_args())
@@ -700,6 +727,9 @@ def printed_metadata(args: argparse.Namespace) -> bool:
if args.print_unsupported_inventory:
print_json_document({"unsupported_inventory": unsupported_inventory()})
printed = True
if args.print_production_readiness:
print_json_document({"production_readiness": production_readiness_inventory()})
printed = True
return printed
@@ -753,11 +753,25 @@ class PyIcebergSmokeConfigTest(unittest.TestCase):
self.assertIn("status", entry)
self.assertIn("roadmap_area", entry)
def test_production_readiness_inventory_tracks_catalog_backing(self) -> None:
inventory = pyiceberg_smoke.production_readiness_inventory()
capabilities = {entry["capability"] for entry in inventory}
self.assertIn("strong-catalog-backing", capabilities)
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")
for entry in inventory:
self.assertIn("status", entry)
self.assertIn("validation", entry)
def test_published_table_catalog_docs_do_not_use_internal_roadmap_labels(self) -> None:
readme = (SCRIPT_DIR / "README.md").read_text(encoding="utf-8")
self.assertIsNone(INTERNAL_ROADMAP_LABEL_RE.search(readme))
self.assertIsNone(INTERNAL_ROADMAP_LABEL_RE.search(str(pyiceberg_smoke.unsupported_inventory())))
self.assertIsNone(INTERNAL_ROADMAP_LABEL_RE.search(str(pyiceberg_smoke.production_readiness_inventory())))
if __name__ == "__main__":