diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 456d38ac7..c296e7ee2 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -241,8 +241,8 @@ pub mod cache { pub mod capacity { pub use crate::core::pools::{ - PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free, path2_bucket_object, - path2_bucket_object_with_base_path, + DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free, + path2_bucket_object, path2_bucket_object_with_base_path, }; pub use crate::store::utils::is_reserved_or_invalid_bucket; } diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 9d8e0d01a..b509ef6ad 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -804,6 +804,88 @@ fn ensure_decommission_generation(meta: &PoolMeta, idx: usize, generation: Offse } } +fn record_decommission_unresolved_entry( + meta: &mut PoolMeta, + idx: usize, + generation: OffsetDateTime, + entry: DecommissionUnresolvedEntry, +) -> Result { + ensure_decommission_generation(meta, idx, generation)?; + if entry.pool_index != idx || entry.source_generation != generation { + return Err(Error::other("decommission unresolved entry does not match the active pool generation")); + } + + let pool_count = meta.pools.len(); + let Some(pool) = meta.pools.get_mut(idx) else { + return Err(invalid_decommission_pool_index_error(pool_count, idx)); + }; + let Some(info) = pool.decommission.as_mut() else { + return Err(decommission_metadata_not_initialized_error("record decommission unresolved entry")); + }; + let existing = info.unresolved_entries.iter_mut().find(|existing| { + existing.bucket == entry.bucket + && existing.object == entry.object + && existing.pool_index == entry.pool_index + && existing.set_index == entry.set_index + && existing.source_generation == entry.source_generation + }); + let changed = match existing { + Some(existing) if existing == &entry => false, + Some(existing) => { + *existing = entry; + true + } + None => { + info.unresolved_entries.push(entry); + true + } + }; + if changed { + pool.last_update = OffsetDateTime::now_utc(); + } + Ok(changed) +} + +fn reconcile_decommission_unresolved_entries_for_completion( + meta: &mut PoolMeta, + idx: usize, + verified_generation: Option, +) -> Result<()> { + let pool_count = meta.pools.len(); + let Some(pool) = meta.pools.get(idx) else { + return Err(invalid_decommission_pool_index_error(pool_count, idx)); + }; + let Some(info) = pool.decommission.as_ref() else { + return Err(decommission_metadata_not_initialized_error("reconcile decommission unresolved entries")); + }; + if info.unresolved_entries.is_empty() { + return Ok(()); + } + + let Some(generation) = verified_generation else { + return Err(Error::other(format!( + "failed to complete decommission for pool {idx}: {} unresolved listing entries remain", + info.unresolved_entries.len() + ))); + }; + ensure_decommission_generation(meta, idx, generation)?; + if info + .unresolved_entries + .iter() + .any(|entry| entry.source_generation != generation) + { + return Err(Error::other(format!( + "failed to complete decommission for pool {idx}: unresolved listing ledger contains a different generation" + ))); + } + + let Some(info) = meta.pools.get_mut(idx).and_then(|pool| pool.decommission.as_mut()) else { + return Err(decommission_metadata_not_initialized_error("reconcile decommission unresolved entries")); + }; + info.unresolved_entries.clear(); + Ok(()) +} + async fn run_decommission_side_effect( rx: &CancellationToken, operation_gate: &Arc>, @@ -981,18 +1063,15 @@ fn resolve_decommission_listing_error(listing_error: Option, entry_error: } } -fn decommission_unresolved_listing_error( - bucket: &str, - prefix: &str, - candidate: Option<&str>, - candidate_count: usize, - disk_error_count: usize, - pool_index: usize, - set_index: usize, -) -> Error { - let location = candidate.unwrap_or(prefix); +fn decommission_unresolved_listing_error(entry: &DecommissionUnresolvedEntry) -> Error { Error::other(format!( - "decommission listing could not resolve metadata for {bucket}/{location} on pool {pool_index} set {set_index} ({candidate_count} candidate(s), {disk_error_count} disk error(s))" + "decommission listing could not resolve metadata for {bucket}/{object} on pool {pool_index} set {set_index} ({candidate_count} candidate(s), {disk_error_count} disk error(s))", + bucket = entry.bucket, + object = entry.object, + pool_index = entry.pool_index, + set_index = entry.set_index, + candidate_count = entry.candidate_count, + disk_error_count = entry.disk_error_count, )) } @@ -1004,22 +1083,32 @@ fn resolve_decommission_partial_listing_entry( disk_error_count: usize, pool_index: usize, set_index: usize, -) -> Result { + source_generation: OffsetDateTime, +) -> std::result::Result { let candidate_count = entries.as_ref().iter().flatten().count(); if let Some(entry) = entries.resolve(resolver) { return Ok(entry); } - let candidate = entries.as_ref().iter().flatten().map(|entry| entry.name.as_str()).next(); - Err(decommission_unresolved_listing_error( - bucket, - prefix, - candidate, + let object = entries + .as_ref() + .iter() + .flatten() + .map(|entry| entry.name.as_str()) + .next() + .unwrap_or(prefix) + .to_string(); + Err(DecommissionUnresolvedEntry { + bucket: bucket.to_string(), + object, candidate_count, disk_error_count, pool_index, set_index, - )) + source_generation, + observed_at: OffsetDateTime::now_utc(), + reason: "metadata_resolution_failed".to_string(), + }) } fn resolve_decommission_pool_meta_reload_result(result: Result<()>, stage: &str) -> Result<()> { @@ -1624,6 +1713,8 @@ struct PersistedPoolDecommissionInfo { pub terminal_reload_attempt_at: Option, #[serde(rename = "terminalReloadFailures", default)] pub terminal_reload_failures: Vec, + #[serde(rename = "unresolvedEntries", default)] + pub unresolved_entries: Vec, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -1755,6 +1846,7 @@ impl TryFrom for PoolDecommissionInfo { bytes_failed: value.bytes_failed, terminal_reload_attempt_at: value.terminal_reload_attempt_at, terminal_reload_failures: value.terminal_reload_failures, + unresolved_entries: value.unresolved_entries, progress_save_item_baseline: value.items_decommissioned.saturating_add(value.items_decommission_failed), progress_save_retry_after: None, }) @@ -1787,6 +1879,7 @@ impl TryFrom for PoolDecommissionInfo { bytes_failed: value.bytes_failed, terminal_reload_attempt_at: None, terminal_reload_failures: Vec::new(), + unresolved_entries: Vec::new(), progress_save_item_baseline: value.items_decommissioned.saturating_add(value.items_decommission_failed), progress_save_retry_after: None, }) @@ -1835,6 +1928,7 @@ impl From<&PoolDecommissionInfo> for PersistedPoolDecommissionInfo { bytes_failed: value.bytes_failed, terminal_reload_attempt_at: value.terminal_reload_attempt_at, terminal_reload_failures: value.terminal_reload_failures.clone(), + unresolved_entries: value.unresolved_entries.clone(), } } } @@ -2435,6 +2529,26 @@ pub fn path2_bucket_object_with_base_path(base_path: &str, path: &str) -> (Strin path_to_bucket_object_with_base_path(base_path, path) } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct DecommissionUnresolvedEntry { + pub bucket: String, + pub object: String, + #[serde(rename = "poolIndex")] + pub pool_index: usize, + #[serde(rename = "setIndex")] + pub set_index: usize, + #[serde(rename = "sourceGeneration", with = "time::serde::rfc3339")] + pub source_generation: OffsetDateTime, + #[serde(rename = "candidateCount")] + pub candidate_count: usize, + #[serde(rename = "diskErrorCount")] + pub disk_error_count: usize, + #[serde(rename = "observedAt", with = "time::serde::rfc3339")] + pub observed_at: OffsetDateTime, + pub reason: String, +} + #[derive(Debug, Clone, Serialize, Deserialize, Default)] pub struct PoolDecommissionInfo { #[serde(rename = "startTime", with = "time::serde::rfc3339::option")] @@ -2480,6 +2594,8 @@ pub struct PoolDecommissionInfo { #[serde(rename = "terminalReloadFailures", default)] pub terminal_reload_failures: Vec, #[serde(skip)] + pub unresolved_entries: Vec, + #[serde(skip)] pub progress_save_item_baseline: usize, #[serde(skip)] pub progress_save_retry_after: Option, @@ -2513,6 +2629,7 @@ impl PoolDecommissionInfo { || self.bytes_failed > 0 || self.terminal_reload_attempt_at.is_some() || !self.terminal_reload_failures.is_empty() + || !self.unresolved_entries.is_empty() } fn counted_items(&self) -> usize { @@ -2830,6 +2947,21 @@ impl ECStore { snapshot.save(self.pools.clone()).await } + async fn persist_decommission_unresolved_entry( + &self, + idx: usize, + generation: OffsetDateTime, + entry: DecommissionUnresolvedEntry, + ) -> Result<()> { + { + let mut pool_meta = self.pool_meta.write().await; + record_decommission_unresolved_entry(&mut pool_meta, idx, generation, entry)?; + } + self.save_current_pool_meta() + .await + .map_err(|err| Error::other(format!("decommission unresolved entry ledger save failed: {err}"))) + } + async fn save_decommission_progress_checkpoint(&self, idx: usize, generation: OffsetDateTime) -> Result { // Lock order: save gate, then the short pool metadata read/write sections. Peer // reloads are intentionally performed by the caller after both locks are released. @@ -3678,6 +3810,7 @@ impl ECStore { let list_bi = bi.clone(); let list_outstanding = outstanding.clone(); let list_entry_error = entry_error.clone(); + let list_store = self.clone(); let mut listing = tokio::spawn(async move { run_decommission_listing_with_retry_and_drain( list_rx.clone(), @@ -3691,8 +3824,9 @@ impl ECStore { let rx = list_rx_for_list.clone(); let bucket = list_bi.clone(); let entry_error = list_entry_error.clone(); + let store = list_store.clone(); async move { - set.list_objects_to_decommission(rx, bucket, callback, entry_error, idx, set_idx) + set.list_objects_to_decommission(store, rx, bucket, callback, entry_error, idx, set_idx, generation) .await } }, @@ -4593,7 +4727,7 @@ impl ECStore { state = "verifying_completion", "Decommission completion verification started" ); - if let Err(err) = self.check_after_decommission(idx).await { + if let Err(err) = self.check_after_decommission(idx, generation).await { resolve_decommission_terminal_mark_result( self.decommission_failed_for_operation(idx, canceler).await, "failed", @@ -4620,7 +4754,7 @@ impl ECStore { "Decommission marking completed state" ); resolve_decommission_terminal_mark_result( - self.complete_decommission_for_operation(idx, canceler).await, + self.complete_decommission_for_operation(idx, canceler, generation).await, "completed", &cmd_line, )?; @@ -4756,14 +4890,25 @@ impl ECStore { #[tracing::instrument(skip(self))] pub async fn complete_decommission(&self, idx: usize) -> Result<()> { - self.complete_decommission_with_owner(idx, None).await + self.complete_decommission_with_owner(idx, None, None).await } - async fn complete_decommission_for_operation(&self, idx: usize, owner: &DecommissionCanceler) -> Result<()> { - self.complete_decommission_with_owner(idx, Some(owner)).await + async fn complete_decommission_for_operation( + &self, + idx: usize, + owner: &DecommissionCanceler, + verified_generation: OffsetDateTime, + ) -> Result<()> { + self.complete_decommission_with_owner(idx, Some(owner), Some(verified_generation)) + .await } - async fn complete_decommission_with_owner(&self, idx: usize, owner: Option<&DecommissionCanceler>) -> Result<()> { + async fn complete_decommission_with_owner( + &self, + idx: usize, + owner: Option<&DecommissionCanceler>, + verified_generation: Option, + ) -> Result<()> { ensure_decommission_terminal_operation_supported(self.single_pool(), "complete decommission")?; let _start_guard = self.start_gate.lock().await; @@ -4775,11 +4920,13 @@ impl ECStore { let previous_pool_meta = pool_meta.clone(); let Some(changed) = update_decommission_for_operation(cancelers.as_slice(), &mut pool_meta, idx, owner, |pool_meta| { - pool_meta.decommission_complete(idx) + reconcile_decommission_unresolved_entries_for_completion(pool_meta, idx, verified_generation)?; + Ok(pool_meta.decommission_complete(idx)) }) else { return Ok(()); }; + let changed = changed?; let terminal_canceler = if let Some(owner) = owner { Some(owner.clone()) } else { @@ -5139,7 +5286,7 @@ impl ECStore { Ok(ret) } - async fn check_after_decommission(self: &Arc, idx: usize) -> Result<()> { + async fn check_after_decommission(self: &Arc, idx: usize, generation: OffsetDateTime) -> Result<()> { let buckets = self.get_buckets_to_decommission().await?; let pool = self.pools[idx].clone(); @@ -5238,7 +5385,16 @@ impl ECStore { }); let list_result = set - .list_objects_to_decommission(callback_rx, bucket_info.clone(), callback, entry_error.clone(), idx, set_index) + .list_objects_to_decommission( + self.clone(), + callback_rx, + bucket_info.clone(), + callback, + entry_error.clone(), + idx, + set_index, + generation, + ) .await; let entry_error = entry_error.lock().await.clone(); resolve_decommission_check_after_list_result(list_result, entry_error)?; @@ -5607,6 +5763,17 @@ mod tests { #[test] fn pool_meta_persists_decommission_resume_queues() { let start_time = OffsetDateTime::now_utc(); + let unresolved_entry = DecommissionUnresolvedEntry { + bucket: "bucket-b".to_string(), + object: "prefix/unresolved.txt".to_string(), + pool_index: 0, + set_index: 1, + source_generation: start_time, + candidate_count: 2, + disk_error_count: 1, + observed_at: start_time, + reason: "metadata_resolution_failed".to_string(), + }; let pool_meta = PoolMeta { version: POOL_META_VERSION, pools: vec![PoolStatus { @@ -5627,6 +5794,7 @@ mod tests { bytes_failed: 128, terminal_reload_attempt_at: Some(start_time), terminal_reload_failures: vec!["complete_decommission: peer node-a failed".to_string()], + unresolved_entries: vec![unresolved_entry.clone()], ..Default::default() }), }], @@ -5666,6 +5834,7 @@ mod tests { restored_decommission.terminal_reload_failures, vec!["complete_decommission: peer node-a failed".to_string()] ); + assert_eq!(restored_decommission.unresolved_entries, vec![unresolved_entry]); assert!(restored_decommission.queued); assert_eq!(restored_decommission.items_since_last_progress_save(), 0); } @@ -5757,6 +5926,7 @@ mod tests { assert!(decommission.bucket.is_empty()); assert!(decommission.prefix.is_empty()); assert!(decommission.object.is_empty()); + assert!(decommission.unresolved_entries.is_empty()); } #[test] @@ -6046,28 +6216,32 @@ async fn record_decommission_entry_error( entry_error: &Arc>>, rx: &CancellationToken, err: Error, -) { +) -> bool { if rx.is_cancelled() { - return; + return false; } let mut first_err = entry_error.lock().await; if first_err.is_none() && !rx.is_cancelled() { *first_err = Some(err); rx.cancel(); + return true; } + false } impl SetDisks { - #[tracing::instrument(skip(self, rx, cb_func, entry_error))] + #[tracing::instrument(skip(self, store, rx, cb_func, entry_error))] async fn list_objects_to_decommission( self: &Arc, + store: Arc, rx: CancellationToken, bucket_info: DecomBucketInfo, cb_func: ListCallback, entry_error: Arc>>, pool_index: usize, set_index: usize, + source_generation: OffsetDateTime, ) -> Result<()> { let (disks, _) = self.get_online_disks_with_healing(false).await; ensure_decommission_listing_disks_available(!disks.is_empty(), &bucket_info.name)?; @@ -6088,6 +6262,8 @@ impl SetDisks { let unresolved_prefix = bucket_info.prefix.clone(); let unresolved_pool_index = pool_index; let unresolved_set_index = set_index; + let unresolved_generation = source_generation; + let unresolved_store = store; list_path_raw( rx, @@ -6107,8 +6283,10 @@ impl SetDisks { let prefix = unresolved_prefix.clone(); let unresolved_error = unresolved_error.clone(); let unresolved_rx = unresolved_rx.clone(); + let unresolved_store = unresolved_store.clone(); let pool_index = unresolved_pool_index; let set_index = unresolved_set_index; + let source_generation = unresolved_generation; let disk_error_count = errs.iter().flatten().count(); if unresolved_rx.is_cancelled() { return Box::pin(async {}); @@ -6122,6 +6300,7 @@ impl SetDisks { disk_error_count, pool_index, set_index, + source_generation, ) { Ok(entry) => { warn!("decommission_pool: list_objects_to_decommission get {}", &entry.name); @@ -6129,10 +6308,11 @@ impl SetDisks { cb_func(entry).await; }) } - Err(err) => Box::pin(async move { + Err(unresolved_entry) => Box::pin(async move { if unresolved_rx.is_cancelled() { return; } + let err = decommission_unresolved_listing_error(&unresolved_entry); warn!( event = EVENT_DECOMMISSION_BUCKET, component = LOG_COMPONENT_ECSTORE, @@ -6143,6 +6323,13 @@ impl SetDisks { error = %err, "Decommission listing failed closed on unresolved metadata" ); + let err = match unresolved_store + .persist_decommission_unresolved_entry(pool_index, source_generation, unresolved_entry) + .await + { + Ok(()) => err, + Err(ledger_err) => Error::other(format!("{err}; {ledger_err}")), + }; record_decommission_entry_error(&unresolved_error, &unresolved_rx, err).await; }), } @@ -6351,47 +6538,52 @@ mod pools_tests { use super::{ DECOMMISSION_ENTRY_CONCURRENCY_DEFAULT_CAP, DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP, DECOMMISSION_ENTRY_QUEUE_HARD_CAP, DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, DecommissionCanceler, - DecommissionEntryEnqueueResult, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, - PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, QueuedDecommissionEntry, apply_decommission_status_space_info, - await_decommission_worker, bind_decommission_cancelers, bind_missing_decommission_cancelers, - cancel_decommission_canceler, clamp_decommission_entry_concurrency, classify_decommission_terminal_state, - count_decommission_item, decommission_cancel_signal_result, decommission_entry_queue_capacity, decommission_item_size, - decommission_meta_bucket_options, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, - default_decommission_entry_concurrency, drain_decommission_entry_queue, enqueue_decommission_entry, - ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_generation, - ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, - ensure_decommission_start_keeps_active_pool, ensure_decommission_start_local_leader, + DecommissionEntryEnqueueResult, DecommissionStartPoolState, DecommissionTerminalState, DecommissionUnresolvedEntry, + ListCallback, POOL_META_VERSION, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, QueuedDecommissionEntry, + apply_decommission_status_space_info, await_decommission_worker, bind_decommission_cancelers, + bind_missing_decommission_cancelers, cancel_decommission_canceler, clamp_decommission_entry_concurrency, + classify_decommission_terminal_state, count_decommission_item, decommission_cancel_signal_result, + decommission_entry_queue_capacity, decommission_item_size, decommission_meta_bucket_options, + decommission_start_pool_state, decommission_unresolved_listing_error, dedup_indices, + default_decommission_bucket_concurrency, default_decommission_entry_concurrency, drain_decommission_entry_queue, + enqueue_decommission_entry, ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, + ensure_decommission_generation, ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, + ensure_decommission_start_allowed, ensure_decommission_start_keeps_active_pool, ensure_decommission_start_local_leader, ensure_decommission_start_pool_states, ensure_decommission_start_rebalance_meta_allowed, ensure_decommission_start_target_capacity, ensure_decommission_terminal_operation_supported, ensure_local_decommission_pool_leaders, ensure_valid_decommission_pool_index, first_resumable_decommission_queue_indices, get_by_index, guard_decommission_cancelers, has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested, load_decommission_entry_versions, local_decommission_queue_prefix, mark_decommission_bucket_done, merge_pool_status_refresh, missing_decommission_worker_prefix, - observe_decommission_terminal_reload_result, pool_meta_has_active_decommission, require_decommission_store, - reserve_decommission_start_cancelers, resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state, - resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result, - resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result, - resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result, - resolve_decommission_partial_listing_entry, resolve_decommission_pool_meta_reload_result, - resolve_decommission_preflight_heal_result, resolve_decommission_progress_save_result, - resolve_decommission_terminal_mark_after_error_result, resolve_decommission_terminal_mark_result, - resolve_decommission_update_after_result, resolve_start_decommission_pool_meta_reload_result, - rollback_start_decommission_pool_meta, run_decommission_buckets_bounded, run_decommission_listing_with_retry, - run_decommission_listing_with_retry_and_drain, run_decommission_side_effect, should_cleanup_decommission_source_entry, - should_continue_decommission_queue, should_count_decommission_version_complete, - should_preserve_decommission_canceled_state, should_reject_decommission_cancel_as_terminal, - should_retry_decommission_cancel_reload, should_retry_decommission_listing, should_skip_canceled_decommission_routine, - spawn_decommission_index_cancelers, split_decommission_buckets, take_and_cancel_decommission_canceler, - take_decommission_canceler, track_decommission_current_object, track_decommission_current_object_stage, - update_decommission_for_operation, validate_start_decommission_request, wait_decommission_listing_retry, - wait_decommission_worker_drain, with_decommission_entry_context, + observe_decommission_terminal_reload_result, pool_meta_has_active_decommission, + reconcile_decommission_unresolved_entries_for_completion, record_decommission_unresolved_entry, + require_decommission_store, reserve_decommission_start_cancelers, resolve_decommission_bucket_done_save_result, + resolve_decommission_bucket_state, resolve_decommission_check_after_list_result, + resolve_decommission_entry_cleanup_delete_result, resolve_decommission_entry_exact_versions, + resolve_decommission_entry_reload_result, resolve_decommission_listing_worker_result, + resolve_decommission_optional_bucket_config_result, resolve_decommission_partial_listing_entry, + resolve_decommission_pool_meta_reload_result, resolve_decommission_preflight_heal_result, + resolve_decommission_progress_save_result, resolve_decommission_terminal_mark_after_error_result, + resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result, + resolve_start_decommission_pool_meta_reload_result, rollback_start_decommission_pool_meta, + run_decommission_buckets_bounded, run_decommission_listing_with_retry, run_decommission_listing_with_retry_and_drain, + run_decommission_side_effect, should_cleanup_decommission_source_entry, should_continue_decommission_queue, + should_count_decommission_version_complete, should_preserve_decommission_canceled_state, + should_reject_decommission_cancel_as_terminal, should_retry_decommission_cancel_reload, + should_retry_decommission_listing, should_skip_canceled_decommission_routine, spawn_decommission_index_cancelers, + split_decommission_buckets, take_and_cancel_decommission_canceler, take_decommission_canceler, + track_decommission_current_object, track_decommission_current_object_stage, update_decommission_for_operation, + validate_start_decommission_request, wait_decommission_listing_retry, wait_decommission_worker_drain, + with_decommission_entry_context, }; + use crate::bucket::metadata_sys; use crate::data_movement; - use crate::disk::endpoint::Endpoint; + use crate::disk::{STORAGE_FORMAT_FILE, endpoint::Endpoint}; use crate::error::{Error, StorageError}; use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::runtime::instance::InstanceContext; use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats}; + use crate::storage_api_contracts::bucket::MakeBucketOptions; use crate::store::ECStore; use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo}; use rustfs_filemeta::{MetaCacheEntries, MetadataResolutionParams}; @@ -7622,7 +7814,8 @@ mod pools_tests { #[test] fn test_resolve_decommission_partial_listing_entry_rejects_unresolved_metadata() { - let err = resolve_decommission_partial_listing_entry( + let generation = OffsetDateTime::now_utc(); + let unresolved_entry = resolve_decommission_partial_listing_entry( MetaCacheEntries(vec![None]), MetadataResolutionParams { dir_quorum: 2, @@ -7635,23 +7828,160 @@ mod pools_tests { 1, 2, 3, + generation, ) .expect_err("unresolved partial listing must fail closed"); - let message = err.to_string(); + assert_eq!(unresolved_entry.bucket, "bucket-a"); + assert_eq!(unresolved_entry.object, "prefix/"); + assert_eq!(unresolved_entry.pool_index, 2); + assert_eq!(unresolved_entry.set_index, 3); + assert_eq!(unresolved_entry.source_generation, generation); + assert_eq!(unresolved_entry.candidate_count, 0); + assert_eq!(unresolved_entry.disk_error_count, 1); + assert_eq!(unresolved_entry.reason, "metadata_resolution_failed"); + + let message = decommission_unresolved_listing_error(&unresolved_entry).to_string(); assert!(message.contains("decommission listing could not resolve metadata")); assert!(message.contains("bucket-a/prefix/")); assert!(message.contains("pool 2 set 3")); assert!(message.contains("1 disk error(s)")); } + #[test] + fn unresolved_entry_ledger_blocks_unverified_completion_and_reconciles_after_verified_sweep() { + let generation = OffsetDateTime::now_utc(); + let mut pool_meta = PoolMeta { + version: POOL_META_VERSION, + pools: vec![PoolStatus { + id: 0, + cmd_line: "/data/pool".to_string(), + last_update: generation, + decommission: Some(PoolDecommissionInfo { + start_time: Some(generation), + ..Default::default() + }), + }], + dont_save: true, + }; + let unresolved_entry = DecommissionUnresolvedEntry { + bucket: "bucket-a".to_string(), + object: "object-a".to_string(), + pool_index: 0, + set_index: 0, + source_generation: generation, + candidate_count: 1, + disk_error_count: 0, + observed_at: generation, + reason: "metadata_resolution_failed".to_string(), + }; + + assert!( + record_decommission_unresolved_entry(&mut pool_meta, 0, generation, unresolved_entry) + .expect("active generation should accept unresolved entry") + ); + let err = reconcile_decommission_unresolved_entries_for_completion(&mut pool_meta, 0, None) + .expect_err("unverified completion must retain unresolved entries"); + assert!(err.to_string().contains("1 unresolved listing entries remain")); + assert_eq!( + pool_meta.pools[0] + .decommission + .as_ref() + .expect("decommission should exist") + .unresolved_entries + .len(), + 1 + ); + + reconcile_decommission_unresolved_entries_for_completion(&mut pool_meta, 0, Some(generation)) + .expect("a successful same-generation final sweep should reconcile stale entries"); + assert!( + pool_meta.pools[0] + .decommission + .as_ref() + .expect("decommission should exist") + .unresolved_entries + .is_empty() + ); + } + + #[tokio::test] + async fn final_sweep_persists_unresolved_entry_from_real_listing() { + let (dirs, store) = metadata_sys::test_support::isolated_store_over_temp_disks().await; + let bucket = "decommission-final-sweep-unresolved"; + let object = "corrupt-object"; + store + .peer_sys + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("test bucket should be created"); + metadata_sys::init_bucket_metadata_sys(store.clone(), vec![bucket.to_string()]).await; + + let generation = OffsetDateTime::now_utc(); + { + let mut pool_meta = store.pool_meta.write().await; + pool_meta.dont_save = false; + pool_meta.pools[0].decommission = Some(PoolDecommissionInfo { + start_time: Some(generation), + ..Default::default() + }); + } + + for (disk_index, dir) in dirs.iter().enumerate() { + let object_dir = dir.path().join(bucket).join(object); + tokio::fs::create_dir_all(&object_dir) + .await + .expect("corrupt object directory should be created"); + tokio::fs::write(object_dir.join(STORAGE_FORMAT_FILE), format!("corrupt-xl-meta-{disk_index}").into_bytes()) + .await + .expect("divergent corrupt metadata should be written"); + } + + let err = store + .check_after_decommission(0, generation) + .await + .expect_err("real final sweep must fail closed on unresolved listing metadata"); + assert!(err.to_string().contains("decommission listing could not resolve metadata")); + + let in_memory_entries = store.pool_meta.read().await.pools[0] + .decommission + .as_ref() + .expect("decommission should remain active") + .unresolved_entries + .clone(); + assert_eq!(in_memory_entries.len(), 1); + let unresolved_entry = &in_memory_entries[0]; + assert_eq!(unresolved_entry.bucket, bucket); + assert_eq!(unresolved_entry.object, object); + assert_eq!(unresolved_entry.pool_index, 0); + assert_eq!(unresolved_entry.set_index, 0); + assert_eq!(unresolved_entry.source_generation, generation); + assert_eq!(unresolved_entry.candidate_count, 4); + assert_eq!(unresolved_entry.disk_error_count, 0); + assert_eq!(unresolved_entry.reason, "metadata_resolution_failed"); + + let mut restored = PoolMeta::default(); + restored + .load_for_startup(store.pools[0].clone()) + .await + .expect("persisted pool metadata should reload after the final-sweep failure"); + assert_eq!( + restored.pools[0] + .decommission + .as_ref() + .expect("reloaded decommission should exist") + .unresolved_entries, + in_memory_entries + ); + } + #[tokio::test] async fn test_record_decommission_entry_error_cancels_listing_and_preserves_first_error() { let entry_error = Arc::new(tokio::sync::Mutex::new(None)); let rx = CancellationToken::new(); - record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await; - record_decommission_entry_error(&entry_error, &rx, Error::OperationCanceled).await; + assert!(record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await); + assert!(!record_decommission_entry_error(&entry_error, &rx, Error::OperationCanceled).await); assert!(rx.is_cancelled()); assert!(matches!(*entry_error.lock().await, Some(Error::SlowDown))); @@ -7663,7 +7993,7 @@ mod pools_tests { let rx = CancellationToken::new(); rx.cancel(); - record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await; + assert!(!record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await); assert!(entry_error.lock().await.is_none()); } diff --git a/rustfs/src/app/admin_usecase.rs b/rustfs/src/app/admin_usecase.rs index 877296f29..5660ec2d9 100644 --- a/rustfs/src/app/admin_usecase.rs +++ b/rustfs/src/app/admin_usecase.rs @@ -16,7 +16,8 @@ use super::storage_api::admin_usecase::admin::get_server_info; use super::storage_api::admin_usecase::capacity::{ - PoolDecommissionInfo, PoolStatus, RebalStatus, get_total_usable_capacity, get_total_usable_capacity_free, + DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus, RebalStatus, get_total_usable_capacity, + get_total_usable_capacity_free, }; use super::storage_api::admin_usecase::contract::StorageAdminApi; use super::storage_api::admin_usecase::contract::bucket::{BucketOperations as _, BucketOptions}; @@ -107,6 +108,8 @@ pub struct AdminPoolDecommissionInfo { pub bytes_failed: usize, #[serde(rename = "waitingReason")] pub waiting_reason: Option, + #[serde(rename = "unresolvedEntries", skip_serializing_if = "Vec::is_empty")] + pub unresolved_entries: Vec, } #[derive(Debug, Clone, serde::Serialize)] @@ -619,6 +622,7 @@ impl DefaultAdminUsecase { bytes_done: info.bytes_done, bytes_failed: info.bytes_failed, waiting_reason, + unresolved_entries: info.unresolved_entries, } } @@ -676,7 +680,7 @@ impl DefaultAdminUsecase { #[cfg(test)] mod tests { - use super::super::storage_api::admin_usecase::capacity::{PoolDecommissionInfo, PoolStatus}; + use super::super::storage_api::admin_usecase::capacity::{DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus}; use super::*; use time::OffsetDateTime; use tracing_subscriber::{Layer, Registry, layer::Context, prelude::*}; @@ -987,6 +991,17 @@ mod tests { items_decommission_failed: 1, bytes_done: 1024, bytes_failed: 64, + unresolved_entries: vec![DecommissionUnresolvedEntry { + bucket: "bucket-a".to_string(), + object: "prefix/unresolved.txt".to_string(), + pool_index: 3, + set_index: 1, + source_generation: OffsetDateTime::UNIX_EPOCH, + candidate_count: 2, + disk_error_count: 1, + observed_at: OffsetDateTime::UNIX_EPOCH, + reason: "metadata_resolution_failed".to_string(), + }], ..Default::default() }), }, @@ -1010,6 +1025,13 @@ mod tests { assert_eq!(value["decommissionInfo"]["objectsDecommissionedFailed"], 1); assert_eq!(value["decommissionInfo"]["bytesDecommissioned"], 1024); assert_eq!(value["decommissionInfo"]["bytesDecommissionedFailed"], 64); + assert_eq!(value["decommissionInfo"]["unresolvedEntries"][0]["bucket"], "bucket-a"); + assert_eq!(value["decommissionInfo"]["unresolvedEntries"][0]["object"], "prefix/unresolved.txt"); + assert_eq!( + value["decommissionInfo"]["unresolvedEntries"][0]["sourceGeneration"], + "1970-01-01T00:00:00Z" + ); + assert_eq!(value["decommissionInfo"]["unresolvedEntries"][0]["reason"], "metadata_resolution_failed"); assert_eq!(value["decommissionInfo"]["waitingReason"], "queued"); } diff --git a/rustfs/src/app/storage_api.rs b/rustfs/src/app/storage_api.rs index a599cfde6..0d96e9a3c 100644 --- a/rustfs/src/app/storage_api.rs +++ b/rustfs/src/app/storage_api.rs @@ -31,6 +31,7 @@ pub(crate) mod admin { } pub(crate) mod capacity { + pub(crate) type DecommissionUnresolvedEntry = crate::storage::storage_api::ecstore_capacity::DecommissionUnresolvedEntry; pub(crate) type PoolDecommissionInfo = crate::storage::storage_api::ecstore_capacity::PoolDecommissionInfo; pub(crate) type PoolStatus = crate::storage::storage_api::ecstore_capacity::PoolStatus; pub(crate) type RebalStatus = crate::storage::storage_api::ecstore_rebalance::RebalStatus; diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 512e255ee..75555790c 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -398,7 +398,7 @@ pub(crate) mod ecstore_bucket { pub(crate) mod ecstore_capacity { pub(crate) use rustfs_ecstore::api::capacity::{ - PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free, + DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free, is_reserved_or_invalid_bucket, }; }