diff --git a/rustfs/src/table_catalog.rs b/rustfs/src/table_catalog.rs index 4df1d8efd..b72b8f7bf 100644 --- a/rustfs/src/table_catalog.rs +++ b/rustfs/src/table_catalog.rs @@ -1658,6 +1658,661 @@ fn table_catalog_backing_manifest( } } +#[cfg(test)] +type StrongNamespaceKey = (String, String); +#[cfg(test)] +type StrongResourceKey = (String, String, String); +#[cfg(test)] +type StrongCommitKey = (String, String, String); + +#[cfg(test)] +#[derive(Default)] +struct StrongTableCatalogState { + table_buckets: BTreeMap, + namespaces: BTreeMap, + tables: BTreeMap, + views: BTreeMap, + commits: BTreeMap, + idempotency: BTreeMap, +} + +#[cfg(test)] +#[derive(Clone)] +pub(crate) struct StrongTableCatalogStore { + object_backend: B, + state: Arc>, +} + +#[cfg(test)] +impl StrongTableCatalogStore +where + B: TableCatalogObjectBackend, +{ + pub fn new(object_backend: B) -> Self { + Self { + object_backend, + state: Arc::new(tokio::sync::Mutex::new(StrongTableCatalogState::default())), + } + } + + fn namespace_key(table_bucket: &str, namespace: &Namespace) -> StrongNamespaceKey { + (table_bucket.to_string(), namespace.public_name()) + } + + fn table_key(table_bucket: &str, namespace: &Namespace, table: &IdentifierSegment) -> StrongResourceKey { + (table_bucket.to_string(), namespace.public_name(), table.as_str().to_string()) + } + + fn commit_key(table_bucket: &str, table_id: &str, commit_id: &str) -> StrongCommitKey { + (table_bucket.to_string(), table_id.to_string(), commit_id.to_string()) + } + + fn idempotency_key(table_bucket: &str, table_id: &str, idempotency_key: &str) -> StrongCommitKey { + (table_bucket.to_string(), table_id.to_string(), idempotency_key.to_string()) + } + + 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}"))); + } + Ok(()) + } + + fn table_commit_recovery_report_for_entry_locked( + state: &StrongTableCatalogState, + entry: &TableEntry, + ) -> TableCommitRecoveryReport { + let mut commits = state + .commits + .iter() + .filter(|((table_bucket, table_id, _), _)| table_bucket == &entry.table_bucket && table_id == &entry.table_id) + .map(|((_, _, _), commit_log)| { + let idempotency_commit = commit_log.idempotency_key.as_deref().and_then(|idempotency_key| { + state + .idempotency + .get(&Self::idempotency_key(&entry.table_bucket, &entry.table_id, idempotency_key)) + }); + table_commit_recovery_entry(entry, commit_log, idempotency_commit) + }) + .collect::>(); + commits.sort_by(|left, right| left.commit_id.cmp(&right.commit_id)); + + let finalization_required_count = commits + .iter() + .filter(|commit| matches!(commit.recovery_state, TableCommitRecoveryState::FinalizationRequired)) + .count(); + let idempotency_repair_required_count = commits + .iter() + .filter(|commit| matches!(commit.recovery_state, TableCommitRecoveryState::IdempotencyIndexRepairRequired)) + .count(); + let staged_before_table_update_count = commits + .iter() + .filter(|commit| matches!(commit.recovery_state, TableCommitRecoveryState::StagedBeforeTableUpdate)) + .count(); + let manual_review_count = commits + .iter() + .filter(|commit| matches!(commit.recovery_state, TableCommitRecoveryState::ManualReview)) + .count(); + let finalized_count = commits + .iter() + .filter(|commit| matches!(commit.recovery_state, TableCommitRecoveryState::Committed)) + .count(); + + TableCommitRecoveryReport { + table_bucket: entry.table_bucket.clone(), + namespace: entry.namespace.clone(), + table: entry.table.clone(), + table_id: entry.table_id.clone(), + current_metadata_location: entry.metadata_location.clone(), + current_version_token: entry.version_token.clone(), + current_generation: entry.generation, + commits, + staged_before_table_update_count, + finalization_required_count, + idempotency_repair_required_count, + manual_review_count, + finalized_count, + } + } + + fn validate_new_table_commit_locked( + state: &StrongTableCatalogState, + key: &StrongResourceKey, + request: &TableCommitRequest, + namespace: &Namespace, + table: &IdentifierSegment, + ) -> TableCatalogStoreResult { + let Some(current) = state.tables.get(key).cloned() else { + return Err(TableCatalogStoreError::NotFound(format!( + "table {}/{}/{}", + request.table_bucket, request.namespace, request.table + ))); + }; + let commit_key = Self::commit_key(&request.table_bucket, ¤t.table_id, &request.commit_id); + let existing_commit = state.commits.get(&commit_key); + let idempotency_key = request + .idempotency_key + .as_deref() + .map(|idempotency_key| Self::idempotency_key(&request.table_bucket, ¤t.table_id, idempotency_key)); + let existing_idempotency_commit = idempotency_key.as_ref().and_then(|key| state.idempotency.get(key)); + + if let Some(existing) = existing_commit { + if !commit_log_matches_request(existing, request, ¤t.table_id) { + return Err(TableCatalogStoreError::Conflict(format!( + "commit id already exists: {}", + request.commit_id + ))); + } + if matches!(existing.status, CommitLogStatus::Committed) || table_matches_committed_log(¤t, existing) { + return Ok(current); + } + return Err(TableCatalogStoreError::Conflict( + "existing commit record does not match current table state".to_string(), + )); + } + if let Some(existing) = existing_idempotency_commit + && !commit_log_matches_request(existing, request, ¤t.table_id) + { + return Err(TableCatalogStoreError::Conflict("idempotency key already exists".to_string())); + } + if existing_idempotency_commit.is_some() { + return Err(TableCatalogStoreError::Conflict( + "idempotency key exists without a recoverable commit record".to_string(), + )); + } + if current.version_token != request.expected_version_token { + return Err(TableCatalogStoreError::Conflict( + "current table version token does not match expected token".to_string(), + )); + } + if current.metadata_location != request.expected_metadata_location { + return Err(TableCatalogStoreError::Conflict( + "current table metadata location does not match expected location".to_string(), + )); + } + if !is_valid_table_metadata_location(namespace, table, &request.new_metadata_location) { + return Err(TableCatalogStoreError::Invalid( + "new metadata location must be inside the table metadata directory".to_string(), + )); + } + Ok(current) + } + + fn committed_existing_result_locked( + state: &mut StrongTableCatalogState, + request: &TableCommitRequest, + current: TableEntry, + ) -> Option { + let commit_key = Self::commit_key(&request.table_bucket, ¤t.table_id, &request.commit_id); + let existing = state.commits.get(&commit_key)?; + if !commit_log_matches_request(existing, request, ¤t.table_id) { + return None; + } + if !matches!(existing.status, CommitLogStatus::Committed) && !table_matches_committed_log(¤t, existing) { + return None; + } + + let mut committed = existing.clone(); + committed.status = CommitLogStatus::Committed; + state.commits.insert(commit_key, committed.clone()); + if let Some(idempotency_key) = committed.idempotency_key.as_deref() { + state.idempotency.insert( + Self::idempotency_key(&request.table_bucket, ¤t.table_id, idempotency_key), + committed.clone(), + ); + } + Some(TableCommitResult { + table: current, + commit_log: committed, + }) + } + + fn apply_commit_locked( + state: &mut StrongTableCatalogState, + request: &TableCommitRequest, + namespace: &Namespace, + table: &IdentifierSegment, + next_warehouse_location: Option, + ) -> TableCatalogStoreResult { + let key = Self::table_key(&request.table_bucket, namespace, table); + let current = Self::validate_new_table_commit_locked(state, &key, request, namespace, table)?; + if let Some(result) = Self::committed_existing_result_locked(state, request, current.clone()) { + return Ok(result); + } + + let commit_log = CommitLogEntry { + version: TABLE_CATALOG_ENTRY_VERSION, + commit_id: request.commit_id.clone(), + idempotency_key: request.idempotency_key.clone(), + table_id: current.table_id.clone(), + operation: request.operation.clone(), + expected_version_token: request.expected_version_token.clone(), + new_version_token: format!("token-{}", Uuid::new_v4()), + previous_metadata_location: current.metadata_location.clone(), + new_metadata_location: request.new_metadata_location.clone(), + requirements: request.requirements.clone(), + status: CommitLogStatus::Committed, + writer: request.writer.clone(), + created_at: None, + updated_at: None, + }; + + let mut next = current; + next.metadata_location = commit_log.new_metadata_location.clone(); + if let Some(warehouse_location) = next_warehouse_location { + next.warehouse_location = warehouse_location; + } + next.version_token = commit_log.new_version_token.clone(); + next.generation = next.generation.saturating_add(1); + + let commit_key = Self::commit_key(&request.table_bucket, &next.table_id, &request.commit_id); + state.commits.insert(commit_key, commit_log.clone()); + if let Some(idempotency_key) = request.idempotency_key.as_deref() { + state.idempotency.insert( + Self::idempotency_key(&request.table_bucket, &next.table_id, idempotency_key), + commit_log.clone(), + ); + } + state.tables.insert(key, next.clone()); + + Ok(TableCommitResult { table: next, commit_log }) + } + + pub(crate) async fn plan_table_commit_recovery( + &self, + table_bucket: &str, + namespace: &str, + table: &str, + ) -> TableCatalogStoreResult { + let namespace = parse_namespace_for_store(namespace)?; + let table = parse_table_for_store(table)?; + let key = Self::table_key(table_bucket, &namespace, &table); + let state = self.state.lock().await; + let Some(entry) = state.tables.get(&key) else { + return Err(TableCatalogStoreError::NotFound(format!( + "table {}/{}/{}", + table_bucket, + namespace.public_name(), + table.as_str() + ))); + }; + Ok(Self::table_commit_recovery_report_for_entry_locked(&state, entry)) + } +} + +#[cfg(test)] +#[async_trait::async_trait] +impl TableCatalogStore for StrongTableCatalogStore +where + B: TableCatalogObjectBackend, +{ + async fn get_table_bucket(&self, table_bucket: &str) -> TableCatalogStoreResult> { + let state = self.state.lock().await; + Ok(state.table_buckets.get(table_bucket).cloned()) + } + + async fn put_table_bucket(&self, entry: TableBucketEntry) -> TableCatalogStoreResult<()> { + validate_catalog_entry_version("table bucket", entry.version)?; + if entry.table_bucket.is_empty() { + return Err(TableCatalogStoreError::Invalid("table bucket name cannot be empty".to_string())); + } + if entry.catalog_type != TABLE_BUCKET_CATALOG_TYPE { + return Err(TableCatalogStoreError::Invalid("unsupported table bucket catalog type".to_string())); + } + + let mut state = self.state.lock().await; + state.table_buckets.insert(entry.table_bucket.clone(), entry); + Ok(()) + } + + async fn create_namespace(&self, entry: NamespaceEntry) -> TableCatalogStoreResult<()> { + validate_catalog_entry_version("namespace", entry.version)?; + let namespace = parse_namespace_for_store(&entry.namespace)?; + let key = Self::namespace_key(&entry.table_bucket, &namespace); + let mut state = self.state.lock().await; + Self::require_table_bucket_in_state(&state, &entry.table_bucket)?; + if state.namespaces.contains_key(&key) { + return Err(TableCatalogStoreError::Conflict(format!( + "catalog object already exists: namespace {}/{}", + entry.table_bucket, entry.namespace + ))); + } + state.namespaces.insert(key, entry); + Ok(()) + } + + async fn list_namespaces(&self, table_bucket: &str) -> TableCatalogStoreResult> { + let state = self.state.lock().await; + let mut entries = state + .namespaces + .iter() + .filter(|((bucket, _), _)| bucket == table_bucket) + .map(|(_, entry)| entry.clone()) + .collect::>(); + entries.sort_by(|left, right| left.namespace.cmp(&right.namespace)); + Ok(entries) + } + + async fn get_namespace(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult> { + let namespace = parse_namespace_for_store(namespace)?; + let state = self.state.lock().await; + Ok(state.namespaces.get(&Self::namespace_key(table_bucket, &namespace)).cloned()) + } + + async fn drop_namespace(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult<()> { + let namespace = parse_namespace_for_store(namespace)?; + let key = Self::namespace_key(table_bucket, &namespace); + let mut state = self.state.lock().await; + if !state.namespaces.contains_key(&key) { + return Err(TableCatalogStoreError::NotFound(format!( + "namespace {}/{}", + table_bucket, + namespace.public_name() + ))); + } + if state + .tables + .keys() + .any(|(bucket, namespace_name, _)| bucket == table_bucket && namespace_name == &namespace.public_name()) + || state + .views + .keys() + .any(|(bucket, namespace_name, _)| bucket == table_bucket && namespace_name == &namespace.public_name()) + { + return Err(TableCatalogStoreError::Conflict(format!( + "namespace {}/{} is not empty", + table_bucket, + namespace.public_name() + ))); + } + state.namespaces.remove(&key); + Ok(()) + } + + async fn create_table(&self, entry: TableEntry) -> TableCatalogStoreResult<()> { + self.register_table(entry).await + } + + async fn register_table(&self, entry: TableEntry) -> TableCatalogStoreResult<()> { + validate_catalog_entry_version("table", entry.version)?; + let namespace = parse_namespace_for_store(&entry.namespace)?; + let table = parse_table_for_store(&entry.table)?; + table_warehouse_object_prefix(&entry)?; + let key = Self::table_key(&entry.table_bucket, &namespace, &table); + let mut state = self.state.lock().await; + Self::require_table_bucket_in_state(&state, &entry.table_bucket)?; + if !state + .namespaces + .contains_key(&Self::namespace_key(&entry.table_bucket, &namespace)) + { + return Err(TableCatalogStoreError::NotFound(format!( + "namespace {}/{}", + entry.table_bucket, entry.namespace + ))); + } + if state.tables.contains_key(&key) { + return Err(TableCatalogStoreError::Conflict(format!( + "catalog object already exists: table {}/{}/{}", + entry.table_bucket, entry.namespace, entry.table + ))); + } + state.tables.insert(key, entry); + Ok(()) + } + + async fn list_tables(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult> { + let namespace = parse_namespace_for_store(namespace)?; + let state = self.state.lock().await; + let mut entries = state + .tables + .iter() + .filter(|((bucket, namespace_name, _), _)| bucket == table_bucket && namespace_name == &namespace.public_name()) + .map(|(_, entry)| entry.clone()) + .collect::>(); + entries.sort_by(|left, right| left.table.cmp(&right.table)); + Ok(entries) + } + + async fn load_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 state = self.state.lock().await; + Ok(state.tables.get(&Self::table_key(table_bucket, &namespace, &table)).cloned()) + } + + async fn commit_table(&self, request: TableCommitRequest) -> TableCatalogStoreResult { + let commit_started = Instant::now(); + record_table_commit_attempt(&request.operation); + let namespace = parse_namespace_for_store(&request.namespace)?; + let table = parse_table_for_store(&request.table)?; + let key = Self::table_key(&request.table_bucket, &namespace, &table); + + { + let mut state = self.state.lock().await; + let current = Self::validate_new_table_commit_locked(&state, &key, &request, &namespace, &table); + match current { + Ok(current) => { + if let Some(result) = Self::committed_existing_result_locked(&mut state, &request, current) { + return table_commit_result( + &request.table_bucket, + &request.namespace, + &request.table, + &request.commit_id, + &request.operation, + commit_started, + Ok(result), + ); + } + } + Err(error) => { + return table_commit_result( + &request.table_bucket, + &request.namespace, + &request.table, + &request.commit_id, + &request.operation, + commit_started, + Err(error), + ); + } + } + } + + let Some(new_metadata_object) = self + .object_backend + .read_object(&request.table_bucket, &request.new_metadata_location) + .await? + else { + return table_commit_result( + &request.table_bucket, + &request.namespace, + &request.table, + &request.commit_id, + &request.operation, + commit_started, + Err(TableCatalogStoreError::NotFound(format!( + "new metadata object {}", + request.new_metadata_location + ))), + ); + }; + let next_warehouse_location = + metadata_warehouse_location(&request.table_bucket, &request.new_metadata_location, &new_metadata_object)?; + + let cas_started = Instant::now(); + let result = { + let mut state = self.state.lock().await; + Self::apply_commit_locked(&mut state, &request, &namespace, &table, next_warehouse_location) + }; + let cas_result = result.as_ref().map(|_| ()).map_err(Clone::clone); + record_table_commit_cas_result(&request.operation, cas_started, &cas_result); + table_commit_result( + &request.table_bucket, + &request.namespace, + &request.table, + &request.commit_id, + &request.operation, + commit_started, + result, + ) + } + + 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 key = Self::table_key(table_bucket, &namespace, &table); + let mut state = self.state.lock().await; + if state.tables.remove(&key).is_none() { + return Err(TableCatalogStoreError::NotFound(format!( + "table {}/{}/{}", + table_bucket, + namespace.public_name(), + table.as_str() + ))); + } + Ok(()) + } + + async fn create_view(&self, entry: ViewEntry) -> TableCatalogStoreResult<()> { + validate_catalog_entry_version("view", entry.version)?; + let namespace = parse_namespace_for_store(&entry.namespace)?; + let view = parse_table_for_store(&entry.view)?; + let key = Self::table_key(&entry.table_bucket, &namespace, &view); + let mut state = self.state.lock().await; + Self::require_table_bucket_in_state(&state, &entry.table_bucket)?; + if !state + .namespaces + .contains_key(&Self::namespace_key(&entry.table_bucket, &namespace)) + { + return Err(TableCatalogStoreError::NotFound(format!( + "namespace {}/{}", + entry.table_bucket, entry.namespace + ))); + } + if state.views.contains_key(&key) { + return Err(TableCatalogStoreError::Conflict(format!( + "catalog object already exists: view {}/{}/{}", + entry.table_bucket, entry.namespace, entry.view + ))); + } + state.views.insert(key, entry); + Ok(()) + } + + async fn list_views(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult> { + let namespace = parse_namespace_for_store(namespace)?; + let state = self.state.lock().await; + let mut entries = state + .views + .iter() + .filter(|((bucket, namespace_name, _), _)| bucket == table_bucket && namespace_name == &namespace.public_name()) + .map(|(_, entry)| entry.clone()) + .collect::>(); + entries.sort_by(|left, right| left.view.cmp(&right.view)); + Ok(entries) + } + + async fn load_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 state = self.state.lock().await; + Ok(state.views.get(&Self::table_key(table_bucket, &namespace, &view)).cloned()) + } + + async fn replace_view(&self, request: ViewCommitRequest) -> TableCatalogStoreResult { + let namespace = parse_namespace_for_store(&request.namespace)?; + let view = parse_table_for_store(&request.view)?; + if !is_valid_view_metadata_location(&namespace, &view, &request.new_metadata_location) { + return Err(TableCatalogStoreError::Invalid( + "new metadata location must be inside the view metadata directory".to_string(), + )); + } + let Some(new_metadata_object) = self + .object_backend + .read_object(&request.table_bucket, &request.new_metadata_location) + .await? + else { + return Err(TableCatalogStoreError::NotFound(format!( + "new view metadata object {}", + request.new_metadata_location + ))); + }; + let next_warehouse_location = + metadata_warehouse_location(&request.table_bucket, &request.new_metadata_location, &new_metadata_object)?; + + let key = Self::table_key(&request.table_bucket, &namespace, &view); + let mut state = self.state.lock().await; + let Some(current) = state.views.get(&key).cloned() else { + return Err(TableCatalogStoreError::NotFound(format!( + "view {}/{}/{}", + request.table_bucket, request.namespace, request.view + ))); + }; + if current.version_token != request.expected_version_token { + return Err(TableCatalogStoreError::Conflict( + "current view version token does not match expected token".to_string(), + )); + } + if current.metadata_location != request.expected_metadata_location { + return Err(TableCatalogStoreError::Conflict( + "current view metadata location does not match expected location".to_string(), + )); + } + + let mut next = current; + next.metadata_location = request.new_metadata_location; + if let Some(warehouse_location) = next_warehouse_location { + next.warehouse_location = warehouse_location; + } + next.version_token = format!("token-{}", Uuid::new_v4()); + next.generation = next.generation.saturating_add(1); + state.views.insert(key, next.clone()); + Ok(ViewCommitResult { view: next }) + } + + 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 key = Self::table_key(table_bucket, &namespace, &view); + let mut state = self.state.lock().await; + if state.views.remove(&key).is_none() { + return Err(TableCatalogStoreError::NotFound(format!( + "view {}/{}/{}", + table_bucket, + namespace.public_name(), + view.as_str() + ))); + } + Ok(()) + } + + async fn get_commit_by_id( + &self, + table_bucket: &str, + table_id: &str, + commit_id: &str, + ) -> TableCatalogStoreResult> { + let state = self.state.lock().await; + Ok(state + .commits + .get(&Self::commit_key(table_bucket, table_id, commit_id)) + .cloned()) + } + + async fn get_commit_by_idempotency_key( + &self, + table_bucket: &str, + table_id: &str, + idempotency_key: &str, + ) -> TableCatalogStoreResult> { + let state = self.state.lock().await; + Ok(state + .idempotency + .get(&Self::idempotency_key(table_bucket, table_id, idempotency_key)) + .cloned()) + } +} + impl ObjectTableCatalogStore where B: TableCatalogObjectBackend, @@ -10690,6 +11345,158 @@ mod tests { assert!(!export.backing_manifest.ha.active_active_supported); } + #[tokio::test] + async fn strong_catalog_backing_commit_is_atomic_with_wal_and_idempotency() { + let backend = TestCatalogObjectBackend::default(); + let store = StrongTableCatalogStore::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"); + + 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; + + let result = 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: Some("client-request".to_string()), + operation: "append".to_string(), + expected_version_token: "token-v1".to_string(), + expected_metadata_location: current_metadata, + new_metadata_location: new_metadata.clone(), + requirements: Vec::new(), + writer: Some("pyiceberg/test".to_string()), + }) + .await + .unwrap(); + + assert_eq!(result.table.metadata_location, new_metadata); + assert_eq!(result.table.generation, 2); + assert_eq!(result.commit_log.status, CommitLogStatus::Committed); + + let loaded = store.load_table(bucket, "sales", "orders").await.unwrap().unwrap(); + assert_eq!(loaded.metadata_location, result.table.metadata_location); + assert_eq!(loaded.version_token, result.table.version_token); + assert_eq!( + store + .get_commit_by_id(bucket, "table-id", "commit-1") + .await + .unwrap() + .unwrap() + .status, + CommitLogStatus::Committed + ); + assert_eq!( + store + .get_commit_by_idempotency_key(bucket, "table-id", "client-request") + .await + .unwrap() + .unwrap() + .status, + CommitLogStatus::Committed + ); + + let recovery = store.plan_table_commit_recovery(bucket, "sales", "orders").await.unwrap(); + assert_eq!(recovery.staged_before_table_update_count, 0); + assert_eq!(recovery.finalization_required_count, 0); + assert_eq!(recovery.idempotency_repair_required_count, 0); + assert_eq!(recovery.manual_review_count, 0); + } + + #[tokio::test] + async fn strong_catalog_backing_commit_conflict_keeps_pointer_and_wal_unchanged() { + let backend = TestCatalogObjectBackend::default(); + let store = StrongTableCatalogStore::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"); + + 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; + + let err = 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: Some("client-request".to_string()), + operation: "append".to_string(), + expected_version_token: "stale-token".to_string(), + expected_metadata_location: current_metadata.clone(), + new_metadata_location: new_metadata, + requirements: Vec::new(), + writer: Some("pyiceberg/test".to_string()), + }) + .await + .unwrap_err(); + + assert_matches!(err, TableCatalogStoreError::Conflict(_)); + let loaded = store.load_table(bucket, "sales", "orders").await.unwrap().unwrap(); + assert_eq!(loaded.metadata_location, current_metadata); + assert_eq!(loaded.version_token, "token-v1"); + assert!( + store + .get_commit_by_id(bucket, "table-id", "commit-1") + .await + .unwrap() + .is_none() + ); + assert!( + store + .get_commit_by_idempotency_key(bucket, "table-id", "client-request") + .await + .unwrap() + .is_none() + ); + } + + #[tokio::test] + async fn strong_catalog_backing_rejects_invalid_table_warehouse_location() { + let backend = TestCatalogObjectBackend::default(); + let store = StrongTableCatalogStore::new(backend); + 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"); + + store.put_table_bucket(test_bucket_entry(bucket)).await.unwrap(); + store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .unwrap(); + let mut entry = test_table_entry(bucket, &namespace, &table, current_metadata); + entry.warehouse_location = "s3://other-bucket/tables/table-id".to_string(); + + let err = store.create_table(entry).await.unwrap_err(); + + assert_matches!(err, TableCatalogStoreError::Invalid(_)); + assert!(store.load_table(bucket, "sales", "orders").await.unwrap().is_none()); + } + #[tokio::test] async fn diagnostics_backing_manifest_requires_recovery_before_migration() { let backend = TestCatalogObjectBackend::default();