diff --git a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs index c566bd0a9..c03a69e6e 100644 --- a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs @@ -122,6 +122,14 @@ fn control_plane_failure(op: &str, bucket: Option<&str>, error_code: Option if error_code == Some(rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode::ControlPlaneErrorNotInitialized as i32) { return Error::RemoteNotInitialized; } + if error_code == Some(rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode::ControlPlaneErrorInvalidArgument as i32) + { + return Error::InvalidArgument( + "control-plane".to_string(), + op.to_string(), + error_info.unwrap_or_else(|| format!("{op}: peer rejected invalid argument without details")), + ); + } match error_info { Some(msg) => Error::other(msg), None => peer_failure_without_details(op, bucket), @@ -2335,6 +2343,29 @@ mod tests { use rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode; assert_eq!(ControlPlaneErrorCode::ControlPlaneErrorUnspecified as i32, 0); assert_eq!(ControlPlaneErrorCode::ControlPlaneErrorNotInitialized as i32, 1); + assert_eq!(ControlPlaneErrorCode::ControlPlaneErrorInvalidArgument as i32, 2); + } + + #[test] + fn control_plane_failure_preserves_typed_invalid_argument_reason() { + use rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode; + + let reason = "durable unresolved-entry recovery requires pool metadata V2 or V3"; + let err = control_plane_failure( + "start_decommission", + None, + Some(ControlPlaneErrorCode::ControlPlaneErrorInvalidArgument as i32), + Some(reason.to_string()), + ); + + assert!( + matches!( + err, + Error::InvalidArgument(ref scope, ref operation, ref actual_reason) + if scope == "control-plane" && operation == "start_decommission" && actual_reason == reason + ), + "forwarded validation failures must remain typed and actionable" + ); } #[test] diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index f25b42f82..e3c1194d8 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -176,6 +176,34 @@ fn pool_meta_v3_writer_enabled() -> bool { ) } +fn ensure_decommission_ledger_persistence_supported_for( + version: u16, + v2_writer_enabled: bool, + v3_writer_enabled: bool, +) -> Result<()> { + if matches!(version, POOL_META_VERSION | POOL_META_GENERATION_VERSION) || v2_writer_enabled || v3_writer_enabled { + return Ok(()); + } + + Err(Error::InvalidArgument( + "decommission".to_string(), + "pool-metadata-version".to_string(), + format!( + "durable unresolved-entry recovery requires pool metadata V2 or V3; enable both {} and {} only after every reader and writer supports V2", + rustfs_config::ENV_POOL_META_V2_WRITE, + rustfs_config::ENV_POOL_META_V2_FLEET_CONFIRMED, + ), + )) +} + +fn ensure_decommission_ledger_persistence_supported(pool_meta: &PoolMeta) -> Result<()> { + ensure_decommission_ledger_persistence_supported_for( + pool_meta.version, + pool_meta_v2_writer_enabled(), + pool_meta_v3_writer_enabled(), + ) +} + #[derive(Clone, Debug)] pub struct DecommissionCanceler { operation: Arc, @@ -2027,6 +2055,7 @@ pub(crate) async fn pause_pool_activation_after_durable_save(pool: &Arc, f #[cfg(test)] struct PoolActivationStartProbeState { kind: PoolActivationStartKind, + preflight_side_effect_attempted: std::sync::atomic::AtomicBool, attempted: std::sync::atomic::AtomicBool, notify: tokio::sync::Notify, } @@ -2045,6 +2074,7 @@ impl PoolActivationStartProbe { pub(crate) fn install(kind: PoolActivationStartKind) -> Self { let state = Arc::new(PoolActivationStartProbeState { kind, + preflight_side_effect_attempted: std::sync::atomic::AtomicBool::new(false), attempted: std::sync::atomic::AtomicBool::new(false), notify: tokio::sync::Notify::new(), }); @@ -2061,6 +2091,14 @@ impl PoolActivationStartProbe { self.state.notify.notified().await; } } + + pub(crate) fn preflight_side_effect_was_attempted(&self) -> bool { + self.state.preflight_side_effect_attempted.load(Ordering::Acquire) + } + + pub(crate) fn activation_was_attempted(&self) -> bool { + self.state.attempted.load(Ordering::Acquire) + } } #[cfg(test)] @@ -2090,6 +2128,21 @@ pub(crate) fn observe_pool_activation_start_attempt(kind: PoolActivationStartKin } } +#[cfg(test)] +fn observe_pool_activation_preflight_side_effect_attempt(kind: PoolActivationStartKind) { + let probes = POOL_ACTIVATION_START_PROBES + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("pool activation start probe should not be poisoned") + .iter() + .filter(|state| state.kind == kind) + .cloned() + .collect::>(); + for state in probes { + state.preflight_side_effect_attempted.store(true, Ordering::Release); + } +} + fn rollback_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: &PoolMeta, indices: &[usize]) { publish_pool_meta_updates(pool_meta, previous_pool_meta, indices); } @@ -6189,6 +6242,7 @@ impl ECStore { ) -> Result<()> { { let mut pool_meta = self.pool_meta.write().await; + ensure_decommission_ledger_persistence_supported(&pool_meta)?; record_decommission_unresolved_entry(&mut pool_meta, idx, generation, entry)?; } self.save_current_pool_meta(&[idx]) @@ -6339,6 +6393,7 @@ impl ECStore { } ensure_decommission_start_pool_states(&latest_pool_meta, indices)?; + ensure_decommission_ledger_persistence_supported(&latest_pool_meta)?; let previous_pool_meta = latest_pool_meta.clone(); let first_idx = indices.first().copied(); @@ -7040,10 +7095,14 @@ impl ECStore { save_guard.ensure_write_safe("decommission cannot be scheduled while pool metadata requires recovery")?; let indices = { let pool_meta = self.pool_meta.read().await; - resumable_decommission_queue_indices(&pool_meta) + let indices = resumable_decommission_queue_indices(&pool_meta) .into_iter() .filter(|idx| indices.contains(idx)) - .collect::>() + .collect::>(); + if !indices.is_empty() { + ensure_decommission_ledger_persistence_supported(&pool_meta)?; + } + indices }; if indices.is_empty() { return Ok(Vec::new()); @@ -9299,8 +9358,11 @@ impl ECStore { { let pool_meta = self.pool_meta.read().await; ensure_decommission_start_pool_states(&pool_meta, &indices)?; + ensure_decommission_ledger_persistence_supported(&pool_meta)?; } + #[cfg(test)] + observe_pool_activation_preflight_side_effect_attempt(PoolActivationStartKind::Decommission); let decom_buckets = self.get_buckets_to_decommission().await?; let mut healed_buckets = HashSet::with_capacity(decom_buckets.len()); @@ -10856,10 +10918,98 @@ mod tests { assert!(!is_pool_activation_fleet_proof_error(&Error::ConfigNotFound)); } + #[tokio::test] + #[serial_test::serial] + async fn decommission_v1_start_preflights_reject_before_metadata_writes() { + let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(None).await; + let baseline = load_pool_meta_replicas(store.pools.clone(), true) + .await + .expect("baseline pool metadata should be readable"); + assert_eq!(baseline.meta.version, POOL_META_V1_VERSION); + *store.pool_meta.write().await = baseline.meta.clone(); + let start_probe = PoolActivationStartProbe::install(PoolActivationStartKind::Decommission); + let err = store + .start_decommission(vec![0]) + .await + .expect_err("the initial V1 start preflight must reject before side effects"); + assert!(matches!(err, Error::InvalidArgument(..))); + assert!( + !start_probe.preflight_side_effect_was_attempted(), + "V1 rejection must precede bucket listing, healing, and metadata-bucket creation" + ); + assert!( + !start_probe.activation_was_attempted(), + "V1 rejection must not enter the authoritative activation save" + ); + let after_early_rejection = load_pool_meta_replicas(store.pools.clone(), true) + .await + .expect("early rejection must leave durable pool metadata readable"); + assert_eq!(after_early_rejection.canonical, baseline.canonical); + assert!( + store + .pool_meta + .read() + .await + .pools + .iter() + .all(|pool| pool.decommission.is_none()) + ); + assert!(store.decommission_cancelers.read().await.iter().all(Option::is_none)); + store + .ensure_pool_meta_side_effects_safe("V1 start preflight") + .await + .expect("a deterministic start rejection must not latch recovery"); + drop(start_probe); + + let err = store + .save_current_pool_meta_for_decommission_start( + &[0], + vec![( + 0, + PoolSpaceInfo { + free: 50, + total: 100, + used: 50, + }, + )], + Vec::new(), + ) + .await + .expect_err("the authoritative V1 start preflight must reject before saving"); + assert!(matches!(err, Error::InvalidArgument(..))); + + let after_authoritative_rejection = load_pool_meta_replicas(store.pools.clone(), true) + .await + .expect("authoritative rejection must leave durable pool metadata readable"); + assert_eq!(after_authoritative_rejection.canonical, baseline.canonical); + assert!( + after_authoritative_rejection + .meta + .pools + .iter() + .all(|pool| pool.decommission.is_none()) + ); + assert!( + store + .pool_meta + .read() + .await + .pools + .iter() + .all(|pool| pool.decommission.is_none()) + ); + assert!(store.decommission_cancelers.read().await.iter().all(Option::is_none)); + store + .ensure_pool_meta_side_effects_safe("authoritative V1 start preflight") + .await + .expect("an authoritative capability rejection must not latch recovery"); + } + #[tokio::test] #[serial_test::serial] async fn decommission_activation_fence_loss_after_durable_save_blocks_publication() { let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(None).await; + crate::services::rebalance::promote_test_pool_meta_to_v2(&store).await; let barrier = PoolActivationDurableSaveBarrier::install(&store.pools[0]); let start_store = Arc::clone(&store); let start_task = tokio::spawn(async move { @@ -10914,6 +11064,7 @@ mod tests { #[serial_test::serial] async fn decommission_activation_adopts_canonical_commit_after_replica_failure() { let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(None).await; + crate::services::rebalance::promote_test_pool_meta_to_v2(&store).await; let barrier = PoolActivationDurableSaveBarrier::install(&store.pools[0]); let start_store = Arc::clone(&store); let start_task = tokio::spawn(async move { @@ -11226,6 +11377,35 @@ mod tests { assert!(pool_meta_v3_writer_enabled_for(true, true)); } + #[test] + fn decommission_ledger_persistence_requires_an_observed_or_confirmed_format() { + for (version, v2_enabled, v3_enabled, expected) in [ + (POOL_META_V1_VERSION, false, false, false), + (POOL_META_V1_VERSION, true, false, true), + (POOL_META_V1_VERSION, false, true, true), + (POOL_META_VERSION, false, false, true), + (super::POOL_META_GENERATION_VERSION, false, false, true), + ] { + let result = super::ensure_decommission_ledger_persistence_supported_for(version, v2_enabled, v3_enabled); + assert_eq!( + result.is_ok(), + expected, + "unexpected capability result for pool metadata version {version}" + ); + } + + let half_confirmed_v2 = pool_meta_v2_writer_enabled_for(true, false); + let half_confirmed_v3 = pool_meta_v3_writer_enabled_for(false, true); + let err = super::ensure_decommission_ledger_persistence_supported_for( + POOL_META_V1_VERSION, + half_confirmed_v2, + half_confirmed_v3, + ) + .expect_err("half-enabled rollout gates must not admit decommission"); + assert!(matches!(err, Error::InvalidArgument(..))); + assert!(err.to_string().contains("durable unresolved-entry recovery")); + } + #[test] fn pool_meta_stale_write_rejection_metric_is_countable() { let recorder = metrics_util::debugging::DebuggingRecorder::new(); @@ -14494,6 +14674,113 @@ mod pools_tests { assert!(store.decommission_cancelers.read().await[0].is_none()); } + #[tokio::test] + async fn test_v1_unresolved_ledger_rejection_keeps_live_state_and_write_gate_safe() { + let generation = OffsetDateTime::UNIX_EPOCH; + let status = decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(generation), + ..Default::default() + }), + ); + let last_update = status.last_update; + let store = decommission_worker_test_store( + PoolMeta { + version: POOL_META_V1_VERSION, + pools: vec![status], + ..Default::default() + }, + vec![None], + ); + let entry = DecommissionUnresolvedEntry { + bucket: "bucket-a".to_string(), + object: "directory/".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(), + }; + + let err = store + .persist_decommission_unresolved_entry(0, generation, entry) + .await + .expect_err("V1 must reject the ledger before changing live state"); + + assert!(matches!(err, Error::InvalidArgument(..))); + let pool_meta = store.pool_meta.read().await; + let status = &pool_meta.pools[0]; + assert_eq!(status.last_update, last_update); + assert!( + status + .decommission + .as_ref() + .expect("active decommission metadata should remain present") + .unresolved_entries + .is_empty() + ); + drop(pool_meta); + store + .pool_meta_save_gate + .lock() + .await + .ensure_write_safe("V1 unresolved-entry preflight") + .expect("a deterministic capability rejection must not latch recovery"); + } + + #[tokio::test] + async fn test_v1_runtime_recovery_rejects_worker_but_keeps_cancel_persistable() { + let generation = OffsetDateTime::UNIX_EPOCH; + let store = decommission_worker_test_store( + PoolMeta { + version: POOL_META_V1_VERSION, + pools: vec![decommission_test_pool_status( + 0, + Some(PoolDecommissionInfo { + start_time: Some(generation), + ..Default::default() + }), + )], + ..Default::default() + }, + vec![None], + ); + + let err = store + .reserve_decommission_routines(&CancellationToken::new(), &[0]) + .await + .err() + .expect("V1 recovery must not install a worker that cannot persist an unresolved ledger"); + assert!(matches!(err, Error::InvalidArgument(..))); + assert!(store.decommission_cancelers.read().await[0].is_none()); + + let save_called = Arc::new(AtomicBool::new(false)); + store + .decommission_cancel_with_owner_and_save(0, None, { + let save_called = save_called.clone(); + move |snapshot, _| async move { + snapshot.encode_config_data_for_v2_gate(false)?; + save_called.store(true, Ordering::SeqCst); + Ok(()) + } + }) + .await + .expect("a rejected V1 recovery must remain cancelable without restart"); + + assert!(save_called.load(Ordering::SeqCst)); + let pool_meta = store.pool_meta.read().await; + let info = pool_meta.pools[0] + .decommission + .as_ref() + .expect("cancel metadata should remain present"); + assert!(info.canceled); + assert!(!info.failed); + assert!(!info.complete); + } + #[tokio::test] async fn test_decommission_transition_waits_without_registered_canceler() { let store = decommission_worker_test_store(PoolMeta::default(), vec![None]); @@ -17487,6 +17774,7 @@ mod pools_tests { #[tokio::test] async fn test_runtime_recovery_reserves_the_startup_resumable_queue() { let meta = PoolMeta { + version: super::POOL_META_VERSION, pools: vec![ decommission_test_pool_status( 0, @@ -17544,6 +17832,7 @@ mod pools_tests { #[tokio::test] async fn test_runtime_recovery_does_not_reserve_behind_active_predecessor() { let meta = PoolMeta { + version: super::POOL_META_VERSION, pools: vec![ decommission_test_pool_status( 0, diff --git a/crates/ecstore/src/services/rebalance/control.rs b/crates/ecstore/src/services/rebalance/control.rs index aeb382803..731a3299b 100644 --- a/crates/ecstore/src/services/rebalance/control.rs +++ b/crates/ecstore/src/services/rebalance/control.rs @@ -1667,6 +1667,8 @@ mod tests { async fn assert_real_activation_start_race(paused_kind: PoolActivationStartKind) { let (_temp_dirs, rebalance_store, decommission_store) = crate::services::rebalance::test_two_pool_stores_with_isolated_node_contexts(None).await; + crate::services::rebalance::promote_test_pool_meta_to_v2(&rebalance_store).await; + crate::services::rebalance::promote_test_pool_meta_to_v2(&decommission_store).await; let disk_stats = vec![ DiskStat { total_space: 100, diff --git a/crates/ecstore/src/services/rebalance/mod.rs b/crates/ecstore/src/services/rebalance/mod.rs index d7ffeb7cd..e36699500 100644 --- a/crates/ecstore/src/services/rebalance/mod.rs +++ b/crates/ecstore/src/services/rebalance/mod.rs @@ -111,6 +111,17 @@ pub(crate) async fn test_two_pool_stores_with_isolated_node_contexts( test_two_pool_stores_with_contexts(rebalance_meta, true).await } +#[cfg(test)] +pub(crate) async fn promote_test_pool_meta_to_v2(store: &std::sync::Arc) { + let mut pool_meta = store.pool_meta.read().await.clone(); + pool_meta.version = crate::core::pools::POOL_META_VERSION; + pool_meta + .save(store.pools.clone()) + .await + .expect("test pool metadata should be promoted to V2"); + *store.pool_meta.write().await = pool_meta; +} + #[cfg(test)] async fn test_two_pool_stores_with_contexts( rebalance_meta: Option, diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index ceaaa7a95..4843e7bf9 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -1553,6 +1553,9 @@ pub enum ControlPlaneErrorCode { /// The peer answered but its storage/IAM layer is not initialized yet /// (legacy string form: "errServerNotInitialized"). ControlPlaneErrorNotInitialized = 1, + /// The peer rejected a control-plane request before changing durable state. + /// error_info carries the actionable validation reason. + ControlPlaneErrorInvalidArgument = 2, } impl ControlPlaneErrorCode { /// String value of the enum field names used in the ProtoBuf definition. @@ -1563,6 +1566,7 @@ impl ControlPlaneErrorCode { match self { Self::ControlPlaneErrorUnspecified => "CONTROL_PLANE_ERROR_UNSPECIFIED", Self::ControlPlaneErrorNotInitialized => "CONTROL_PLANE_ERROR_NOT_INITIALIZED", + Self::ControlPlaneErrorInvalidArgument => "CONTROL_PLANE_ERROR_INVALID_ARGUMENT", } } /// Creates an enum from field names used in the ProtoBuf definition. @@ -1570,6 +1574,7 @@ impl ControlPlaneErrorCode { match value { "CONTROL_PLANE_ERROR_UNSPECIFIED" => Some(Self::ControlPlaneErrorUnspecified), "CONTROL_PLANE_ERROR_NOT_INITIALIZED" => Some(Self::ControlPlaneErrorNotInitialized), + "CONTROL_PLANE_ERROR_INVALID_ARGUMENT" => Some(Self::ControlPlaneErrorInvalidArgument), _ => None, } } diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index a51ab51dc..7347e17ca 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -30,6 +30,9 @@ enum ControlPlaneErrorCode { // The peer answered but its storage/IAM layer is not initialized yet // (legacy string form: "errServerNotInitialized"). CONTROL_PLANE_ERROR_NOT_INITIALIZED = 1; + // The peer rejected a control-plane request before changing durable state. + // error_info carries the actionable validation reason. + CONTROL_PLANE_ERROR_INVALID_ARGUMENT = 2; } message PingRequest { diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index b9dfac111..c6b7a740f 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -123,6 +123,21 @@ fn verify_node_mutation_body(request: &Request, ope .map_err(|err| Status::permission_denied(format!("{operation} authentication failed: {err}"))) } +fn start_decommission_failure_response(err: Error) -> StartDecommissionResponse { + match err { + Error::InvalidArgument(_, _, reason) => StartDecommissionResponse { + success: false, + error_info: Some(reason), + error_code: Some(ControlPlaneErrorCode::ControlPlaneErrorInvalidArgument as i32), + }, + err => StartDecommissionResponse { + success: false, + error_info: Some(err.to_string()), + error_code: None, + }, + } +} + fn supports_dynamic_config_rpc(sub_system: &str) -> bool { NOTIFY_SUB_SYSTEMS.contains(&sub_system) || matches!( @@ -2334,11 +2349,7 @@ impl Node for NodeService { success: true, error_info: None, })), - Err(err) => Ok(Response::new(StartDecommissionResponse { - error_code: None, - success: false, - error_info: Some(err.to_string()), - })), + Err(err) => Ok(Response::new(start_decommission_failure_response(err))), } } @@ -2456,7 +2467,7 @@ mod tests { initialize_heal_topology_fingerprint, initialize_heal_topology_fingerprint_with_probe, legacy_scanner_activity_response, make_heal_control_server, make_heal_control_server_with_cache, make_server, make_server_for_context, make_tier_mutation_control_server_for_context, previous_scanner_activity_response, remove_heal_control_replay, - scanner_activity_response_v7, stop_rebalance_response, + scanner_activity_response_v7, start_decommission_failure_response, stop_rebalance_response, }; use crate::storage::rpc::node_service::heal::heal_topology_fingerprint; use crate::storage::storage_api::rpc_consumer::node_service::{DiskError, HealBucketInfo}; @@ -2479,14 +2490,14 @@ mod tests { use rustfs_protos::models::PingBodyBuilder; use rustfs_protos::proto_gen::node_service::{ BackgroundHealStatusRequest, BatchGenerallyLockRequest, CancelDecommissionRequest, CheckPartsRequest, - ClearDecommissionRequest, DeleteBucketMetadataRequest, DeleteBucketRequest, DeletePathsRequest, DeletePolicyRequest, - DeleteRequest, DeleteServiceAccountRequest, DeleteUserRequest, DeleteVersionRequest, DeleteVersionsRequest, - DeleteVolumeRequest, DiskInfoRequest, DownloadProfileDataRequest, GenerallyLockRequest, GetAllBucketStatsRequest, - GetBucketInfoRequest, GetBucketStatsDataRequest, GetCpusRequest, GetMemInfoRequest, GetMetacacheListingRequest, - GetMetricsRequest, GetNetInfoRequest, GetOsInfoRequest, GetPartitionsRequest, GetProcInfoRequest, GetSeLinuxInfoRequest, - GetSrMetricsDataRequest, GetSysConfigRequest, GetSysErrorsRequest, HealBucketRequest, HealControlRequest, - ListBucketRequest, ListDirRequest, ListVolumesRequest, LoadBucketMetadataRequest, LoadGroupRequest, - LoadPolicyMappingRequest, LoadPolicyRequest, LoadRebalanceMetaRequest, LoadServiceAccountRequest, + ClearDecommissionRequest, ControlPlaneErrorCode, DeleteBucketMetadataRequest, DeleteBucketRequest, DeletePathsRequest, + DeletePolicyRequest, DeleteRequest, DeleteServiceAccountRequest, DeleteUserRequest, DeleteVersionRequest, + DeleteVersionsRequest, DeleteVolumeRequest, DiskInfoRequest, DownloadProfileDataRequest, GenerallyLockRequest, + GetAllBucketStatsRequest, GetBucketInfoRequest, GetBucketStatsDataRequest, GetCpusRequest, GetMemInfoRequest, + GetMetacacheListingRequest, GetMetricsRequest, GetNetInfoRequest, GetOsInfoRequest, GetPartitionsRequest, + GetProcInfoRequest, GetSeLinuxInfoRequest, GetSrMetricsDataRequest, GetSysConfigRequest, GetSysErrorsRequest, + HealBucketRequest, HealControlRequest, ListBucketRequest, ListDirRequest, ListVolumesRequest, LoadBucketMetadataRequest, + LoadGroupRequest, LoadPolicyMappingRequest, LoadPolicyRequest, LoadRebalanceMetaRequest, LoadServiceAccountRequest, LoadTransitionTierConfigRequest, LoadUserRequest, LocalStorageInfoRequest, MakeBucketRequest, MakeVolumeRequest, MakeVolumesRequest, Mss, PingRequest, PreparePartTransactionRequest, ReadAllRequest, ReadAtRequest, ReadMultipleRequest, ReadVersionRequest, ReadXlRequest, ReloadPoolMetaRequest, ReloadSiteReplicationConfigRequest, RenameDataRequest, @@ -2576,6 +2587,20 @@ mod tests { } } + #[test] + fn start_decommission_failure_response_preserves_invalid_argument_reason() { + let reason = "durable unresolved-entry recovery requires pool metadata V2 or V3"; + let response = start_decommission_failure_response(Error::InvalidArgument( + "decommission".to_string(), + "pool-metadata-version".to_string(), + reason.to_string(), + )); + + assert!(!response.success); + assert_eq!(response.error_info.as_deref(), Some(reason)); + assert_eq!(response.error_code, Some(ControlPlaneErrorCode::ControlPlaneErrorInvalidArgument as i32)); + } + struct HealControlMockStorage; #[async_trait::async_trait]