fix(rebalance): retry fleet proof before start activation (#8049)

Retry the transient missing fleet capability proof in both the admin start flow and node RPC activation. Return a retryable 503 with Retry-After after rollback when readiness does not converge.
This commit is contained in:
cxymds
2026-09-21 19:27:58 +08:00
committed by GitHub
parent cd2fa50e50
commit 8ba3309624
2 changed files with 319 additions and 20 deletions
+222 -16
View File
@@ -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<String>) -> 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<T, F, Fut>(
operation: &'static str,
rebalance_id: &str,
request_id: &str,
actor: &str,
remote_addr: &str,
deadline: Instant,
mut operation_fn: F,
) -> Result<T, StorageError>
where
F: FnMut() -> Fut,
Fut: Future<Output = Result<T, StorageError>>,
{
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<ECStore>,
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<AdminOperation>) -> 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![
+97 -4
View File
@@ -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<F, Fut>(mut operation: F) -> StorageResult<()>
where
F: FnMut() -> Fut,
Fut: Future<Output = StorageResult<()>>,
{
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")));