From ae8d0796586e3bfb112440a94e04f2f09a7f8694 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sun, 13 Sep 2026 10:55:46 +0800 Subject: [PATCH] fix(heal): reject invalid admin selectors before admission (#7724) --- .../src/cluster/rpc/peer_rest_client.rs | 40 ++- rustfs/src/admin/handlers/heal.rs | 88 +++++- rustfs/src/storage/rpc/node_service.rs | 287 ++++++++++++++---- rustfs/src/storage/rpc/node_service/heal.rs | 71 ++++- 4 files changed, 419 insertions(+), 67 deletions(-) diff --git a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs index 7bea7eb0f..34ecaebe6 100644 --- a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs @@ -122,6 +122,13 @@ fn peer_failure_without_details(op: &str, bucket: Option<&str>) -> Error { } } +fn heal_control_status_error(status: tonic::Status) -> Error { + if status.code() == tonic::Code::InvalidArgument { + return Error::InvalidArgument("heal".to_string(), "control".to_string(), status.message().to_owned()); + } + Error::from(status) +} + /// Decode a control-plane response failure. Peers at or above the typed /// `ControlPlaneErrorCode` change (backlog#1845) carry a machine-readable /// discriminant beside the legacy `error_info` string; prefer it, then fall @@ -1831,7 +1838,11 @@ impl PeerRestClient { }); request.set_timeout(rustfs_protos::heal_control_execution_timeout()); set_tonic_canonical_body_digest(&mut request, &canonical_body)?; - let response = client.heal_control(request).await?.into_inner(); + let response = client + .heal_control(request) + .await + .map_err(heal_control_status_error)? + .into_inner(); if !response.success { return Err(Error::other( response @@ -2846,6 +2857,7 @@ mod tests { use super::*; use crate::config::com::STORAGE_CLASS_SUB_SYS; use crate::disk::error::DiskError; + use crate::disk::error_reduce::reduce_errs; use crate::layout::{disks_layout::DisksLayout, endpoints::SetupType}; use rustfs_config::{ENV_KUBERNETES_SERVICE_HOST, ENV_LOCAL_ENDPOINT_HOST, ENV_STARTUP_TOPOLOGY_WAIT_MODE}; @@ -2856,6 +2868,32 @@ mod tests { use temp_env::async_with_vars; use tracing_subscriber::{Registry, fmt::MakeWriter, layer::SubscriberExt}; + #[test] + fn heal_selector_rpc_error_preserves_invalid_argument_without_retry() { + let err = heal_control_status_error(tonic::Status::invalid_argument("heal pool index 99 is out of range")); + assert!(matches!(&err, Error::InvalidArgument(_, _, message) if message == "heal pool index 99 is out of range")); + assert!(heal_control_retry_action(&err, false, false).is_none()); + assert!(!PeerRestClient::is_network_like_error(&err)); + + for code in [ + tonic::Code::Unavailable, + tonic::Code::Internal, + tonic::Code::FailedPrecondition, + ] { + let err = heal_control_status_error(tonic::Status::new(code, "invalid selector")); + assert!(!matches!(err, Error::InvalidArgument(..)), "classify by code, not message"); + let Error::Io(io_error) = &err else { + panic!("non-argument RPC failures must retain their typed transport status: {err:?}"); + }; + assert_eq!( + embedded_tonic_status(io_error) + .expect("transport status should be preserved") + .code(), + code + ); + } + } + #[test] fn control_plane_failure_prefers_typed_not_initialized_code() { use rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode; diff --git a/rustfs/src/admin/handlers/heal.rs b/rustfs/src/admin/handlers/heal.rs index a33eda6ad..8148e4c01 100644 --- a/rustfs/src/admin/handlers/heal.rs +++ b/rustfs/src/admin/handlers/heal.rs @@ -17,12 +17,15 @@ use crate::admin::router::{AdminOperation, Operation, S3Router}; use crate::admin::runtime_sources::app_context_from_req; use crate::admin::storage_api::bucket::is_reserved_or_invalid_bucket; use crate::admin::storage_api::bucket::utils::is_valid_object_prefix; +use crate::admin::storage_api::error::StorageError; +use crate::admin::storage_api::runtime::EndpointServerPools; +use crate::admin::storage_api::s3::{S3ErrorCode, error as admin_error}; use crate::error::ApiError; use crate::server::ADMIN_PREFIX; use crate::server::RemoteAddr; use crate::storage::rpc::node_service::heal::{ HealControlCoordinator, NodeHealStatusSnapshot, capture_node_heal_status, decode_node_heal_status, - decode_node_replacement_recovery_status, heal_control_coordinator, heal_topology_fingerprint, + decode_node_replacement_recovery_status, heal_control_coordinator, heal_topology_fingerprint, validate_heal_selector, }; use bytes::Bytes; use futures_util::future::join_all; @@ -849,6 +852,7 @@ fn cluster_heal_control_unavailable(reason: &str) -> s3s::S3Error { } struct PreparedHealControlRoute { + endpoints: EndpointServerPools, remote_grid_hosts: Vec, fingerprint: String, coordinator_epoch: u64, @@ -873,6 +877,7 @@ fn prepare_heal_control_route(context: &crate::admin::runtime_sources::AppContex .map(|node| node.grid_host) .collect(); Ok(PreparedHealControlRoute { + endpoints, remote_grid_hosts, fingerprint, coordinator_epoch, @@ -942,9 +947,12 @@ async fn route_cluster_heal_control( coordinator_capability_verified: bool, ) -> S3Result { let response = if route.coordinator.is_local { - crate::storage::rpc::node_service::execute_heal_control_envelope(envelope, route.coordinator_epoch) + crate::storage::rpc::node_service::execute_heal_control_envelope(envelope, route.coordinator_epoch, &route.endpoints) .await .map_err(|err| { + if err.code() == tonic::Code::InvalidArgument { + return admin_error(S3ErrorCode::InvalidArgument, err.message().to_owned()); + } warn!( event = EVENT_ADMIN_REQUEST_FAILED, component = LOG_COMPONENT_ADMIN_API, @@ -986,6 +994,9 @@ async fn route_cluster_heal_control( .heal_control(rustfs_protos::HEAL_CONTROL_PROTOCOL_VERSION, route.fingerprint.clone(), command) .await .map_err(|err| { + if let StorageError::InvalidArgument(_, _, message) = &err { + return admin_error(S3ErrorCode::InvalidArgument, message.clone()); + } warn!( event = EVENT_ADMIN_REQUEST_FAILED, component = LOG_COMPONENT_ADMIN_API, @@ -1016,13 +1027,20 @@ async fn route_cluster_heal_control( }) } -async fn execute_after_heal_control_capability(probe: P, execute: E) -> S3Result +async fn execute_after_heal_start_preflight( + endpoints: &EndpointServerPools, + options: &HealOpts, + probe: P, + execute: E, +) -> S3Result where P: FnOnce() -> PF, PF: Future>, E: FnOnce() -> EF, EF: Future>, { + validate_heal_selector(endpoints, options.pool, options.set) + .map_err(|err| admin_error(S3ErrorCode::InvalidArgument, err.to_string()))?; probe().await?; execute().await } @@ -1032,7 +1050,9 @@ async fn submit_cluster_heal_start( hip: &HealInitParams, ) -> S3Result { let route = prepare_heal_control_route(&context)?; - let result = execute_after_heal_control_capability( + let result = execute_after_heal_start_preflight( + &route.endpoints, + &hip.hs, || require_cluster_heal_control_capability(&context, &route), || async { let heal_request = build_heal_channel_request(hip); @@ -1634,7 +1654,7 @@ mod tests { BackgroundHealCoverage, BackgroundHealCoverageReason, BackgroundHealProgress, HealInitParams, HealResp, HealRuntimeState, aggregate_cluster_heal_status, aggregate_replacement_recovery_cluster_status, background_heal_runtime_state, build_heal_channel_request, build_replacement_recovery_status_response, encode_background_heal_status, - encode_heal_control_path, encode_heal_start_success, encode_heal_task_status, execute_after_heal_control_capability, + encode_heal_control_path, encode_heal_start_success, encode_heal_task_status, execute_after_heal_start_preflight, heal_channel_response_items, heal_channel_response_progress, heal_channel_response_summary, heal_control_response_id, json_response, map_heal_response, merge_peer_heal_statuses, peer_topology_complete, query_peer_heal_status, query_peer_replacement_recovery_status, read_cluster_heal_status, reject_heal_admission, validate_heal_request_mode, @@ -1810,10 +1830,58 @@ mod tests { } } + #[tokio::test] + async fn heal_selector_preflight_rejects_before_probe_and_token_creation() { + use crate::admin::storage_api::runtime::PoolEndpoints; + + let endpoints = super::EndpointServerPools::from(vec![PoolEndpoints { + legacy: false, + set_count: 1, + drives_per_set: 4, + endpoints: Default::default(), + cmd_line: String::new(), + platform: String::new(), + }]); + for (pool, set) in [(99, 99), (1, 0), (0, 1)] { + let hip = HealInitParams { + hs: HealOpts { + pool: Some(pool), + set: Some(set), + recursive: true, + ..Default::default() + }, + force_start: true, + ..Default::default() + }; + let probed = AtomicBool::new(false); + let mut token = None; + let error = execute_after_heal_start_preflight( + &endpoints, + &hip.hs, + || async { + probed.store(true, Ordering::SeqCst); + Ok(()) + }, + || async { + token = Some(build_heal_channel_request(&hip).id); + Ok(()) + }, + ) + .await + .expect_err("invalid selector must be rejected synchronously"); + assert_eq!(error.code(), &S3ErrorCode::InvalidArgument); + assert_eq!(error.code().status_code(), Some(StatusCode::BAD_REQUEST)); + assert!(!probed.load(Ordering::SeqCst)); + assert!(token.is_none(), "rejected START must not allocate a client token"); + } + } + #[tokio::test] async fn cluster_capability_gate_runs_before_execution() { let executed = AtomicBool::new(false); - let rejected = execute_after_heal_control_capability( + let rejected = execute_after_heal_start_preflight( + &super::EndpointServerPools::default(), + &HealOpts::default(), || async { Err(super::cluster_heal_control_unavailable("test_capability_failure")) }, || async { executed.store(true, Ordering::SeqCst); @@ -1824,7 +1892,9 @@ mod tests { assert!(rejected.is_err()); assert!(!executed.load(Ordering::SeqCst)); - execute_after_heal_control_capability( + execute_after_heal_start_preflight( + &super::EndpointServerPools::default(), + &HealOpts::default(), || async { Ok(()) }, || async { executed.store(true, Ordering::SeqCst); @@ -1846,7 +1916,9 @@ mod tests { for attempt in 0..3 { let executed_ids = &mut request_ids; let request_params = &hip; - let result = execute_after_heal_control_capability( + let result = execute_after_heal_start_preflight( + &super::EndpointServerPools::default(), + &HealOpts::default(), || async { if attempt < 2 { Err(super::cluster_heal_control_unavailable("test_capability_failure")) diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index 41a383f11..7a9418f17 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -252,7 +252,10 @@ fn heal_control_remaining(expires_at_unix_ms: i64, now_unix_ms: i64) -> Result Result<(), Status> { +fn validate_admin_heal_control_start( + request: &rustfs_heal_contracts::heal_channel::HealChannelRequest, + endpoints: &EndpointServerPools, +) -> Result<(), Status> { if request.source != rustfs_heal_contracts::heal_channel::HealRequestSource::Admin { return Err(Status::permission_denied("heal control start source must be admin")); } @@ -261,9 +264,8 @@ fn validate_admin_heal_control_start(request: &rustfs_heal_contracts::heal_chann "admin heal control start cannot contain automatic replacement endpoints", )); } - if request.pool_index.is_some() != request.set_index.is_some() { - return Err(Status::invalid_argument("heal control start requires both pool and set")); - } + heal::validate_heal_selector(endpoints, request.pool_index, request.set_index) + .map_err(|err| Status::invalid_argument(err.to_string()))?; if request.bucket.is_empty() { if request.object_prefix.as_deref().is_some_and(|prefix| !prefix.is_empty()) { return Err(Status::invalid_argument("root heal control start cannot contain an object prefix")); @@ -275,6 +277,14 @@ fn validate_admin_heal_control_start(request: &rustfs_heal_contracts::heal_chann if request.disk.is_some() != erasure_set_target { return Err(Status::invalid_argument("root erasure-set heal control target is inconsistent")); } + if let Some(disk) = &request.disk { + let selector = rustfs_heal::heal::utils::normalize_set_disk_id(disk) + .and_then(|normalized| rustfs_heal::heal::utils::parse_set_disk_id(&normalized).ok()) + .ok_or_else(|| Status::invalid_argument("invalid root erasure-set heal control target"))?; + if (Some(selector.0), Some(selector.1)) != (request.pool_index, request.set_index) { + return Err(Status::invalid_argument("root erasure-set heal control target is inconsistent")); + } + } } else if request.disk.is_some() { return Err(Status::invalid_argument( "bucket heal control start cannot contain an erasure-set disk target", @@ -763,14 +773,16 @@ async fn initialize_heal_topology_fingerprint_with_probe( pub(crate) async fn execute_heal_control_envelope( envelope: rustfs_protos::heal_control::Envelope, expected_coordinator_epoch: u64, + endpoints: &EndpointServerPools, ) -> Result, Status> { - execute_heal_control_envelope_with_manager(envelope, expected_coordinator_epoch, None).await + execute_heal_control_envelope_with_manager(envelope, expected_coordinator_epoch, None, endpoints).await } async fn execute_heal_control_envelope_with_manager( envelope: rustfs_protos::heal_control::Envelope, expected_coordinator_epoch: u64, manager: Option>, + endpoints: &EndpointServerPools, ) -> Result, Status> { let now = heal_control_now_unix_ms()?; envelope @@ -781,6 +793,12 @@ async fn execute_heal_control_envelope_with_manager( let command_digest = Sha256::digest(&canonical_envelope).into(); let (request_id, coordinator_epoch, command) = envelope.into_execution().map_err(Status::invalid_argument)?; + // Reject invalid targets before replay admission or forceStart can mutate + // task ownership. Both local and forwarded starts use this boundary. + if let rustfs_protos::heal_control::ExecutableCommand::Start { request } = &command { + validate_admin_heal_control_start(request, endpoints)?; + } + let replay_cache = HEAL_CONTROL_REPLAY_CACHE.get_or_init(|| tokio::sync::Mutex::new(HashMap::new())); let replay_entry = { let mut replay_cache = timeout(heal_control_remaining(expires_at_unix_ms, now)?, replay_cache.lock()) @@ -799,9 +817,6 @@ async fn execute_heal_control_envelope_with_manager( return Ok(cached.clone()); } - if let rustfs_protos::heal_control::ExecutableCommand::Start { request } = &command { - validate_admin_heal_control_start(request)?; - } let retain_completed_result = !matches!(&command, rustfs_protos::heal_control::ExecutableCommand::Query { .. }); let manager = manager @@ -1197,7 +1212,7 @@ impl heal_control_service_server::HealControlService for HealControlRpcService { rustfs_protos::heal_control::decode_envelope(&request.get_ref().command).map_err(Status::invalid_argument)?; let coordinator_epoch = rustfs_protos::heal_control_coordinator_epoch(fingerprint).map_err(Status::failed_precondition)?; - let result = execute_heal_control_envelope(envelope, coordinator_epoch).await?; + let result = execute_heal_control_envelope(envelope, coordinator_epoch, &endpoints).await?; let canonical_response = rustfs_protos::canonical_heal_control_response_body( request.get_ref().version, &request.get_ref().topology_fingerprint, @@ -3090,7 +3105,7 @@ mod tests { request.recursive = Some(true); request.heal_endpoints = vec!["/mnt/replacement".to_string()]; - let err = validate_admin_heal_control_start(&request) + let err = validate_admin_heal_control_start(&request, &heal_control_test_endpoints_with_coordinator("node-d", true)) .expect_err("admin heal-control must not accept automatic replacement targets"); assert_eq!(err.code(), tonic::Code::InvalidArgument); } @@ -3153,9 +3168,13 @@ mod tests { } let envelope = rustfs_protos::heal_control::decode_envelope(&command).map_err(Status::invalid_argument)?; - let result = - execute_heal_control_envelope_with_manager(envelope, self.coordinator_epoch, Some(Arc::clone(&self.manager))) - .await?; + let result = execute_heal_control_envelope_with_manager( + envelope, + self.coordinator_epoch, + Some(Arc::clone(&self.manager)), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await?; if matches!(self.fault, HealControlTransportFault::DropAfterAdmission) { return Err(Status::unavailable("transport failed after heal admission")); } @@ -3280,21 +3299,126 @@ mod tests { .expect("heal-control start response should carry a matching canonical receipt") } + #[tokio::test] + async fn heal_selector_rejection_preserves_existing_task_and_replay_state() { + let (manager, request, metadata) = heal_start_retry_fixture(); + let endpoints = heal_control_test_endpoints_with_coordinator("node-d", true); + let original_id = request.id.clone(); + let original = rustfs_protos::heal_control::Envelope::start(request.clone(), metadata).expect("valid original start"); + execute_heal_control_envelope_with_manager(original, metadata.coordinator_epoch, Some(manager.clone()), &endpoints) + .await + .expect("original heal should be admitted"); + for (pool, set) in [(99, 99), (1, 0), (0, 2)] { + for root in [false, true] { + let mut invalid = request.clone(); + invalid.id = Uuid::new_v4().to_string(); + invalid.pool_index = Some(pool); + invalid.set_index = Some(set); + if root { + invalid.bucket.clear(); + invalid.object_prefix = None; + invalid.disk = Some(rustfs_heal::heal::utils::format_set_disk_id(pool, set)); + } + let invalid_id = invalid.id.clone(); + let envelope = rustfs_protos::heal_control::Envelope::start(invalid, metadata).expect("invalid selector encodes"); + let error = execute_heal_control_envelope_with_manager( + envelope, + metadata.coordinator_epoch, + Some(manager.clone()), + &endpoints, + ) + .await + .expect_err("invalid forceStart must be rejected before cancelling or admitting work"); + assert_eq!(error.code(), tonic::Code::InvalidArgument); + assert_eq!(manager.operations_snapshot().await.queue_length, 1); + manager + .get_task_status(&original_id) + .await + .expect("original task must retain ownership"); + assert!(matches!( + manager.get_task_status(&invalid_id).await, + Err(rustfs_heal::Error::TaskNotFound { .. }) + )); + let cache = super::HEAL_CONTROL_REPLAY_CACHE + .get() + .expect("original request initialized replay cache") + .lock() + .await; + assert!(!cache.contains_key(&invalid_id), "rejected request must not consume replay admission"); + } + } + } + + #[test] + fn heal_selector_rpc_target_must_match_validated_selector() { + let (_, mut request, _) = heal_start_retry_fixture(); + let endpoints = heal_control_test_endpoints_with_coordinator("node-d", true); + request.bucket.clear(); + request.object_prefix = None; + request.pool_index = Some(0); + request.set_index = Some(0); + for disk in ["pool_99_set_99", "pool_0_set_1", "0_1", "invalid"] { + request.disk = Some(disk.to_string()); + let err = validate_admin_heal_control_start(&request, &endpoints).expect_err("RPC target must match selector"); + assert_eq!(err.code(), tonic::Code::InvalidArgument); + } + for disk in ["pool_0_set_0", "0_0"] { + request.disk = Some(disk.to_string()); + validate_admin_heal_control_start(&request, &endpoints).expect("matching canonical or compact selector is valid"); + } + } + + #[tokio::test] + async fn heal_selector_transport_returns_invalid_argument_without_admission() { + let (manager, mut request, metadata) = heal_start_retry_fixture(); + request.pool_index = Some(99); + request.set_index = Some(99); + let request_id = request.id.clone(); + let command = encode_transport_start(request, metadata); + let fingerprint = "selector-validation-transport"; + let mut client = connect_faulty_heal_control_client( + manager.clone(), + fingerprint, + metadata.coordinator_epoch, + HealControlTransportFault::None, + ) + .await + .expect("loopback listener is required for selector transport regression"); + let error = call_heal_control_transport(&mut client, fingerprint, command) + .await + .expect_err("peer must reject selector"); + assert_eq!(error.code(), tonic::Code::InvalidArgument); + assert_eq!(manager.operations_snapshot().await.queue_length, 0); + assert!(matches!( + manager.get_task_status(&request_id).await, + Err(rustfs_heal::Error::TaskNotFound { .. }) + )); + } + #[tokio::test] async fn heal_start_retry_exact_forced_envelope_returns_cached_admission() { let (manager, request, metadata) = heal_start_retry_fixture(); let request_id = request.id.clone(); let envelope = rustfs_protos::heal_control::Envelope::start(request, metadata).expect("valid forced start"); - let lost_response = - execute_heal_control_envelope_with_manager(envelope.clone(), metadata.coordinator_epoch, Some(manager.clone())) - .await - .expect("first request is admitted before its response is lost"); + let lost_response = execute_heal_control_envelope_with_manager( + envelope.clone(), + metadata.coordinator_epoch, + Some(manager.clone()), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .expect("first request is admitted before its response is lost"); assert_eq!(manager.operations_snapshot().await.queue_length, 1); // The caller sees no first response, but retries the original envelope. - let replayed = execute_heal_control_envelope_with_manager(envelope, metadata.coordinator_epoch, Some(manager.clone())) - .await - .expect("an exact envelope replay must recover its receipt"); + let replayed = execute_heal_control_envelope_with_manager( + envelope, + metadata.coordinator_epoch, + Some(manager.clone()), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .expect("an exact envelope replay must recover its receipt"); assert_eq!(replayed, lost_response); assert_eq!( manager.operations_snapshot().await.queue_length, @@ -3395,9 +3519,14 @@ mod tests { let (manager, request, metadata) = heal_start_retry_fixture(); let first_id = request.id.clone(); let first = rustfs_protos::heal_control::Envelope::start(request.clone(), metadata).expect("first start"); - let _lost_response = execute_heal_control_envelope_with_manager(first, metadata.coordinator_epoch, Some(manager.clone())) - .await - .expect("first admission"); + let _lost_response = execute_heal_control_envelope_with_manager( + first, + metadata.coordinator_epoch, + Some(manager.clone()), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .expect("first admission"); // A fresh HTTP forceStart request intentionally requests another start. let mut next_request = request; @@ -3408,9 +3537,14 @@ mod tests { ..metadata }; let next = rustfs_protos::heal_control::Envelope::start(next_request, next_metadata).expect("new forced start"); - let response = execute_heal_control_envelope_with_manager(next, metadata.coordinator_epoch, Some(manager.clone())) - .await - .expect("forceStart preserves its explicit admission semantics"); + let response = execute_heal_control_envelope_with_manager( + next, + metadata.coordinator_epoch, + Some(manager.clone()), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .expect("forceStart preserves its explicit admission semantics"); let outcome = rustfs_protos::heal_control::decode_result(&response) .and_then(|result| result.into_outcome(&next_id, metadata.coordinator_epoch)) .expect("new receipt"); @@ -3428,10 +3562,14 @@ mod tests { async fn heal_start_retry_same_id_with_changed_envelope_conflicts_before_admission() { let (manager, request, metadata) = heal_start_retry_fixture(); let original = rustfs_protos::heal_control::Envelope::start(request.clone(), metadata).expect("original start"); - let receipt = - execute_heal_control_envelope_with_manager(original.clone(), metadata.coordinator_epoch, Some(manager.clone())) - .await - .expect("original admission"); + let receipt = execute_heal_control_envelope_with_manager( + original.clone(), + metadata.coordinator_epoch, + Some(manager.clone()), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .expect("original admission"); let mut changed_options = request.clone(); changed_options.remove_corrupted = Some(true); let changed_metadata = rustfs_protos::heal_control::RequestMetadata { @@ -3442,16 +3580,26 @@ mod tests { rustfs_protos::heal_control::Envelope::start(changed_options, metadata).expect("changed options"), rustfs_protos::heal_control::Envelope::start(request, changed_metadata).expect("changed nonce"), ] { - let error = execute_heal_control_envelope_with_manager(changed, metadata.coordinator_epoch, Some(manager.clone())) - .await - .expect_err("one request ID cannot identify different envelope bytes"); + let error = execute_heal_control_envelope_with_manager( + changed, + metadata.coordinator_epoch, + Some(manager.clone()), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .expect_err("one request ID cannot identify different envelope bytes"); assert_eq!(error.code(), tonic::Code::AlreadyExists); assert_eq!(manager.operations_snapshot().await.queue_length, 1); } assert_eq!( - execute_heal_control_envelope_with_manager(original, metadata.coordinator_epoch, Some(manager)) - .await - .expect("conflicts must preserve the original receipt"), + execute_heal_control_envelope_with_manager( + original, + metadata.coordinator_epoch, + Some(manager), + &heal_control_test_endpoints_with_coordinator("node-d", true) + ) + .await + .expect("conflicts must preserve the original receipt"), receipt ); } @@ -3461,9 +3609,14 @@ mod tests { let (manager, request, metadata) = heal_start_retry_fixture(); let request_id = request.id.clone(); let envelope = rustfs_protos::heal_control::Envelope::start(request, metadata).expect("start envelope"); - let error = execute_heal_control_envelope_with_manager(envelope, metadata.coordinator_epoch + 1, Some(manager.clone())) - .await - .expect_err("a different coordinator epoch cannot accept the request"); + let error = execute_heal_control_envelope_with_manager( + envelope, + metadata.coordinator_epoch + 1, + Some(manager.clone()), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .expect_err("a different coordinator epoch cannot accept the request"); assert_eq!(error.code(), tonic::Code::FailedPrecondition); assert_eq!(manager.operations_snapshot().await.queue_length, 0); assert!(matches!( @@ -3763,9 +3916,14 @@ mod tests { let canonical_token = uuid::Uuid::new_v4().to_string(); let first = rustfs_protos::heal_control::Envelope::start(start(canonical_token.clone()), metadata()).unwrap(); - let first_result = execute_heal_control_envelope_with_manager(first, coordinator_epoch, Some(Arc::clone(&manager))) - .await - .unwrap(); + let first_result = execute_heal_control_envelope_with_manager( + first, + coordinator_epoch, + Some(Arc::clone(&manager)), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .unwrap(); let first_outcome = rustfs_protos::heal_control::decode_result(&first_result) .and_then(|result| result.into_outcome(&canonical_token, coordinator_epoch)) .unwrap(); @@ -3779,10 +3937,14 @@ mod tests { let duplicate_id = uuid::Uuid::new_v4().to_string(); let duplicate = rustfs_protos::heal_control::Envelope::start(start(duplicate_id.clone()), metadata()).unwrap(); - let duplicate_result = - execute_heal_control_envelope_with_manager(duplicate, coordinator_epoch, Some(Arc::clone(&manager))) - .await - .unwrap(); + let duplicate_result = execute_heal_control_envelope_with_manager( + duplicate, + coordinator_epoch, + Some(Arc::clone(&manager)), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .unwrap(); let duplicate_outcome = rustfs_protos::heal_control::decode_result(&duplicate_result) .and_then(|result| result.into_outcome(&duplicate_id, coordinator_epoch)) .unwrap(); @@ -3803,9 +3965,14 @@ mod tests { None, ) .unwrap(); - let query_result = execute_heal_control_envelope_with_manager(query, coordinator_epoch, Some(Arc::clone(&manager))) - .await - .unwrap(); + let query_result = execute_heal_control_envelope_with_manager( + query, + coordinator_epoch, + Some(Arc::clone(&manager)), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .unwrap(); let query_outcome = rustfs_protos::heal_control::decode_result(&query_result) .and_then(|result| result.into_outcome(&query_id, coordinator_epoch)) .unwrap(); @@ -3822,9 +3989,14 @@ mod tests { canonical_token.clone(), ) .unwrap(); - let cancel_result = execute_heal_control_envelope_with_manager(cancel, coordinator_epoch, Some(Arc::clone(&manager))) - .await - .unwrap(); + let cancel_result = execute_heal_control_envelope_with_manager( + cancel, + coordinator_epoch, + Some(Arc::clone(&manager)), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .unwrap(); let cancel_outcome = rustfs_protos::heal_control::decode_result(&cancel_result) .and_then(|result| result.into_outcome(&cancel_id, coordinator_epoch)) .unwrap(); @@ -3842,9 +4014,14 @@ mod tests { None, ) .unwrap(); - let stopped_result = execute_heal_control_envelope_with_manager(stopped_query, coordinator_epoch, Some(manager)) - .await - .unwrap(); + let stopped_result = execute_heal_control_envelope_with_manager( + stopped_query, + coordinator_epoch, + Some(manager), + &heal_control_test_endpoints_with_coordinator("node-d", true), + ) + .await + .unwrap(); let stopped_outcome = rustfs_protos::heal_control::decode_result(&stopped_result) .and_then(|result| result.into_outcome(&stopped_query_id, coordinator_epoch)) .unwrap(); diff --git a/rustfs/src/storage/rpc/node_service/heal.rs b/rustfs/src/storage/rpc/node_service/heal.rs index fe9637453..96a2306b1 100644 --- a/rustfs/src/storage/rpc/node_service/heal.rs +++ b/rustfs/src/storage/rpc/node_service/heal.rs @@ -63,6 +63,36 @@ pub(crate) struct HealControlCoordinator { pub is_local: bool, } +#[derive(Debug, PartialEq, Eq, thiserror::Error)] +pub(crate) enum HealSelectorError { + #[error("heal start requires both pool and set")] + Incomplete, + #[error("heal pool index {pool} is out of range")] + InvalidPool { pool: usize }, + #[error("heal set index {set} is out of range for pool {pool}")] + InvalidSet { pool: usize, set: usize }, +} + +/// Validate configured positions independently of disk health: an offline set +/// remains a valid target for repair. +pub(crate) fn validate_heal_selector( + endpoints: &EndpointServerPools, + pool: Option, + set: Option, +) -> Result<(), HealSelectorError> { + match (pool, set) { + (None, None) => Ok(()), + (Some(pool), Some(set)) => { + let selected_pool = endpoints.as_ref().get(pool).ok_or(HealSelectorError::InvalidPool { pool })?; + if set >= selected_pool.set_count { + return Err(HealSelectorError::InvalidSet { pool, set }); + } + Ok(()) + } + _ => Err(HealSelectorError::Incomplete), + } +} + pub(crate) fn heal_control_coordinator(endpoint_pools: &EndpointServerPools) -> Result { endpoint_pools .get_nodes() @@ -469,10 +499,10 @@ pub(crate) fn decode_node_replacement_recovery_status(data: &[u8]) -> Result