fix(ecstore): recover remote disks after transient stalls

Merge the approved fix from pull request #8149.
This commit is contained in:
Chris
2026-09-28 18:40:30 +08:00
committed by GitHub
parent cc7d5e3f5f
commit 151103a609
6 changed files with 221 additions and 26 deletions
+99 -6
View File
@@ -1244,13 +1244,14 @@ impl RemoteDisk {
{
return;
}
// Own the flag before spawning: shutdown can drop the task without polling it.
let lease = RecoveryMonitorLease {
active: Arc::clone(&active),
};
let span = Self::recovery_monitor_span(&addr, &endpoint, handle_id);
super::spawn_background_monitor(span, async move {
#[cfg(test)]
test_state.start_count.fetch_add(1, Ordering::AcqRel);
let lease = RecoveryMonitorLease {
active: Arc::clone(&active),
};
Self::monitor_remote_disk_recovery(addr.clone(), endpoint.clone(), Arc::clone(&health), cancel_token.clone()).await;
#[cfg(test)]
if let Some(hook) = test_state.teardown_hook.lock().await.take() {
@@ -1556,7 +1557,8 @@ impl RemoteDisk {
};
if evict_cached_connection {
evict_failed_connection(addr).await;
// Cache contention must not hold the recovery lease indefinitely.
let _ = timeout(get_drive_active_check_timeout(), evict_failed_connection(addr)).await;
}
result
@@ -1715,7 +1717,7 @@ impl RemoteDisk {
if timeout_duration == Duration::ZERO {
let operation_result = operation().await;
if operation_result.is_ok() {
self.health.log_success();
self.health.record_operation_success(&self.endpoint, "operation_success");
}
self.handle_network_like_error(op, timeout_duration, &operation_result, failure_health_action)
.await;
@@ -1729,7 +1731,7 @@ impl RemoteDisk {
Ok(operation_result) => {
// Log success; the waiting guard balances every exit path.
if operation_result.is_ok() {
self.health.log_success();
self.health.record_operation_success(&self.endpoint, "operation_success");
}
self.handle_network_like_error(op, timeout_duration, &operation_result, failure_health_action)
.await;
@@ -8204,6 +8206,97 @@ mod tests {
assert!(result.is_ok());
}
#[tokio::test]
async fn successful_operation_recovers_suspect_remote_disk_without_monitor() {
let endpoint = Endpoint {
url: url::Url::parse("http://remote-recovery:9000/data").expect("valid endpoint"),
is_local: false,
pool_idx: 0,
set_idx: 0,
disk_idx: 0,
};
let disk = RemoteDisk::new(
&endpoint,
&DiskOption {
cleanup: false,
health_check: false,
},
Arc::new(TcpHttpInternodeDataTransport),
)
.await
.expect("create remote disk");
for duration in [Duration::ZERO, Duration::from_secs(1)] {
disk.health.mark_failure(&endpoint, "read_operation_deadline");
assert_eq!(disk.runtime_state(), RuntimeDriveHealthState::Suspect);
assert!(!disk.recovery_monitor_is_active());
let result = disk
.execute_with_timeout(|| async { Err::<(), _>(DiskError::FileNotFound) }, duration)
.await;
assert!(matches!(result, Err(DiskError::FileNotFound)));
assert_eq!(
disk.runtime_state(),
RuntimeDriveHealthState::Suspect,
"a failed operation must not restore readiness"
);
disk.execute_with_timeout(|| async { Ok(()) }, duration)
.await
.expect("successful disk RPC");
assert_eq!(
disk.runtime_state(),
RuntimeDriveHealthState::Online,
"successful traffic must restore readiness without a recovery monitor"
);
assert!(disk.offline_duration_secs().is_none());
assert_eq!(disk.health.waiting.load(Ordering::Acquire), 0);
}
disk.health.mark_offline(&endpoint, "test_offline");
let result = disk.execute_with_timeout(|| async { Ok(()) }, Duration::from_secs(1)).await;
assert!(
matches!(result, Err(DiskError::FaultyDisk)),
"offline handles still require recovery probes before data I/O"
);
assert_eq!(disk.runtime_state(), RuntimeDriveHealthState::Offline);
}
#[test]
fn recovery_monitor_releases_lease_when_dropped_before_first_poll() {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("create test runtime");
let endpoint = Endpoint {
url: url::Url::parse("http://remote-unpolled:9000/data").expect("valid endpoint"),
is_local: false,
pool_idx: 0,
set_idx: 0,
disk_idx: 0,
};
let active = Arc::new(AtomicBool::new(false));
let start_count = Arc::new(AtomicU32::new(0));
{
let _guard = runtime.enter();
RemoteDisk::schedule_recovery_monitor(
"http://remote-unpolled:9000".to_string(),
endpoint,
Uuid::new_v4(),
Arc::new(DiskHealthTracker::new()),
CancellationToken::new(),
Arc::clone(&active),
RecoveryMonitorTestState {
start_count: Arc::clone(&start_count),
teardown_hook: Arc::new(tokio::sync::Mutex::new(None)),
},
);
assert!(active.load(Ordering::Acquire));
}
drop(runtime);
assert_eq!(start_count.load(Ordering::Acquire), 0, "task must never be polled");
assert!(!active.load(Ordering::Acquire), "dropping an unpolled monitor must release its lease");
}
#[tokio::test]
async fn test_execute_with_timeout_marks_remote_disk_faulty() {
let url = url::Url::parse("http://remote-timeout:9000").expect("operation should succeed");
+11
View File
@@ -1236,6 +1236,17 @@ impl DiskHealthTracker {
record_drive_recovery_class(classify_drive_recovery(duration));
}
self.offline_since_unix_secs.store(0, Ordering::Release);
info!(
event = EVENT_DISK_RECOVERY_PROBE_STATE,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_DISK,
endpoint = %endpoint,
state = "recovered",
previous_state = current.as_str(),
runtime_state = next.as_str(),
reason,
"Disk recovered"
);
} else if let Some(duration) = self.offline_duration() {
record_drive_offline_duration(endpoint, duration);
}
+2
View File
@@ -69,6 +69,8 @@ The existing `details.storage.ready` boolean and `connected` / `disconnected` st
Node readiness additionally reports `details.storage.readQuorum`, `details.storage.writeQuorum`, and `details.poolMetadata.ready`. `details.storage.ready` follows read quorum. `details.lock.ready` on this probe follows shared-lock quorum; on `/minio/health/cluster` it follows exclusive-lock quorum. The metadata component's status is `writable` or `unavailable`. A blocked metadata writer is visible there and on the cluster-write probe; it does not by itself make `/health/ready` return 503.
`details.storage.unavailableDrives` lists inventoried handles excluded from the node's quorum snapshot. Each entry has zero-based `poolIndex`, `setIndex`, and `diskIndex`, its `runtimeState`, and `hostOnline` from the lock reachability observation (always true for local drives). A `suspect` drive with `hostOnline: true` is reachable but not counted; an `online` drive with `hostOnline: false` is excluded because its host did not answer. Internal addresses and filesystem paths are omitted. This list does not enumerate missing inventory slots, and an empty list does not override a failed inventory or quorum check. It is omitted with the other details in minimal responses; HEAD responses have no body. Admin storage info still reports each owner's disk view, which can differ from these node-local handle states.
Node storage quorum uses configured drives per set, all configured pools/sets, and their Standard storage-class data/parity layout. Missing, duplicate, unreachable, or unhealthy disk observations cannot supply extra quorum votes. The read quorum is the data-drive count; the write quorum is that count plus one when data and parity counts are equal. These are observations of available storage slots, not guarantees that a particular object's metadata, shards, or required locks are available.
For a healthy IAM and metadata writer in a four-node, one-drive-per-node EC 2+2 set, after startup has published `FullReady`:
+34 -1
View File
@@ -355,11 +355,12 @@ pub(crate) fn build_health_response_parts(
kms_ready,
include_dependency_details,
});
if let Some(details) = readiness_report.and_then(|report| report.storage_details)
if let Some(details) = readiness_report.and_then(|report| report.storage_details.as_ref())
&& payload.get("details").is_some()
{
payload["details"]["storage"]["readQuorum"] = json!(details.read_quorum_ready);
payload["details"]["storage"]["writeQuorum"] = json!(details.write_quorum_ready);
payload["details"]["storage"]["unavailableDrives"] = json!(details.unavailable_drives);
payload["details"]["poolMetadata"] = json!({
"ready": details.pool_metadata_write_ready,
"status": if details.pool_metadata_write_ready { "writable" } else { "unavailable" },
@@ -478,6 +479,7 @@ mod tests {
read_quorum_ready: read_quorum,
write_quorum_ready: write_quorum,
pool_metadata_write_ready: metadata_ready,
..Default::default()
});
let parts = build_health_response_parts(Method::GET, HealthProbe::Readiness, Some(&report), "rustfs", None, None);
let expected_ready = read_quorum && lock_ready;
@@ -510,6 +512,36 @@ mod tests {
});
}
#[test]
#[serial]
fn node_storage_details_explain_excluded_drives_without_exposing_endpoints() {
with_var(rustfs_config::ENV_HEALTH_MINIMAL_RESPONSE_ENABLE, Some("false"), || {
let mut report = ready_report();
report.readiness.storage_ready = false;
report.storage_details = Some(crate::shared_types::StorageReadinessDetails {
pool_metadata_write_ready: true,
unavailable_drives: vec![crate::shared_types::UnavailableReadinessDrive {
pool_index: 0,
set_index: 1,
disk_index: 2,
runtime_state: "suspect".to_string(),
host_online: true,
}],
..Default::default()
});
let parts = build_health_response_parts(Method::GET, HealthProbe::Readiness, Some(&report), "rustfs", None, None);
assert_eq!(parts.status_code, StatusCode::SERVICE_UNAVAILABLE);
let payload = parts.payload.expect("readiness GET body");
assert_eq!(
payload["details"]["storage"]["unavailableDrives"],
json!([{
"poolIndex": 0, "setIndex": 1, "diskIndex": 2,
"runtimeState": "suspect", "hostOnline": true,
}])
);
});
}
#[test]
#[serial]
fn node_storage_details_do_not_expand_minimal_or_liveness_payloads() {
@@ -518,6 +550,7 @@ mod tests {
read_quorum_ready: true,
write_quorum_ready: true,
pool_metadata_write_ready: true,
..Default::default()
});
with_var(rustfs_config::ENV_HEALTH_MINIMAL_RESPONSE_ENABLE, Some("true"), || {
let parts = build_health_response_parts(Method::GET, HealthProbe::Readiness, Some(&report), "rustfs", None, None);
+61 -18
View File
@@ -77,7 +77,9 @@ const METRIC_RUNTIME_READINESS_READY: &str = "rustfs_runtime_readiness_ready";
const METRIC_RUNTIME_READINESS_DEGRADED_TOTAL: &str = "rustfs_runtime_readiness_degraded_total";
const METRIC_POOL_METADATA_CHECK_TIMEOUT_TOTAL: &str = "rustfs_pool_metadata_check_timeouts_total";
pub use crate::shared_types::{DependencyReadiness, DependencyReadinessReport, ReadinessDegradedReason, StorageReadinessDetails};
pub use crate::shared_types::{
DependencyReadiness, DependencyReadinessReport, ReadinessDegradedReason, StorageReadinessDetails, UnavailableReadinessDrive,
};
/// ReadinessGateLayer ensures that the system components (IAM, Storage)
/// are fully initialized before allowing any request to proceed.
@@ -920,7 +922,8 @@ pub async fn collect_node_readiness_report() -> DependencyReadinessReport {
let pool_metadata_status = store.pool_meta_write_status().await;
storage = apply_node_pool_metadata_timeout_policy(pool_metadata_write_readiness(pool_metadata_status), Instant::now());
match node_storage_snapshot(store.as_ref(), &lock_observation.online_hosts).await {
Ok(info) => {
Ok((info, unavailable_drives)) => {
details.unavailable_drives = unavailable_drives;
details.read_quorum_ready = storage_read_ready_from_runtime_state(&info);
details.write_quorum_ready = storage_ready_from_runtime_state(&info);
}
@@ -947,7 +950,10 @@ pub async fn collect_node_readiness_report() -> DependencyReadinessReport {
report
}
async fn node_storage_snapshot<S>(store: &S, online_hosts: &HashSet<String>) -> Result<StorageInfo, StorageError>
async fn node_storage_snapshot<S>(
store: &S,
online_hosts: &HashSet<String>,
) -> Result<(StorageInfo, Vec<UnavailableReadinessDrive>), StorageError>
where
S: StorageAdminApi<BackendInfo = BackendInfo, Disk = DiskStore, Error = StorageError>,
{
@@ -956,8 +962,9 @@ where
backend: store.backend_info().await,
..Default::default()
};
let mut unavailable_drives = Vec::new();
if configured_readiness_topology(&info).is_none() {
return Ok(info);
return Ok((info, unavailable_drives));
}
// Inventory and runtime health are local snapshots. Reuse the lock probe's
@@ -971,26 +978,38 @@ where
.into_iter()
.flatten()
{
info.disks.push(node_disk_snapshot(
disk_endpoint_snapshot(&disk),
disk.runtime_state().as_str(),
online_hosts,
));
let (snapshot, unavailable) =
node_disk_snapshot(disk_endpoint_snapshot(&disk), disk.runtime_state().as_str(), online_hosts);
if let Some(unavailable) = unavailable {
unavailable_drives.push(unavailable);
}
info.disks.push(snapshot);
}
}
}
Ok(info)
Ok((info, unavailable_drives))
})
.await
.map_err(|_| StorageError::Timeout)?
}
fn node_disk_snapshot(endpoint: Endpoint, runtime_state: &str, online_hosts: &HashSet<String>) -> Disk {
fn node_disk_snapshot(
endpoint: Endpoint,
runtime_state: &str,
online_hosts: &HashSet<String>,
) -> (Disk, Option<UnavailableReadinessDrive>) {
let reachable = endpoint.is_local || online_hosts.contains(&endpoint.host_port());
// Returning drives can still reject data I/O as faulty. Without a fresh
// disk-info probe, only an Online runtime observation can supply quorum.
let online = reachable && runtime_state == rustfs_madmin::ITEM_ONLINE;
Disk {
let unavailable = (!online).then(|| UnavailableReadinessDrive {
pool_index: endpoint.pool_idx,
set_index: endpoint.set_idx,
disk_index: endpoint.disk_idx,
runtime_state: runtime_state.to_string(),
host_online: reachable,
});
let disk = Disk {
endpoint: endpoint.to_string(),
drive_path: endpoint.get_file_path(),
pool_index: endpoint.pool_idx,
@@ -999,7 +1018,8 @@ fn node_disk_snapshot(endpoint: Endpoint, runtime_state: &str, online_hosts: &Ha
state: if online { DISK_STATE_OK } else { "offline" }.to_string(),
runtime_state: Some(runtime_state.to_string()),
..Default::default()
}
};
(disk, unavailable)
}
async fn collect_cluster_health_report_with<LoadFn, Fut>(
@@ -1392,10 +1412,16 @@ mod tests {
let write_quorum = data + usize::from(data == parity);
for survivors in (0..=drive_count).rev().chain(std::iter::once(drive_count)) {
let online_hosts = (0..survivors).map(|idx| format!("node-{idx}:9000")).collect();
let info = node_storage_snapshot(&store, &online_hosts)
let (info, unavailable_drives) = node_storage_snapshot(&store, &online_hosts)
.await
.expect("read local runtime inventory");
assert_eq!(info.disks.len(), drive_count, "offline members retain their topology slots");
assert_eq!(unavailable_drives.len(), drive_count - survivors);
assert!(
unavailable_drives
.iter()
.all(|disk| !disk.host_online && disk.runtime_state == "online")
);
assert_eq!(
storage_read_ready_from_runtime_state(&info),
survivors >= data,
@@ -1428,7 +1454,22 @@ mod tests {
disk_idx: 2,
};
let mut disks = online_readiness_disks(0, 2);
disks.push(node_disk_snapshot(endpoint, runtime_state, &online_hosts));
let (disk, unavailable) = node_disk_snapshot(endpoint, runtime_state, &online_hosts);
if (is_local || reachable) && runtime_state == "online" {
assert!(unavailable.is_none());
} else {
assert_eq!(
unavailable,
Some(UnavailableReadinessDrive {
pool_index: 0,
set_index: 0,
disk_index: 2,
runtime_state: runtime_state.to_string(),
host_online: is_local || reachable,
})
);
}
disks.push(disk);
let info = StorageInfo {
backend: BackendInfo {
total_sets: vec![1],
@@ -1455,21 +1496,21 @@ mod tests {
async fn node_storage_snapshot_requires_each_configured_set_and_distinct_drives() {
let mut store = runtime_inventory(&[(2, 4, 2), (1, 8, 2)]).await;
let online_hosts = (0..8).map(|idx| format!("node-{idx}:9000")).collect();
let info = node_storage_snapshot(&store, &online_hosts)
let (info, _) = node_storage_snapshot(&store, &online_hosts)
.await
.expect("healthy mixed layout");
assert!(storage_ready_from_runtime_state(&info));
let selector = DiskSetSelector::new(0, 1);
let original = store.disks.remove(&selector).expect("second configured set");
let info = node_storage_snapshot(&store, &online_hosts)
let (info, _) = node_storage_snapshot(&store, &online_hosts)
.await
.expect("missing set snapshot");
assert!(!storage_read_ready_from_runtime_state(&info));
assert!(!storage_ready_from_runtime_state(&info));
store.disks.insert(selector, vec![original[0].clone(); 4]);
let info = node_storage_snapshot(&store, &online_hosts)
let (info, _) = node_storage_snapshot(&store, &online_hosts)
.await
.expect("duplicate drive snapshot");
assert!(!storage_read_ready_from_runtime_state(&info));
@@ -1480,6 +1521,7 @@ mod tests {
&node_storage_snapshot(&store, &online_hosts)
.await
.expect("restored inventory")
.0
));
}
@@ -1502,6 +1544,7 @@ mod tests {
&node_storage_snapshot(&store, &online_hosts)
.await
.expect("inspection recovered")
.0
));
}
+14 -1
View File
@@ -85,11 +85,24 @@ pub struct DependencyReadinessReport {
pub storage_details: Option<StorageReadinessDetails>,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct StorageReadinessDetails {
pub read_quorum_ready: bool,
pub write_quorum_ready: bool,
pub pool_metadata_write_ready: bool,
pub unavailable_drives: Vec<UnavailableReadinessDrive>,
}
/// Node-local reasons for excluding an inventoried drive from readiness quorum.
/// Use topology indices instead of exposing internal addresses or filesystem paths.
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct UnavailableReadinessDrive {
pub pool_index: i32,
pub set_index: i32,
pub disk_index: i32,
pub runtime_state: String,
pub host_online: bool,
}
pub(crate) fn convert_ecstore_object_info(object: StorageObjectInfo) -> NotifyObjectInfo {