mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-02 05:14:35 +08:00
fix(ecstore): recover remote disks after transient stalls
Merge the approved fix from pull request #8149.
This commit is contained in:
@@ -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");
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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`:
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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
|
||||
));
|
||||
}
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user