fix(ecstore): require durable decommission ledger format (#6871)

This commit is contained in:
cxymds
2026-08-30 08:47:09 +08:00
committed by GitHub
parent 1e8c8d4cd5
commit 21e5b3dc64
7 changed files with 382 additions and 16 deletions
@@ -122,6 +122,14 @@ fn control_plane_failure(op: &str, bucket: Option<&str>, error_code: Option<i32>
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]
+291 -2
View File
@@ -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<DecommissionOperation>,
@@ -2027,6 +2055,7 @@ pub(crate) async fn pause_pool_activation_after_durable_save<S>(pool: &Arc<S>, 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::<Vec<_>>();
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::<Vec<_>>()
.collect::<Vec<_>>();
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,
@@ -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,
@@ -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<crate::store::ECStore>) {
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<RebalanceMeta>,
@@ -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,
}
}
+3
View File
@@ -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 {
+39 -14
View File
@@ -123,6 +123,21 @@ fn verify_node_mutation_body<T: CanonicalMutationBody>(request: &Request<T>, 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]