diff --git a/rustfs/src/admin/handlers/rebalance.rs b/rustfs/src/admin/handlers/rebalance.rs index 81eb853de..557228264 100644 --- a/rustfs/src/admin/handlers/rebalance.rs +++ b/rustfs/src/admin/handlers/rebalance.rs @@ -43,13 +43,21 @@ use s3s::{ s3_error, }; use serde::{Deserialize, Serialize}; -use std::{sync::Arc, time::Duration}; +use std::{ + future::Future, + sync::Arc, + time::{Duration, Instant}, +}; use time::OffsetDateTime; use tracing::{error, info, warn}; const LOG_COMPONENT_ADMIN: &str = "admin"; const LOG_SUBSYSTEM_REBALANCE: &str = "rebalance"; const EVENT_ADMIN_REBALANCE_STATE: &str = "admin_rebalance_state"; +const REBALANCE_START_FLEET_PROOF_MARKER: &str = "pool activation requires a live fleet capability proof"; +const REBALANCE_START_FLEET_PROOF_RETRY_BUDGET: Duration = Duration::from_secs(30); +const REBALANCE_START_FLEET_PROOF_RETRY_DELAY: Duration = Duration::from_secs(5); +const REBALANCE_START_RETRY_AFTER_SECS: &str = "5"; fn admin_request_id(headers: &HeaderMap) -> Option<&str> { headers @@ -105,6 +113,80 @@ fn rebalance_internal_error(message: impl Into) -> S3Error { S3Error::with_message(S3ErrorCode::InternalError, message.into()) } +fn rebalance_fleet_proof_unavailable_error() -> S3Error { + let mut err = S3Error::with_message( + S3ErrorCode::ServiceUnavailable, + "rebalance start is waiting for cluster capability readiness; retry later", + ); + let mut headers = HeaderMap::new(); + headers.insert(http::header::RETRY_AFTER, HeaderValue::from_static(REBALANCE_START_RETRY_AFTER_SECS)); + err.set_headers(headers); + err +} + +fn is_rebalance_start_fleet_proof_retryable(err: &StorageError) -> bool { + crate::storage_api::capacity::is_pool_activation_fleet_proof_error(err) + && err.to_string().contains(REBALANCE_START_FLEET_PROOF_MARKER) +} + +fn rebalance_start_failure_error(start_err: &StorageError, rollback_result: &Result<(), String>) -> S3Error { + if crate::storage_api::capacity::is_pool_activation_fleet_proof_error(start_err) { + if rollback_result.is_ok() { + return rebalance_fleet_proof_unavailable_error(); + } + return rebalance_internal_error("failed to roll back rebalance start after fleet capability readiness failure"); + } + + rebalance_internal_error(rebalance_start_rollback_error(&start_err.to_string(), rollback_result)) +} + +async fn retry_rebalance_start_operation( + operation: &'static str, + rebalance_id: &str, + request_id: &str, + actor: &str, + remote_addr: &str, + deadline: Instant, + mut operation_fn: F, +) -> Result +where + F: FnMut() -> Fut, + Fut: Future>, +{ + let mut attempt = 1_u32; + loop { + match operation_fn().await { + Ok(value) => return Ok(value), + Err(err) if is_rebalance_start_fleet_proof_retryable(&err) => { + let now = Instant::now(); + if now >= deadline { + return Err(err); + } + let delay = REBALANCE_START_FLEET_PROOF_RETRY_DELAY.min(deadline.saturating_duration_since(now)); + warn!( + event = EVENT_ADMIN_REBALANCE_STATE, + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_REBALANCE, + action = "start", + state = "fleet_proof_retry_scheduled", + result = "retrying", + operation, + attempt, + retry_delay_ms = delay.as_millis(), + request_id = %request_id, + actor = %actor, + remote_addr = %remote_addr, + rebalance_id = %rebalance_id, + "admin rebalance state" + ); + tokio::time::sleep(delay).await; + attempt = attempt.saturating_add(1); + } + Err(err) => return Err(err), + } + } +} + fn rebalance_rollback_stop_failure_message(rebalance_id: &str, failures: &[String]) -> String { format!("cluster stop_rebalance rollback for {rebalance_id} partial: {}", failures.join("; ")) } @@ -209,7 +291,7 @@ async fn rollback_rebalance_start_for_admin( store: &Arc, notification_sys: Option<&NotificationSys>, rebalance_id: &str, - start_err: &str, + start_err: &StorageError, request_id: &str, actor: &str, remote_addr: &str, @@ -246,7 +328,7 @@ async fn rollback_rebalance_start_for_admin( ), } - Err(rebalance_internal_error(rebalance_start_rollback_error(start_err, &rollback_result))) + Err(rebalance_start_failure_error(start_err, &rollback_result)) } pub fn register_rebalance_route(r: &mut S3Router) -> std::io::Result<()> { @@ -562,6 +644,7 @@ impl Operation for RebalanceStart { "admin rebalance state" ); let notification_sys = current_notification_system(); + let fleet_proof_deadline = Instant::now() + REBALANCE_START_FLEET_PROOF_RETRY_BUDGET; for step in rebalance_start_steps(notification_sys.is_some()) { match step { RebalanceStartStep::PropagateFence => { @@ -592,13 +675,11 @@ impl Operation for RebalanceStart { error = %err, "admin rebalance state" ); - - let start_err = err.to_string(); rollback_rebalance_start_for_admin( &store, Some(notification_sys), &id, - &start_err, + &err, &request_id, &actor, &remote_addr, @@ -608,7 +689,17 @@ impl Operation for RebalanceStart { } } RebalanceStartStep::StartLocal => { - if let Err(err) = store.start_rebalance_for_id(&id).await { + let start_result = retry_rebalance_start_operation( + "start_rebalance_for_id", + &id, + &request_id, + &actor, + &remote_addr, + fleet_proof_deadline, + || store.start_rebalance_for_id(&id), + ) + .await; + if let Err(err) = start_result { error!( event = EVENT_ADMIN_REBALANCE_STATE, component = LOG_COMPONENT_ADMIN, @@ -682,6 +773,9 @@ impl Operation for RebalanceStart { ))); } } + if crate::storage_api::capacity::is_pool_activation_fleet_proof_error(&err) { + return Err(rebalance_fleet_proof_unavailable_error()); + } return Err(rebalance_internal_error(format!( "failed to start rebalance after metadata initialized for {id}; local metadata was finalized as failed: {start_err}" ))); @@ -701,7 +795,17 @@ impl Operation for RebalanceStart { rebalance_id = %id, "admin rebalance state" ); - if let Err(err) = notification_sys.load_rebalance_meta(true).await { + let worker_result = retry_rebalance_start_operation( + "load_rebalance_meta(start=true)", + &id, + &request_id, + &actor, + &remote_addr, + fleet_proof_deadline, + || notification_sys.load_rebalance_meta(true), + ) + .await; + if let Err(err) = worker_result { error!( event = EVENT_ADMIN_REBALANCE_STATE, component = LOG_COMPONENT_ADMIN, @@ -715,13 +819,11 @@ impl Operation for RebalanceStart { error = %err, "admin rebalance state" ); - - let start_err = err.to_string(); rollback_rebalance_start_for_admin( &store, Some(notification_sys), &id, - &start_err, + &err, &request_id, &actor, &remote_addr, @@ -1059,12 +1161,14 @@ mod rebalance_handler_tests { use super::calculate_rebalance_progress; use super::{ Body, HeaderMap, Method, Operation, Params, RebalPoolProgress, RebalanceAdminStatus, RebalancePoolStatus, RebalanceStart, - RebalanceStartStep, RebalanceStatus, RebalanceStop, RebalanceStopPropagationStatus, S3ErrorCode, S3Request, Uri, - build_rebalance_admin_status, build_rebalance_pool_statuses, build_rebalance_stop_propagation_status, - rebalance_pool_used, rebalance_query_present, rebalance_remaining_buckets, rebalance_rollback_failure_message, - rebalance_rollback_stop_failure_message, rebalance_start_rollback_error, rebalance_start_steps, rebalance_stop_target_id, - rebalance_used_pct, rollback_result_label, stop_rebalance_admission_first, + RebalanceStartStep, RebalanceStatus, RebalanceStop, RebalanceStopPropagationStatus, S3ErrorCode, S3Request, StatusCode, + Uri, build_rebalance_admin_status, build_rebalance_pool_statuses, build_rebalance_stop_propagation_status, + is_rebalance_start_fleet_proof_retryable, rebalance_pool_used, rebalance_query_present, rebalance_remaining_buckets, + rebalance_rollback_failure_message, rebalance_rollback_stop_failure_message, rebalance_start_failure_error, + rebalance_start_rollback_error, rebalance_start_steps, rebalance_stop_target_id, rebalance_used_pct, + retry_rebalance_start_operation, rollback_result_label, stop_rebalance_admission_first, }; + use crate::admin::storage_api::error::StorageError; use crate::admin::storage_api::rebalance::{ DiskStat, RebalSaveOpt, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo, RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, encode_rebalance_stop_propagation_record, @@ -1328,6 +1432,108 @@ mod rebalance_handler_tests { assert!(message.contains("rollback error: peer b stop_rebalance failed: timeout")); } + #[test] + fn test_rebalance_start_fleet_proof_maps_to_retryable_503() { + let start_err = StorageError::other(super::REBALANCE_START_FLEET_PROOF_MARKER); + + let err = rebalance_start_failure_error(&start_err, &Ok(())); + + assert_eq!(err.code(), &S3ErrorCode::ServiceUnavailable); + assert_eq!(err.status_code(), Some(StatusCode::SERVICE_UNAVAILABLE)); + assert_eq!( + err.headers() + .and_then(|headers| headers.get(http::header::RETRY_AFTER)) + .and_then(|value| value.to_str().ok()), + Some(super::REBALANCE_START_RETRY_AFTER_SECS) + ); + assert!( + !err.message() + .expect("retryable rebalance start error should have a message") + .contains(super::REBALANCE_START_FLEET_PROOF_MARKER) + ); + } + + #[test] + fn test_rebalance_start_fleet_proof_rollback_failure_is_generic_internal_error() { + let start_err = StorageError::other(super::REBALANCE_START_FLEET_PROOF_MARKER); + let rollback_result = Err("peer rollback failed".to_string()); + + let err = rebalance_start_failure_error(&start_err, &rollback_result); + + assert_eq!(err.code(), &S3ErrorCode::InternalError); + assert!( + !err.message() + .expect("internal rebalance start error should have a message") + .contains(super::REBALANCE_START_FLEET_PROOF_MARKER) + ); + } + + #[test] + fn test_rebalance_start_unrelated_failure_remains_internal_error() { + let start_err = StorageError::other("disk read failed"); + + let err = rebalance_start_failure_error(&start_err, &Ok(())); + + assert_eq!(err.code(), &S3ErrorCode::InternalError); + assert!(err.message().is_some_and(|message| message.contains("disk read failed"))); + } + + #[test] + fn test_rebalance_start_fleet_proof_retry_only_matches_missing_live_proof() { + let missing = StorageError::other(super::REBALANCE_START_FLEET_PROOF_MARKER); + let expired = StorageError::other(format!( + "rebalance meta save failed during start_rebalance: {}", + "pool activation fleet capability proof expired before commit" + )); + + assert!(is_rebalance_start_fleet_proof_retryable(&missing)); + assert!(!is_rebalance_start_fleet_proof_retryable(&expired)); + } + + #[tokio::test(start_paused = true)] + async fn test_rebalance_start_retry_waits_for_fleet_proof() { + use std::sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }; + + let attempts = Arc::new(AtomicUsize::new(0)); + let task = tokio::spawn({ + let attempts = Arc::clone(&attempts); + async move { + retry_rebalance_start_operation( + "load_rebalance_meta(start=true)", + "rebalance-id", + "request-id", + "actor", + "remote", + std::time::Instant::now() + super::REBALANCE_START_FLEET_PROOF_RETRY_BUDGET, + move || { + let attempts = Arc::clone(&attempts); + async move { + let attempt = attempts.fetch_add(1, Ordering::SeqCst); + if attempt == 0 { + Err(StorageError::other(super::REBALANCE_START_FLEET_PROOF_MARKER)) + } else { + Ok(()) + } + } + }, + ) + .await + } + }); + + tokio::task::yield_now().await; + assert_eq!(attempts.load(Ordering::SeqCst), 1); + + tokio::time::advance(super::REBALANCE_START_FLEET_PROOF_RETRY_DELAY).await; + task.await + .expect("retry task should not panic") + .expect("fleet proof retry should eventually succeed"); + assert_eq!(attempts.load(Ordering::SeqCst), 2); + } + #[test] fn test_rebalance_rollback_stop_failure_message_lists_failures() { let failures = vec![ diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index 72247fc08..45e11792a 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -52,9 +52,11 @@ use serde::Deserialize; use sha2::{Digest, Sha256}; use std::{ collections::HashMap, + future::Future, io::Cursor, pin::Pin, sync::{Arc, LazyLock, OnceLock}, + time::Instant, }; use time::OffsetDateTime; use tokio::spawn; @@ -78,6 +80,9 @@ const EVENT_RPC_REQUEST_FAILED: &str = "rpc_request_failed"; const EVENT_RPC_RESPONSE_EMITTED: &str = "rpc_response_emitted"; const EVENT_RPC_BACKGROUND_TASK_SPAWNED: &str = "rpc_background_task_spawned"; const EVENT_RPC_BACKGROUND_TASK_FAILED: &str = "rpc_background_task_failed"; +const REBALANCE_START_FLEET_PROOF_MARKER: &str = "pool activation requires a live fleet capability proof"; +const REBALANCE_START_FLEET_PROOF_RETRY_BUDGET: Duration = Duration::from_secs(30); +const REBALANCE_START_FLEET_PROOF_RETRY_DELAY: Duration = Duration::from_secs(5); const HEAL_CONTROL_REPLAY_CACHE_MAX_ENTRIES: usize = 4096; const TIER_MUTATION_PEER_STATE_UNSPECIFIED_WIRE: i32 = 0; const TIER_MUTATION_PEER_STATE_PREPARED_WIRE: i32 = 1; @@ -501,6 +506,48 @@ fn background_rebalance_start_error_message(result: StorageResult<()>) -> Option result.err().map(|err| format!("start_rebalance failed: {err}")) } +fn is_rebalance_start_fleet_proof_retryable(err: &Error) -> bool { + crate::storage::storage_api::ecstore_capacity::is_pool_activation_fleet_proof_error(err) + && err.to_string().contains(REBALANCE_START_FLEET_PROOF_MARKER) +} + +async fn retry_rebalance_start_fleet_proof(mut operation: F) -> StorageResult<()> +where + F: FnMut() -> Fut, + Fut: Future>, +{ + let deadline = Instant::now() + REBALANCE_START_FLEET_PROOF_RETRY_BUDGET; + let mut attempt = 1_u32; + + loop { + match operation().await { + Ok(()) => return Ok(()), + Err(err) if is_rebalance_start_fleet_proof_retryable(&err) => { + let now = Instant::now(); + if now >= deadline { + return Err(err); + } + let delay = REBALANCE_START_FLEET_PROOF_RETRY_DELAY.min(deadline.saturating_duration_since(now)); + warn!( + event = EVENT_RPC_BACKGROUND_TASK_FAILED, + component = LOG_COMPONENT_STORAGE, + subsystem = LOG_SUBSYSTEM_REBALANCE, + operation = "start_rebalance", + state = "fleet_proof_retry_scheduled", + result = "retrying", + attempt, + retry_delay_ms = delay.as_millis(), + error = %err, + "node rpc background task retry" + ); + tokio::time::sleep(delay).await; + attempt = attempt.saturating_add(1); + } + Err(err) => return Err(err), + } + } +} + fn stop_rebalance_response(result: StorageResult<()>) -> StopRebalanceResponse { match result { Ok(_) => StopRebalanceResponse { @@ -2800,7 +2847,9 @@ impl Node for NodeService { if start_rebalance { log_background_rebalance_task_spawned!(start_rebalance); - if let Some(message) = background_rebalance_start_error_message(store.start_rebalance().await) { + if let Some(message) = + background_rebalance_start_error_message(retry_rebalance_start_fleet_proof(|| store.start_rebalance()).await) + { error!( event = EVENT_RPC_BACKGROUND_TASK_FAILED, component = LOG_COMPONENT_STORAGE, @@ -2975,9 +3024,10 @@ mod tests { SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION, SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION, SCANNER_PUBLICATION_LEASE_TTL_MS, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, STORAGE_CLASS_SUB_SYS, admit_heal_control_replay, background_rebalance_start_error_message, execute_heal_control_envelope_with_manager, - 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, + initialize_heal_topology_fingerprint, initialize_heal_topology_fingerprint_with_probe, + is_rebalance_start_fleet_proof_retryable, 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, retry_rebalance_start_fleet_proof, scanner_activity_response_v7, start_decommission_failure_response, stop_rebalance_response, validate_admin_heal_control_start, }; @@ -8121,6 +8171,49 @@ mod tests { assert!(message.contains("boom")); } + #[test] + fn test_rebalance_start_retry_ignores_expired_fleet_proof() { + let expired = Error::other("pool activation fleet capability proof expired before commit"); + + assert!(!is_rebalance_start_fleet_proof_retryable(&expired)); + } + + #[tokio::test(start_paused = true)] + async fn test_retry_rebalance_start_waits_for_fleet_proof() { + use std::sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }; + + let attempts = Arc::new(AtomicUsize::new(0)); + let task = tokio::spawn({ + let attempts = Arc::clone(&attempts); + async move { + retry_rebalance_start_fleet_proof(move || { + let attempts = Arc::clone(&attempts); + async move { + let attempt = attempts.fetch_add(1, Ordering::SeqCst); + if attempt == 0 { + Err(Error::other(super::REBALANCE_START_FLEET_PROOF_MARKER)) + } else { + Ok(()) + } + } + }) + .await + } + }); + + tokio::task::yield_now().await; + assert_eq!(attempts.load(Ordering::SeqCst), 1); + + tokio::time::advance(super::REBALANCE_START_FLEET_PROOF_RETRY_DELAY).await; + task.await + .expect("retry task should not panic") + .expect("fleet proof retry should eventually succeed"); + assert_eq!(attempts.load(Ordering::SeqCst), 2); + } + #[test] fn test_stop_rebalance_response_reports_local_stop_error() { let response = stop_rebalance_response(Err(Error::other("boom")));