diff --git a/rustfs/src/table_catalog.rs b/rustfs/src/table_catalog.rs index 4355d6829..f82a53be5 100644 --- a/rustfs/src/table_catalog.rs +++ b/rustfs/src/table_catalog.rs @@ -1670,6 +1670,12 @@ pub(crate) struct TableCatalogObject { pub mod_time: Option, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct TableCatalogObjectMetadata { + pub etag: Option, + pub mod_time: Option, +} + #[derive(Debug, Clone, PartialEq, Eq)] pub(crate) enum TableCatalogPutPrecondition { Any, @@ -1685,6 +1691,16 @@ pub(crate) trait TableCatalogObjectBackend: Clone + Send + Sync + 'static { self.read_object(bucket, object).await } + async fn object_metadata(&self, bucket: &str, object: &str) -> TableCatalogStoreResult> { + Ok(self + .read_object(bucket, object) + .await? + .map(|object| TableCatalogObjectMetadata { + etag: object.etag, + mod_time: object.mod_time, + })) + } + async fn object_exists(&self, bucket: &str, object: &str) -> TableCatalogStoreResult; async fn put_object( @@ -2033,6 +2049,7 @@ fn table_catalog_backing_manifest( type StrongNamespaceKey = (String, String); type StrongResourceKey = (String, String, String); type StrongCommitKey = (String, String, String); +type StrongWarehouseIndex = BTreeMap>; #[derive(Clone, Default)] struct StrongTableCatalogState { @@ -2044,6 +2061,7 @@ struct StrongTableCatalogState { views: BTreeMap, commits: BTreeMap, idempotency: BTreeMap, + warehouse_index: StrongWarehouseIndex, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -2142,6 +2160,52 @@ where } } + fn rebuild_warehouse_index_locked(state: &mut StrongTableCatalogState) -> TableCatalogStoreResult<()> { + let mut warehouse_index: StrongWarehouseIndex = BTreeMap::new(); + for ((table_bucket, namespace, table), entry) in &state.tables { + if entry.state != TableCatalogEntryState::Active { + continue; + } + if !state + .table_buckets + .get(table_bucket) + .is_some_and(|entry| entry.state == TableCatalogEntryState::Active) + { + continue; + } + if !state + .namespaces + .get(&(table_bucket.clone(), namespace.clone())) + .is_some_and(|entry| entry.state == TableCatalogEntryState::Active) + { + continue; + } + let Ok(warehouse_object_prefix) = table_warehouse_object_prefix(entry) else { + continue; + }; + let table_key = (table_bucket.clone(), namespace.clone(), table.clone()); + if let Some(existing_key) = warehouse_index + .entry(table_bucket.clone()) + .or_default() + .insert(warehouse_object_prefix.clone(), table_key.clone()) + { + return Err(TableCatalogStoreError::Invalid(format!( + "duplicate active table warehouse location in strong catalog snapshot: {warehouse_object_prefix} is owned by {}/{}/{} and {}/{}/{}", + existing_key.0, existing_key.1, existing_key.2, table_key.0, table_key.1, table_key.2 + ))); + } + } + state.warehouse_index = warehouse_index; + Ok(()) + } + + fn snapshot_from_mutated_state_locked( + state: &mut StrongTableCatalogState, + ) -> TableCatalogStoreResult { + Self::rebuild_warehouse_index_locked(state)?; + Ok(Self::snapshot_from_state_locked(state)) + } + fn state_from_snapshot( snapshot: StrongTableCatalogSnapshot, snapshot_etag: Option, @@ -2193,6 +2257,7 @@ where record.commit, ); } + Self::rebuild_warehouse_index_locked(&mut state)?; Ok(state) } @@ -2205,12 +2270,31 @@ where }) } - fn snapshot_write_context_locked(state: &StrongTableCatalogState) -> (TableCatalogPutPrecondition, StrongTableCatalogState) { + fn snapshot_draft_context_locked(state: &StrongTableCatalogState) -> (TableCatalogPutPrecondition, StrongTableCatalogState) { (Self::snapshot_write_precondition_locked(state), state.clone()) } async fn hydrate_state(&self) -> TableCatalogStoreResult<()> { - self.reload_state_from_durable().await + let Some(current_snapshot_etag) = ({ + let state = self.state.lock().await; + if state.hydrated { + Some(state.snapshot_etag.clone()) + } else { + None + } + }) else { + return self.reload_state_from_durable().await; + }; + + let snapshot_metadata = self + .object_backend + .object_metadata(RUSTFS_META_BUCKET, &Self::snapshot_object_path()) + .await?; + match (snapshot_metadata, current_snapshot_etag.as_deref()) { + (None, None) => Ok(()), + (Some(metadata), Some(current_etag)) if metadata.etag.as_deref() == Some(current_etag) => Ok(()), + _ => self.reload_state_from_durable().await, + } } async fn reload_state_from_durable(&self) -> TableCatalogStoreResult<()> { @@ -2248,15 +2332,11 @@ where &self, snapshot: StrongTableCatalogSnapshot, precondition: TableCatalogPutPrecondition, - rollback_state: StrongTableCatalogState, ) -> TableCatalogStoreResult<()> { match self.persist_snapshot(snapshot, precondition).await { Ok(()) => self.reload_state_from_durable().await, Err(err) => { - if self.reload_state_from_durable().await.is_err() { - let mut state = self.state.lock().await; - *state = rollback_state; - } + let _ = self.reload_state_from_durable().await; Err(err) } } @@ -2581,13 +2661,13 @@ where return Err(TableCatalogStoreError::Invalid("unsupported table bucket catalog type".to_string())); } - let (snapshot, precondition, rollback_state) = { - let mut state = self.state.lock().await; - let (precondition, rollback_state) = Self::snapshot_write_context_locked(&state); - state.table_buckets.insert(entry.table_bucket.clone(), entry); - (Self::snapshot_from_state_locked(&state), precondition, rollback_state) + let (snapshot, precondition) = { + let state = self.state.lock().await; + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + draft_state.table_buckets.insert(entry.table_bucket.clone(), entry); + (Self::snapshot_from_mutated_state_locked(&mut draft_state)?, precondition) }; - self.finalize_snapshot_write(snapshot, precondition, rollback_state).await + self.finalize_snapshot_write(snapshot, precondition).await } async fn create_namespace(&self, entry: NamespaceEntry) -> TableCatalogStoreResult<()> { @@ -2596,8 +2676,8 @@ where 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 (snapshot, precondition, rollback_state) = { - let mut state = self.state.lock().await; + let (snapshot, precondition) = { + let 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!( @@ -2605,11 +2685,11 @@ where entry.table_bucket, entry.namespace ))); } - let (precondition, rollback_state) = Self::snapshot_write_context_locked(&state); - state.namespaces.insert(key, entry); - (Self::snapshot_from_state_locked(&state), precondition, rollback_state) + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + draft_state.namespaces.insert(key, entry); + (Self::snapshot_from_mutated_state_locked(&mut draft_state)?, precondition) }; - self.finalize_snapshot_write(snapshot, precondition, rollback_state).await + self.finalize_snapshot_write(snapshot, precondition).await } async fn list_namespaces(&self, table_bucket: &str) -> TableCatalogStoreResult> { @@ -2637,8 +2717,8 @@ where self.hydrate_state().await?; let namespace = parse_namespace_for_store(namespace)?; let key = Self::namespace_key(table_bucket, &namespace); - let (snapshot, precondition, rollback_state) = { - let mut state = self.state.lock().await; + let (snapshot, precondition) = { + let state = self.state.lock().await; if !state.namespaces.contains_key(&key) { return Err(TableCatalogStoreError::NotFound(format!( "namespace {}/{}", @@ -2661,11 +2741,11 @@ where namespace.public_name() ))); } - let (precondition, rollback_state) = Self::snapshot_write_context_locked(&state); - state.namespaces.remove(&key); - (Self::snapshot_from_state_locked(&state), precondition, rollback_state) + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + draft_state.namespaces.remove(&key); + (Self::snapshot_from_mutated_state_locked(&mut draft_state)?, precondition) }; - self.finalize_snapshot_write(snapshot, precondition, rollback_state).await + self.finalize_snapshot_write(snapshot, precondition).await } async fn create_table(&self, entry: TableEntry) -> TableCatalogStoreResult<()> { @@ -2680,8 +2760,8 @@ where let table = parse_table_for_store(&entry.table)?; table_warehouse_object_prefix(&entry)?; let key = Self::table_key(&entry.table_bucket, &namespace, &table); - let (snapshot, precondition, rollback_state) = { - let mut state = self.state.lock().await; + let (snapshot, precondition) = { + let state = self.state.lock().await; Self::require_table_bucket_in_state(&state, &entry.table_bucket)?; if !state .namespaces @@ -2699,11 +2779,11 @@ where ))); } Self::ensure_table_warehouse_prefix_available_locked(&state, &entry, &key)?; - let (precondition, rollback_state) = Self::snapshot_write_context_locked(&state); - state.tables.insert(key, entry); - (Self::snapshot_from_state_locked(&state), precondition, rollback_state) + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + draft_state.tables.insert(key, entry); + (Self::snapshot_from_mutated_state_locked(&mut draft_state)?, precondition) }; - self.finalize_snapshot_write(snapshot, precondition, rollback_state).await + self.finalize_snapshot_write(snapshot, precondition).await } async fn list_tables(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult> { @@ -2728,6 +2808,42 @@ where Ok(state.tables.get(&Self::table_key(table_bucket, &namespace, &table)).cloned()) } + async fn resolve_table_data_plane_resource( + &self, + table_bucket: &str, + object: &str, + ) -> TableCatalogStoreResult> { + if table_bucket.is_empty() || object.is_empty() { + return Ok(None); + } + + self.hydrate_state().await?; + let state = self.state.lock().await; + let Some(bucket_entry) = state.table_buckets.get(table_bucket) else { + return Ok(None); + }; + if bucket_entry.state != TableCatalogEntryState::Active { + return Ok(None); + } + + let Some(bucket_index) = state.warehouse_index.get(table_bucket) else { + return Ok(None); + }; + + for warehouse_object_prefix in warehouse_index_candidate_prefixes(object) { + if let Some(table_key) = bucket_index.get(warehouse_object_prefix) { + let Some(table) = state.tables.get(table_key) else { + continue; + }; + return Ok(Some(table_data_plane_resource_from_entry( + table.clone(), + warehouse_object_prefix.to_string(), + ))); + } + } + Ok(None) + } + async fn commit_table(&self, request: TableCommitRequest) -> TableCatalogStoreResult { let _write_guard = self.write_lock.lock().await; self.hydrate_state().await?; @@ -2738,17 +2854,15 @@ where let key = Self::table_key(&request.table_bucket, &namespace, &table); let committed_existing_result = { - let mut state = self.state.lock().await; + let state = self.state.lock().await; let current = Self::validate_new_table_commit_locked(&state, &key, &request, &namespace, &table); match current { Ok(current) => { - let (precondition, rollback_state) = Self::snapshot_write_context_locked(&state); - if let Some(result) = Self::committed_existing_result_locked(&mut state, &request, current) { - let snapshot = Self::snapshot_from_state_locked(&state); - Some((result, snapshot, precondition, rollback_state)) - } else { - None - } + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + Self::committed_existing_result_locked(&mut draft_state, &request, current).map(|result| { + Self::snapshot_from_mutated_state_locked(&mut draft_state) + .map(|snapshot| (result, snapshot, precondition)) + }) } Err(error) => { return table_commit_result( @@ -2763,11 +2877,13 @@ where } } }; - if let Some((result, snapshot, precondition, rollback_state)) = committed_existing_result { - let result = self - .finalize_snapshot_write(snapshot, precondition, rollback_state) - .await - .map(|_| result); + if let Some(prepared_result) = committed_existing_result { + let result = match prepared_result { + Ok((result, snapshot, precondition)) => { + self.finalize_snapshot_write(snapshot, precondition).await.map(|_| result) + } + Err(err) => Err(err), + }; return table_commit_result( &request.table_bucket, &request.namespace, @@ -2801,21 +2917,19 @@ where table_metadata_warehouse_location(&request.table_bucket, &request.new_metadata_location, &new_metadata_object)?; let cas_started = Instant::now(); - let (result, snapshot, precondition, rollback_state) = { - let mut state = self.state.lock().await; - let (precondition, rollback_state) = Self::snapshot_write_context_locked(&state); - let result = Self::apply_commit_locked(&mut state, &request, &namespace, &table, next_warehouse_location); - let snapshot = result.as_ref().ok().map(|_| Self::snapshot_from_state_locked(&state)); - let precondition = result.as_ref().ok().map(|_| precondition); - let rollback_state = result.as_ref().ok().map(|_| rollback_state); - (result, snapshot, precondition, rollback_state) + let prepared_result = { + let state = self.state.lock().await; + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + match Self::apply_commit_locked(&mut draft_state, &request, &namespace, &table, next_warehouse_location) { + Ok(result) => { + Self::snapshot_from_mutated_state_locked(&mut draft_state).map(|snapshot| (result, snapshot, precondition)) + } + Err(err) => Err(err), + } }; - let result = match (result, snapshot, precondition, rollback_state) { - (Ok(result), Some(snapshot), Some(precondition), Some(rollback_state)) => self - .finalize_snapshot_write(snapshot, precondition, rollback_state) - .await - .map(|_| result), - (result, _, _, _) => result, + let result = match prepared_result { + Ok((result, snapshot, precondition)) => self.finalize_snapshot_write(snapshot, precondition).await.map(|_| result), + Err(err) => Err(err), }; let cas_result = result.as_ref().map(|_| ()).map_err(Clone::clone); record_table_commit_cas_result(&request.operation, cas_started, &cas_result); @@ -2836,10 +2950,9 @@ where let namespace = parse_namespace_for_store(namespace)?; let table = parse_table_for_store(table)?; let key = Self::table_key(table_bucket, &namespace, &table); - let (snapshot, precondition, rollback_state) = { - let mut state = self.state.lock().await; - let (precondition, rollback_state) = Self::snapshot_write_context_locked(&state); - if state.tables.remove(&key).is_none() { + let (snapshot, precondition) = { + let state = self.state.lock().await; + if !state.tables.contains_key(&key) { return Err(TableCatalogStoreError::NotFound(format!( "table {}/{}/{}", table_bucket, @@ -2847,9 +2960,11 @@ where table.as_str() ))); } - (Self::snapshot_from_state_locked(&state), precondition, rollback_state) + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + draft_state.tables.remove(&key); + (Self::snapshot_from_mutated_state_locked(&mut draft_state)?, precondition) }; - self.finalize_snapshot_write(snapshot, precondition, rollback_state).await + self.finalize_snapshot_write(snapshot, precondition).await } async fn create_view(&self, entry: ViewEntry) -> TableCatalogStoreResult<()> { @@ -2860,8 +2975,8 @@ where 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 (snapshot, precondition, rollback_state) = { - let mut state = self.state.lock().await; + let (snapshot, precondition) = { + let state = self.state.lock().await; Self::require_table_bucket_in_state(&state, &entry.table_bucket)?; if !state .namespaces @@ -2878,11 +2993,11 @@ where entry.table_bucket, entry.namespace, entry.view ))); } - let (precondition, rollback_state) = Self::snapshot_write_context_locked(&state); - state.views.insert(key, entry); - (Self::snapshot_from_state_locked(&state), precondition, rollback_state) + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + draft_state.views.insert(key, entry); + (Self::snapshot_from_mutated_state_locked(&mut draft_state)?, precondition) }; - self.finalize_snapshot_write(snapshot, precondition, rollback_state).await + self.finalize_snapshot_write(snapshot, precondition).await } async fn list_views(&self, table_bucket: &str, namespace: &str) -> TableCatalogStoreResult> { @@ -2931,36 +3046,37 @@ where view_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 (precondition, rollback_state) = Self::snapshot_write_context_locked(&state); - 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 (snapshot, precondition, next) = { + let 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()); - let snapshot = Self::snapshot_from_state_locked(&state); - drop(state); - self.finalize_snapshot_write(snapshot, precondition, rollback_state).await?; + 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); + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + draft_state.views.insert(key, next.clone()); + (Self::snapshot_from_mutated_state_locked(&mut draft_state)?, precondition, next) + }; + self.finalize_snapshot_write(snapshot, precondition).await?; Ok(ViewCommitResult { view: next }) } @@ -2970,10 +3086,9 @@ where let namespace = parse_namespace_for_store(namespace)?; let view = parse_table_for_store(view)?; let key = Self::table_key(table_bucket, &namespace, &view); - let (snapshot, precondition, rollback_state) = { - let mut state = self.state.lock().await; - let (precondition, rollback_state) = Self::snapshot_write_context_locked(&state); - if state.views.remove(&key).is_none() { + let (snapshot, precondition) = { + let state = self.state.lock().await; + if !state.views.contains_key(&key) { return Err(TableCatalogStoreError::NotFound(format!( "view {}/{}/{}", table_bucket, @@ -2981,9 +3096,11 @@ where view.as_str() ))); } - (Self::snapshot_from_state_locked(&state), precondition, rollback_state) + let (precondition, mut draft_state) = Self::snapshot_draft_context_locked(&state); + draft_state.views.remove(&key); + (Self::snapshot_from_mutated_state_locked(&mut draft_state)?, precondition) }; - self.finalize_snapshot_write(snapshot, precondition, rollback_state).await + self.finalize_snapshot_write(snapshot, precondition).await } async fn get_commit_by_id( @@ -6699,6 +6816,17 @@ where .await } + async fn object_metadata(&self, bucket: &str, object: &str) -> TableCatalogStoreResult> { + match self.store.get_object_info(bucket, object, &ObjectOptions::default()).await { + Ok(info) => Ok(Some(TableCatalogObjectMetadata { + etag: info.etag, + mod_time: info.mod_time, + })), + Err(err) if is_missing_storage_error(&err) => Ok(None), + Err(err) => Err(storage_error_to_catalog("stat catalog object", err)), + } + } + async fn object_exists(&self, bucket: &str, object: &str) -> TableCatalogStoreResult { match self.store.get_object_info(bucket, object, &ObjectOptions::default()).await { Ok(_) => Ok(true), @@ -10558,17 +10686,35 @@ mod tests { type TestCatalogObjectLock = Arc>; type TestCatalogObjectLocks = Arc>>; + #[derive(Clone, Default)] + struct TestCatalogObjectPutPause { + started: Arc, + release: Arc, + } + + impl TestCatalogObjectPutPause { + async fn wait_started(&self) { + self.started.notified().await; + } + + fn release(&self) { + self.release.notify_one(); + } + } + #[derive(Default)] struct TestCatalogObjectState { objects: BTreeMap<(String, String), TestCatalogObjectRecord>, fail_read_attempts: BTreeMap<(String, String), BTreeSet>, read_attempts: BTreeMap<(String, String), usize>, fail_put_attempts: BTreeMap<(String, String), BTreeSet>, + pause_put_attempts: BTreeMap<(String, String), BTreeMap>, fail_delete_attempts: BTreeMap<(String, String), BTreeSet>, put_attempts: BTreeMap<(String, String), usize>, delete_attempts: BTreeMap<(String, String), usize>, write_lock_acquisitions: BTreeMap<(String, String), usize>, read_calls: usize, + metadata_calls: usize, list_calls: usize, next_etag: u64, } @@ -10620,9 +10766,14 @@ mod tests { self.state.lock().await.read_calls } + async fn metadata_call_count(&self) -> usize { + self.state.lock().await.metadata_calls + } + async fn reset_call_counts(&self) { let mut state = self.state.lock().await; state.read_calls = 0; + state.metadata_calls = 0; state.list_calls = 0; } @@ -10649,6 +10800,19 @@ mod tests { let next_attempt = state.put_attempts.get(&key).copied().unwrap_or_default() + 1; state.fail_put_attempts.entry(key).or_default().insert(next_attempt); } + + async fn pause_next_put(&self, bucket: &str, object: &str) -> TestCatalogObjectPutPause { + let mut state = self.state.lock().await; + let key = (bucket.to_string(), object.to_string()); + let next_attempt = state.put_attempts.get(&key).copied().unwrap_or_default() + 1; + let pause = TestCatalogObjectPutPause::default(); + state + .pause_put_attempts + .entry(key) + .or_default() + .insert(next_attempt, pause.clone()); + pause + } } impl TestCatalogObjectState { @@ -10983,6 +11147,22 @@ mod tests { })) } + async fn object_metadata( + &self, + bucket: &str, + object: &str, + ) -> TableCatalogStoreResult> { + let mut state = self.state.lock().await; + state.metadata_calls += 1; + Ok(state + .objects + .get(&(bucket.to_string(), object.to_string())) + .map(|record| TableCatalogObjectMetadata { + etag: Some(record.etag.clone()), + mod_time: record.mod_time, + })) + } + async fn object_exists(&self, bucket: &str, object: &str) -> TableCatalogStoreResult { let state = self.state.lock().await; Ok(state.objects.contains_key(&(bucket.to_string(), object.to_string()))) @@ -10995,22 +11175,34 @@ mod tests { data: Vec, precondition: TableCatalogPutPrecondition, ) -> TableCatalogStoreResult<()> { - let mut state = self.state.lock().await; let key = (bucket.to_string(), object.to_string()); - let attempt = { - let attempts = state.put_attempts.entry(key.clone()).or_default(); - *attempts += 1; - *attempts + let pause = { + let mut state = self.state.lock().await; + let attempt = { + let attempts = state.put_attempts.entry(key.clone()).or_default(); + *attempts += 1; + *attempts + }; + if state + .fail_put_attempts + .get(&key) + .is_some_and(|attempts| attempts.contains(&attempt)) + { + return Err(TableCatalogStoreError::Internal(format!( + "injected put failure for {object} attempt {attempt}" + ))); + } + state + .pause_put_attempts + .get_mut(&key) + .and_then(|attempts| attempts.remove(&attempt)) }; - if state - .fail_put_attempts - .get(&key) - .is_some_and(|attempts| attempts.contains(&attempt)) - { - return Err(TableCatalogStoreError::Internal(format!( - "injected put failure for {object} attempt {attempt}" - ))); + if let Some(pause) = pause { + pause.started.notify_one(); + pause.release.notified().await; } + + let mut state = self.state.lock().await; match precondition { TableCatalogPutPrecondition::IfAbsent if state.objects.contains_key(&key) => { return Err(TableCatalogStoreError::Conflict(format!("object already exists: {object}"))); @@ -15828,6 +16020,177 @@ mod tests { assert_eq!(loaded_after_commit.version_token, result.table.version_token); } + #[tokio::test] + async fn strong_catalog_backing_skips_snapshot_body_when_etag_is_unchanged() { + let backend = TestCatalogObjectBackend::default(); + let store = StrongTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + + store + .put_table_bucket(test_bucket_entry(bucket)) + .await + .expect("table bucket should be created"); + + backend.reset_call_counts().await; + let loaded = store + .get_table_bucket(bucket) + .await + .expect("table bucket load should succeed") + .expect("table bucket should exist"); + + assert_eq!(loaded.table_bucket, bucket); + assert_eq!(backend.metadata_call_count().await, 1); + assert_eq!(backend.read_call_count().await, 0); + assert_eq!(backend.list_call_count().await, 0); + } + + #[tokio::test] + async fn strong_catalog_backing_resolves_data_plane_resource_without_catalog_scan() { + let backend = TestCatalogObjectBackend::default(); + let store = StrongTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").expect("namespace should parse"); + let table = IdentifierSegment::parse("orders").expect("table should parse"); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + + store + .put_table_bucket(test_bucket_entry(bucket)) + .await + .expect("table bucket should be created"); + store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .expect("namespace should be created"); + store + .create_table(test_table_entry(bucket, &namespace, &table, current_metadata)) + .await + .expect("table should be created"); + + backend.reset_call_counts().await; + let resource = table_data_plane_resource_for_object(&store, bucket, "tables/table-id/data/file.parquet") + .await + .expect("data-plane lookup should succeed") + .expect("data-plane resource should resolve"); + + assert_eq!(resource.table_bucket, bucket); + assert_eq!(resource.namespace, namespace.public_name()); + assert_eq!(resource.table, table.as_str()); + assert_eq!(resource.table_id, "table-id"); + assert_eq!(resource.warehouse_object_prefix, "tables/table-id/"); + assert_eq!(backend.metadata_call_count().await, 1); + assert_eq!(backend.read_call_count().await, 0); + assert_eq!(backend.list_call_count().await, 0); + } + + #[tokio::test] + async fn strong_catalog_backing_does_not_publish_warehouse_index_before_snapshot_persist() { + let backend = TestCatalogObjectBackend::default(); + let store = StrongTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").expect("namespace should parse"); + let table = IdentifierSegment::parse("orders").expect("table should parse"); + let current_metadata = default_table_metadata_file_path(&namespace, &table, "00001.metadata.json"); + let snapshot_path = StrongTableCatalogStore::::snapshot_object_path(); + + store + .put_table_bucket(test_bucket_entry(bucket)) + .await + .expect("table bucket should be created"); + store + .create_namespace(test_namespace_entry(bucket, &namespace)) + .await + .expect("namespace should be created"); + + let pause = backend.pause_next_put(RUSTFS_META_BUCKET, &snapshot_path).await; + let create_store = store.clone(); + let create_task = tokio::spawn(async move { + let namespace = Namespace::parse("sales").expect("namespace should parse"); + let table = IdentifierSegment::parse("orders").expect("table should parse"); + create_store + .create_table(test_table_entry("analytics", &namespace, &table, current_metadata)) + .await + }); + pause.wait_started().await; + + let resource_before_commit = table_data_plane_resource_for_object(&store, bucket, "tables/table-id/data/file.parquet") + .await + .expect("data-plane lookup should succeed while snapshot persist is paused"); + assert!(resource_before_commit.is_none()); + assert!( + store + .load_table(bucket, "sales", "orders") + .await + .expect("table load should succeed while snapshot persist is paused") + .is_none() + ); + + pause.release(); + create_task + .await + .expect("create task should join") + .expect("table should be created"); + + let resource = table_data_plane_resource_for_object(&store, bucket, "tables/table-id/data/file.parquet") + .await + .expect("data-plane lookup should succeed after snapshot persist") + .expect("new table warehouse index should be visible after commit"); + assert_eq!(resource.table_bucket, bucket); + assert_eq!(resource.namespace, namespace.public_name()); + assert_eq!(resource.table, table.as_str()); + assert_eq!(resource.table_id, "table-id"); + assert_eq!(resource.warehouse_object_prefix, "tables/table-id/"); + } + + #[tokio::test] + async fn strong_catalog_backing_rejects_duplicate_snapshot_warehouse_index_entries() { + let backend = TestCatalogObjectBackend::default(); + let store = StrongTableCatalogStore::new(backend.clone()); + let bucket = "analytics"; + let namespace = Namespace::parse("sales").expect("namespace should parse"); + let orders = IdentifierSegment::parse("orders").expect("table should parse"); + let returns = IdentifierSegment::parse("returns").expect("table should parse"); + let mut returns_entry = test_table_entry( + bucket, + &namespace, + &returns, + default_table_metadata_file_path(&namespace, &returns, "00001.metadata.json"), + ); + returns_entry.table_id = "table-id-2".to_string(); + returns_entry.table_uuid = "table-uuid-2".to_string(); + + let snapshot = StrongTableCatalogSnapshot { + version: STRONG_TABLE_CATALOG_SNAPSHOT_VERSION, + table_buckets: vec![test_bucket_entry(bucket)], + namespaces: vec![test_namespace_entry(bucket, &namespace)], + tables: vec![ + test_table_entry( + bucket, + &namespace, + &orders, + default_table_metadata_file_path(&namespace, &orders, "00001.metadata.json"), + ), + returns_entry, + ], + views: Vec::new(), + commits: Vec::new(), + idempotency: Vec::new(), + }; + backend + .seed_object( + RUSTFS_META_BUCKET, + &StrongTableCatalogStore::::snapshot_object_path(), + serde_json::to_vec(&snapshot).expect("strong snapshot should encode"), + ) + .await; + + let err = store + .get_table_bucket(bucket) + .await + .expect_err("duplicate warehouse prefix should fail snapshot hydration"); + + assert_matches!(err, TableCatalogStoreError::Invalid(message) if message.contains("duplicate active table warehouse location")); + } + #[tokio::test] async fn strong_catalog_backing_rolls_back_cache_when_persist_and_reload_fail() { let backend = TestCatalogObjectBackend::default();