diff --git a/crates/ecstore/src/diagnostics/admin_server_info.rs b/crates/ecstore/src/diagnostics/admin_server_info.rs index 76fc5b318..31cc656c9 100644 --- a/crates/ecstore/src/diagnostics/admin_server_info.rs +++ b/crates/ecstore/src/diagnostics/admin_server_info.rs @@ -21,7 +21,8 @@ use crate::data_usage::load_data_usage_cache; use crate::storage_api_contracts::admin::StorageAdminApi; use rustfs_common::heal_channel::DriveState; use rustfs_madmin::{ - BackendDisks, Disk, ErasureSetInfo, ITEM_INITIALIZING, ITEM_OFFLINE, ITEM_ONLINE, InfoMessage, MemStats, ServerProperties, + BackendDisks, Disk, ErasureSetInfo, ITEM_INITIALIZING, ITEM_OFFLINE, ITEM_ONLINE, ITEM_UNKNOWN, InfoMessage, MemStats, + ServerProperties, }; use rustfs_protos::{ models::{PingBody, PingBodyBuilder}, @@ -276,7 +277,7 @@ pub async fn get_server_info(get_pools: bool) -> InfoMessage { for server in servers.iter() { all_disks.extend(server.disks.clone()); } - let (online_disks, offline_disks) = get_online_offline_disks_stats(&all_disks); + let (online_disks, offline_disks, unknown_disks) = get_online_offline_disks_stats(&all_disks); let after5 = OffsetDateTime::now_utc(); @@ -285,6 +286,7 @@ pub async fn get_server_info(get_pools: bool) -> InfoMessage { backend_type: rustfs_madmin::BackendType::ErasureType, online_disks: online_disks.sum(), offline_disks: offline_disks.sum(), + unknown_disks: unknown_disks.sum(), standard_sc_parity: backend_info.standard_sc_parity, rr_sc_parity: backend_info.rr_sc_parity, total_sets: backend_info.total_sets, @@ -318,19 +320,32 @@ pub async fn get_server_info(get_pools: bool) -> InfoMessage { } } -fn get_online_offline_disks_stats(disks_info: &[Disk]) -> (BackendDisks, BackendDisks) { +/// Classify every drive into online / offline / unknown buckets. +/// +/// `unknown` holds drives synthesized for a member whose properties RPC could +/// not be answered this cycle but which is not confirmed offline. Keeping them +/// out of the `offline` bucket means a transient probe miss no longer inflates +/// the offline count for a healthy member, while `online + offline + unknown` +/// still sums to the pool's total drive count (rustfs/backlog#1049). +fn get_online_offline_disks_stats(disks_info: &[Disk]) -> (BackendDisks, BackendDisks, BackendDisks) { let mut online_disks: HashMap = HashMap::new(); let mut offline_disks: HashMap = HashMap::new(); + let mut unknown_disks: HashMap = HashMap::new(); for disk in disks_info { let ep = &disk.endpoint; offline_disks.entry(ep.clone()).or_insert(0); online_disks.entry(ep.clone()).or_insert(0); + unknown_disks.entry(ep.clone()).or_insert(0); } for disk in disks_info { let ep = &disk.endpoint; let state = &disk.state; + if *state == ITEM_UNKNOWN { + *unknown_disks.get_mut(ep).expect("endpoint should be in disk map") += 1; + continue; + } if *state != DriveState::Ok.to_string() && *state != DriveState::Unformatted.to_string() { *offline_disks.get_mut(ep).expect("endpoint should be in disk map") += 1; continue; @@ -345,8 +360,11 @@ fn get_online_offline_disks_stats(disks_info: &[Disk]) -> (BackendDisks, Backend } } - if disks_info.len() == (root_disk_count + offline_disks.values().sum::()) { - return (BackendDisks(online_disks), BackendDisks(offline_disks)); + // When every non-offline, non-unknown drive is a root mount, leave the + // online tally as-is instead of demoting all of them (matches the prior + // behavior; the unknown bucket is simply carried through untouched). + if disks_info.len() == (root_disk_count + offline_disks.values().sum::() + unknown_disks.values().sum::()) { + return (BackendDisks(online_disks), BackendDisks(offline_disks), BackendDisks(unknown_disks)); } for disk in disks_info { @@ -357,7 +375,7 @@ fn get_online_offline_disks_stats(disks_info: &[Disk]) -> (BackendDisks, Backend } } - (BackendDisks(online_disks), BackendDisks(offline_disks)) + (BackendDisks(online_disks), BackendDisks(offline_disks), BackendDisks(unknown_disks)) } async fn get_pools_info(all_disks: &[Disk]) -> Result>> { @@ -413,8 +431,45 @@ mod tests { use serial_test::serial; use crate::runtime::sources as runtime_sources; + use rustfs_madmin::{Disk, ITEM_OFFLINE, ITEM_UNKNOWN}; - use super::{get_local_server_property, get_server_info}; + use super::{get_local_server_property, get_online_offline_disks_stats, get_server_info}; + + fn disk_with_state(endpoint: &str, state: &str) -> Disk { + Disk { + endpoint: endpoint.to_string(), + state: state.to_string(), + ..Default::default() + } + } + + #[test] + fn disk_stats_split_unknown_into_its_own_bucket() { + // A member whose properties RPC could not be answered contributes + // drives tagged `unknown`. They must land in the unknown bucket, not + // inflate `offline`, and the three buckets must still account for every + // drive so the summary stays balanced (rustfs/backlog#1049). + // + // A live drive reports the DriveState string "ok"; only "ok"/"unformatted" + // count as online. + let disks = vec![ + disk_with_state("http://n1:9000/data", "ok"), + disk_with_state("http://n2:9000/data", "ok"), + disk_with_state("http://n3:9000/data", ITEM_OFFLINE), + disk_with_state("http://n4:9000/data", ITEM_UNKNOWN), + ]; + + let (online, offline, unknown) = get_online_offline_disks_stats(&disks); + + assert_eq!(online.sum(), 2, "the two healthy drives are online"); + assert_eq!(offline.sum(), 1, "only the confirmed-offline drive is offline"); + assert_eq!(unknown.sum(), 1, "the unreachable member's drive is unknown, not offline"); + assert_eq!( + online.sum() + offline.sum() + unknown.sum(), + disks.len(), + "online + offline + unknown must equal the total drive count" + ); + } #[serial] #[tokio::test] diff --git a/crates/ecstore/src/services/notification_sys.rs b/crates/ecstore/src/services/notification_sys.rs index e9a4b8de2..72578d73e 100644 --- a/crates/ecstore/src/services/notification_sys.rs +++ b/crates/ecstore/src/services/notification_sys.rs @@ -14,6 +14,7 @@ use crate::cluster::rpc::PeerRestClient; use crate::diagnostics::admin_server_info::get_commit_id; +use crate::disk::DiskAPI; use crate::error::{Error, Result}; use crate::layout::endpoints::EndpointServerPools; use crate::runtime::sources as runtime_sources; @@ -46,6 +47,10 @@ struct PeerAdminCache { last_server_info: Option, storage_failures: u32, server_failures: u32, + /// When the last successful server_info probe landed. Used to stop a stale + /// cached `online` snapshot from being served indefinitely while a peer is + /// actually down (rustfs/backlog#1049 P2). + last_server_success: Option, } impl PeerAdminCache { @@ -55,10 +60,16 @@ impl PeerAdminCache { last_server_info: None, storage_failures: 0, server_failures: 0, + last_server_success: None, } } } +/// A cached `online` snapshot older than this is no longer trusted on a probe +/// failure: rather than reporting a stale `online`, the member falls through to +/// the live unknown/degraded/offline classification (rustfs/backlog#1049 P2). +const SERVER_INFO_CACHE_MAX_AGE: Duration = Duration::from_secs(60); + lazy_static! { pub static ref GLOBAL_NOTIFICATION_SYS: OnceLock = OnceLock::new(); } @@ -346,25 +357,55 @@ impl NotificationSys { let endpoints = endpoints.clone(); let cache = self.peer_admin_caches.get(idx); futures.push(async move { - if let Some(client) = client { - let host = client.host.to_string(); - match timeout(peer_timeout, client.server_info()).await { - Ok(Ok(info)) => { - update_server_info_cache(cache, &host, &info); - info - } - Ok(Err(err)) => { - warn!("peer {} server_info failed: {}", host, err); - handle_server_info_failure(cache, &host, &endpoints) - } - Err(_) => { - warn!("peer {} server_info timed out after {:?}", host, peer_timeout); - client.evict_connection().await; - handle_server_info_failure(cache, &host, &endpoints) - } + // `peer_clients` comes from `new_clients`, which only ever pushes + // `Some(client)` (local hosts are excluded, not slotted as + // `None`), so this branch is unreachable in practice. Kept as a + // defensive fallback: report an explicit `unknown` state rather + // than a blank `default()` entry, so it can never contribute a + // hollow row to `servers[]` (rustfs/backlog#1049 P3). + let Some(client) = client else { + return ServerProperties { + state: ItemState::Unknown.to_string().to_owned(), + ..Default::default() + }; + }; + let host = client.host.to_string(); + + // First attempt. A single evicted or half-open internode channel + // is enough to fail one probe and, before retrying, would drop + // the member to unknown/offline for this whole snapshot. So on any + // first-attempt failure we evict the channel and re-dial once + // before falling back (rustfs/backlog#1049, P1-B). + match timeout(peer_timeout, client.server_info()).await { + Ok(Ok(info)) => { + update_server_info_cache(cache, &host, &info); + return info; + } + Ok(Err(err)) => debug!("peer {host} server_info failed (attempt 1/2): {err}"), + Err(_) => debug!("peer {host} server_info timed out (attempt 1/2) after {peer_timeout:?}"), + } + + // Drop the suspect channel so the retry re-dials on a fresh + // connection instead of reusing the bad one. + client.evict_connection().await; + + // Second and final attempt on the fresh channel. + match timeout(peer_timeout, client.server_info()).await { + Ok(Ok(info)) => { + update_server_info_cache(cache, &host, &info); + info + } + Ok(Err(err)) => { + warn!("peer {host} server_info failed after retry: {err}"); + let health = peer_disk_health(&host).await; + handle_server_info_failure(cache, &host, &endpoints, health.as_ref()) + } + Err(_) => { + warn!("peer {host} server_info timed out after retry ({peer_timeout:?})"); + client.evict_connection().await; + let health = peer_disk_health(&host).await; + handle_server_info_failure(cache, &host, &endpoints, health.as_ref()) } - } else { - ServerProperties::default() } }); } @@ -1047,7 +1088,7 @@ fn handle_peer_failure( ); } return Some(StorageInfo { - disks: get_offline_disks(host, endpoints), + disks: synthesized_disks(host, endpoints, ItemState::Offline), ..Default::default() }); } @@ -1080,15 +1121,91 @@ fn update_storage_info_cache(cache: Option<&Mutex>, host: &str, c.storage_failures = 0; } -/// Handle a peer failure for server_info: return cached data if available, -/// or mark offline only after consecutive failures exceed the threshold. +/// Independent liveness evidence for a peer, gathered from the local disk-health +/// heartbeat rather than the admin RPC path. `any_online` is true when at least +/// one of the peer's drives is still answering the ~15s health check; `disks` +/// carries a per-drive entry (state `"ok"` when online, `"offline"` when the +/// heartbeat marks it faulty) so a `degraded` member's drives are counted for +/// real. See rustfs/backlog#1049 (P0-B). +struct PeerDiskHealth { + any_online: bool, + disks: Vec, +} + +/// Consult the local disk-health state for `host` without issuing any RPC. +/// +/// On the aggregating node a peer's drives are remote-disk handles whose +/// `is_online()` is a pure atomic read of the heartbeat tracker (independent of +/// the admin `server_info` RPC that just failed). Returns `None` when the store +/// is not initialized or the host owns no drives in the topology. +async fn peer_disk_health(host: &str) -> Option { + let store = runtime_sources::object_store_handle()?; + + let mut disks = Vec::new(); + let mut any_online = false; + for sets in store.pools.iter() { + for set in sets.disk_set.iter() { + let guard = set.disks.read().await; + for (idx, slot) in guard.iter().enumerate() { + let Some(ep) = set.set_endpoints.get(idx) else { + continue; + }; + if host != ep.host_port() { + continue; + } + let online = match slot { + Some(disk) => disk.is_online().await, + None => false, + }; + any_online |= online; + // A live drive is counted online via the DriveState "ok" string; + // a faulty one is counted offline. This keeps a degraded member's + // drives in the real online/offline buckets. + disks.push(rustfs_madmin::Disk { + endpoint: ep.to_string(), + state: if online { + rustfs_common::heal_channel::DriveState::Ok.to_string() + } else { + ItemState::Offline.to_string().to_owned() + }, + pool_index: ep.pool_idx, + set_index: ep.set_idx, + disk_index: ep.disk_idx, + ..Default::default() + }); + } + } + } + + if disks.is_empty() { + None + } else { + Some(PeerDiskHealth { any_online, disks }) + } +} + +/// Handle a peer failure for server_info: return cached data if available, or +/// classify the member as `unknown` / `degraded` / `offline` depending on how +/// many consecutive probes have failed and whether the peer's drives are still +/// answering the local disk-health heartbeat. +/// +/// - Below the failure threshold with no cached snapshot: `unknown` (probe +/// missed this cycle but the member is not confirmed down). +/// - At/after the threshold with drives still online: `degraded` (the admin RPC +/// is stuck but the node is alive and serving data) — this is what stops a +/// healthy node from rotating through a false `offline` (rustfs/backlog#1049). +/// - At/after the threshold with drives also offline: `offline` (confirmed). +/// +/// Synthesized entries always carry one drive per endpoint so the pool's drive +/// totals stay balanced. fn handle_server_info_failure( cache: Option<&Mutex>, host: &str, endpoints: &EndpointServerPools, + peer_health: Option<&PeerDiskHealth>, ) -> ServerProperties { let Some(cache) = cache else { - return initializing_server_properties(host); + return unknown_server_properties(host, endpoints); }; let mut c = match cache.lock() { @@ -1103,17 +1220,51 @@ fn handle_server_info_failure( if let Some(ref cached) = c.last_server_info && c.server_failures < CONSECUTIVE_FAILURE_THRESHOLD { + if cached_snapshot_is_fresh(c.last_server_success) { + debug!( + event = "peer_probe_failure", + peer = host, + consecutive_failures = c.server_failures, + threshold = CONSECUTIVE_FAILURE_THRESHOLD, + "peer server_info probe failed; returning cached state until the offline threshold is reached" + ); + return cached.clone(); + } + // The cached snapshot is too old to keep reporting as `online`; fall + // through to the live unknown/degraded/offline classification below + // instead of masking a down peer with a stale success (P2). debug!( - event = "peer_probe_failure", + event = "peer_cache_stale", peer = host, - consecutive_failures = c.server_failures, - threshold = CONSECUTIVE_FAILURE_THRESHOLD, - "peer server_info probe failed; returning cached state until the offline threshold is reached" + max_age_secs = SERVER_INFO_CACHE_MAX_AGE.as_secs(), + "cached server_info snapshot is stale; reclassifying from live signals instead of reporting stale online" ); - return cached.clone(); } if c.server_failures >= CONSECUTIVE_FAILURE_THRESHOLD { + // Drives still answering the heartbeat: the node is alive, only its + // admin surface is unreachable — report `degraded`, not `offline`, so a + // stuck admin path does not read as an ejected node. + if let Some(health) = peer_health.filter(|h| h.any_online) { + if c.server_failures == CONSECUTIVE_FAILURE_THRESHOLD { + warn!( + event = "peer_marked_degraded", + peer = host, + consecutive_failures = c.server_failures, + threshold = CONSECUTIVE_FAILURE_THRESHOLD, + "peer admin server_info keeps failing but its drives are online; reporting degraded (not offline)" + ); + } else { + debug!( + event = "peer_still_degraded", + peer = host, + consecutive_failures = c.server_failures, + "peer admin server_info still failing while its drives remain online" + ); + } + return degraded_server_properties(host, &health.disks); + } + // Log the transition exactly once (at the crossing) so the console's // "node offline" verdict has a matching WARN in the observer's logs // (rustfs/backlog#888: nodes were marked offline with no log naming @@ -1139,7 +1290,7 @@ fn handle_server_info_failure( return offline_server_properties(host, endpoints); } - initializing_server_properties(host) + unknown_server_properties(host, endpoints) } fn update_server_info_cache(cache: Option<&Mutex>, host: &str, info: &ServerProperties) { @@ -1163,13 +1314,31 @@ fn update_server_info_cache(cache: Option<&Mutex>, host: &str, i ); } c.last_server_info = Some(info.clone()); + c.last_server_success = Some(SystemTime::now()); c.server_failures = 0; } -fn initializing_server_properties(host: &str) -> ServerProperties { +/// Whether a cached server_info snapshot is recent enough to still report as +/// `online` on a probe failure. A missing timestamp means no age information is +/// available (e.g. a snapshot set without going through the success path in a +/// test); such a snapshot is treated as fresh to preserve the prior behavior, +/// while a clock that went backwards is treated as stale. See P2 in +/// rustfs/backlog#1049. +fn cached_snapshot_is_fresh(last_success: Option) -> bool { + match last_success { + Some(at) => at.elapsed().map(|age| age < SERVER_INFO_CACHE_MAX_AGE).unwrap_or(false), + None => true, + } +} + +/// A member that could not be probed this cycle and is not confirmed down. +/// Carries the endpoint's drives (marked `unknown`) so the pool's drive totals +/// stay balanced instead of the member's drives vanishing from the summary. +fn unknown_server_properties(host: &str, endpoints: &EndpointServerPools) -> ServerProperties { ServerProperties { endpoint: host.to_string(), - state: ItemState::Initializing.to_string().to_owned(), + state: ItemState::Unknown.to_string().to_owned(), + disks: synthesized_disks(host, endpoints, ItemState::Unknown), ..Default::default() } } @@ -1180,20 +1349,37 @@ fn offline_server_properties(host: &str, endpoints: &EndpointServerPools) -> Ser version: get_commit_id(), endpoint: host.to_string(), state: ItemState::Offline.to_string().to_owned(), - disks: get_offline_disks(host, endpoints), + disks: synthesized_disks(host, endpoints, ItemState::Offline), ..Default::default() } } -fn get_offline_disks(offline_host: &str, endpoints: &EndpointServerPools) -> Vec { - let mut offline_disks = Vec::new(); +/// A member whose admin RPC is unreachable but whose drives are still online. +/// Carries the per-drive health observed from the local heartbeat so the drives +/// land in the real online/offline buckets while the member reads as degraded. +fn degraded_server_properties(host: &str, disks: &[rustfs_madmin::Disk]) -> ServerProperties { + ServerProperties { + uptime: runtime_sources::boot_uptime_secs(), + version: get_commit_id(), + endpoint: host.to_string(), + state: ItemState::Degraded.to_string().to_owned(), + disks: disks.to_vec(), + ..Default::default() + } +} + +/// Enumerate the drives a host owns from the pool topology, tagged with the +/// given member state. Used to synthesize drive entries for a member whose +/// properties RPC could not be answered, so summary counters stay complete. +fn synthesized_disks(host: &str, endpoints: &EndpointServerPools, state: ItemState) -> Vec { + let mut disks = Vec::new(); for pool in endpoints.as_ref() { for ep in pool.endpoints.as_ref() { - if (offline_host.is_empty() && ep.is_local) || offline_host == ep.host_port() { - offline_disks.push(rustfs_madmin::Disk { + if (host.is_empty() && ep.is_local) || host == ep.host_port() { + disks.push(rustfs_madmin::Disk { endpoint: ep.to_string(), - state: ItemState::Offline.to_string().to_owned(), + state: state.to_string().to_owned(), pool_index: ep.pool_idx, set_index: ep.set_idx, disk_index: ep.disk_idx, @@ -1203,7 +1389,7 @@ fn get_offline_disks(offline_host: &str, endpoints: &EndpointServerPools) -> Vec } } - offline_disks + disks } fn aggregate_notification_failures(operation: &str, failures: Vec) -> Result<()> { @@ -1415,6 +1601,7 @@ mod tests { last_server_info: None, storage_failures: 0, server_failures: 0, + last_server_success: None, }); let endpoints = EndpointServerPools::default(); @@ -1442,6 +1629,7 @@ mod tests { last_server_info: None, storage_failures: CONSECUTIVE_FAILURE_THRESHOLD - 1, server_failures: 0, + last_server_success: None, }); let endpoints = EndpointServerPools::default(); @@ -1464,23 +1652,69 @@ mod tests { last_server_info: Some(cached_props), storage_failures: 0, server_failures: 0, + last_server_success: None, }); let endpoints = EndpointServerPools::default(); - let result = handle_server_info_failure(Some(&cache), "peer-1", &endpoints); + let result = handle_server_info_failure(Some(&cache), "peer-1", &endpoints, None); assert_eq!(result.endpoint, "peer-1"); assert_eq!(result.state, "online"); assert_eq!(cache.lock().unwrap().server_failures, 1); } #[test] - fn handle_server_info_failure_returns_initializing_before_threshold_without_cache() { + fn handle_server_info_failure_does_not_serve_stale_cached_online() { + // A single failure with a cached snapshot would normally return the + // cached `online`, but if that snapshot is older than the max age we + // must not keep reporting online — fall through to `unknown` instead of + // masking a possibly-down peer (rustfs/backlog#1049 P2). + let cached_props = ServerProperties { + endpoint: "peer-1".to_string(), + state: "online".to_string(), + ..Default::default() + }; + let stale_at = SystemTime::now() + .checked_sub(SERVER_INFO_CACHE_MAX_AGE + Duration::from_secs(1)) + .expect("test clock underflow"); + + let cache = Mutex::new(PeerAdminCache { + last_storage_info: None, + last_server_info: Some(cached_props), + storage_failures: 0, + server_failures: 0, + last_server_success: Some(stale_at), + }); + let endpoints = EndpointServerPools::default(); + + let result = handle_server_info_failure(Some(&cache), "peer-1", &endpoints, None); + assert_eq!(result.state, ItemState::Unknown.to_string()); + assert_eq!(cache.lock().unwrap().server_failures, 1); + } + + #[test] + fn cached_snapshot_freshness_respects_age_and_missing_timestamp() { + assert!(cached_snapshot_is_fresh(None), "no timestamp is treated as fresh"); + assert!(cached_snapshot_is_fresh(Some(SystemTime::now())), "a just-now success is fresh"); + let stale = SystemTime::now() + .checked_sub(SERVER_INFO_CACHE_MAX_AGE + Duration::from_secs(1)) + .expect("test clock underflow"); + assert!(!cached_snapshot_is_fresh(Some(stale)), "an old success is stale"); + } + + #[test] + fn handle_server_info_failure_returns_unknown_before_threshold_without_cache() { let cache = Mutex::new(PeerAdminCache::new()); let endpoints = EndpointServerPools::default(); - let result = handle_server_info_failure(Some(&cache), "peer-1", &endpoints); + let result = handle_server_info_failure(Some(&cache), "peer-1", &endpoints, None); assert_eq!(result.endpoint, "peer-1"); - assert_eq!(result.state, ItemState::Initializing.to_string()); + // A probe miss below the threshold is "unknown" (not confirmed down, + // and not the misleading "initializing"): rustfs/backlog#1049. + assert_eq!(result.state, ItemState::Unknown.to_string()); + // The default (empty) pool has no topology entry for this host, so no + // drives are synthesized here; the drive-synthesis and counter-balance + // behavior is exercised by the get_online_offline_disks_stats tests in + // admin_server_info. assert!(result.disks.is_empty()); assert_eq!(cache.lock().unwrap().server_failures, 1); } @@ -1498,14 +1732,69 @@ mod tests { last_server_info: Some(cached_props), storage_failures: 0, server_failures: CONSECUTIVE_FAILURE_THRESHOLD - 1, + last_server_success: None, }); let endpoints = EndpointServerPools::default(); - let result = handle_server_info_failure(Some(&cache), "peer-1", &endpoints); + let result = handle_server_info_failure(Some(&cache), "peer-1", &endpoints, None); assert_eq!(result.state, ItemState::Offline.to_string()); assert_eq!(cache.lock().unwrap().server_failures, CONSECUTIVE_FAILURE_THRESHOLD); } + #[test] + fn handle_server_info_failure_returns_degraded_when_disks_online_past_threshold() { + // Past the threshold but the peer's drives still answer the heartbeat: + // the node is alive, only its admin RPC is stuck — report degraded (with + // the real per-drive health), not offline (rustfs/backlog#1049 P0-B). + let cache = Mutex::new(PeerAdminCache { + last_storage_info: None, + last_server_info: None, + storage_failures: 0, + server_failures: CONSECUTIVE_FAILURE_THRESHOLD - 1, + last_server_success: None, + }); + let endpoints = EndpointServerPools::default(); + let health = PeerDiskHealth { + any_online: true, + disks: vec![rustfs_madmin::Disk { + endpoint: "http://peer-1:9000/data".to_string(), + state: "ok".to_string(), + ..Default::default() + }], + }; + + let result = handle_server_info_failure(Some(&cache), "peer-1", &endpoints, Some(&health)); + assert_eq!(result.state, ItemState::Degraded.to_string()); + assert_eq!(result.disks.len(), 1); + assert_eq!(result.disks[0].state, "ok"); + assert_eq!(cache.lock().unwrap().server_failures, CONSECUTIVE_FAILURE_THRESHOLD); + } + + #[test] + fn handle_server_info_failure_stays_offline_when_disks_also_offline() { + // Past the threshold and the heartbeat also reports the drives down: + // this is a genuine offline, degraded must not mask it. + let cache = Mutex::new(PeerAdminCache { + last_storage_info: None, + last_server_info: None, + storage_failures: 0, + server_failures: CONSECUTIVE_FAILURE_THRESHOLD - 1, + last_server_success: None, + }); + let endpoints = EndpointServerPools::default(); + let health = PeerDiskHealth { + any_online: false, + disks: vec![rustfs_madmin::Disk { + endpoint: "http://peer-1:9000/data".to_string(), + state: ItemState::Offline.to_string().to_owned(), + ..Default::default() + }], + }; + + let result = handle_server_info_failure(Some(&cache), "peer-1", &endpoints, Some(&health)); + assert_eq!(result.state, ItemState::Offline.to_string()); + } + #[test] fn success_resets_failure_counters_independently() { let cache = Mutex::new(PeerAdminCache { @@ -1513,6 +1802,7 @@ mod tests { last_server_info: None, storage_failures: 2, server_failures: 2, + last_server_success: None, }); { @@ -1537,13 +1827,14 @@ mod tests { }), storage_failures: CONSECUTIVE_FAILURE_THRESHOLD - 1, server_failures: 0, + last_server_success: None, }); let endpoints = EndpointServerPools::default(); let storage_result = handle_peer_failure(Some(&cache), "peer-1", &endpoints); assert!(storage_result.is_some()); - let server_result = handle_server_info_failure(Some(&cache), "peer-1", &endpoints); + let server_result = handle_server_info_failure(Some(&cache), "peer-1", &endpoints, None); assert_eq!(server_result.state, "online"); assert_eq!(cache.lock().unwrap().server_failures, 1); } @@ -1566,9 +1857,9 @@ mod tests { let storage_result = handle_peer_failure(Some(&storage_cache), "peer-1", &endpoints); assert!(storage_result.is_none()); - let server_result = handle_server_info_failure(Some(&server_cache), "peer-1", &endpoints); + let server_result = handle_server_info_failure(Some(&server_cache), "peer-1", &endpoints, None); assert_eq!(server_result.endpoint, "peer-1"); - assert_eq!(server_result.state, ItemState::Initializing.to_string()); + assert_eq!(server_result.state, ItemState::Unknown.to_string()); } #[test] @@ -1578,12 +1869,14 @@ mod tests { last_server_info: None, storage_failures: CONSECUTIVE_FAILURE_THRESHOLD - 1, server_failures: 0, + last_server_success: None, }); let server_cache = Mutex::new(PeerAdminCache { last_storage_info: None, last_server_info: None, storage_failures: 0, server_failures: CONSECUTIVE_FAILURE_THRESHOLD - 1, + last_server_success: None, }); let endpoints = EndpointServerPools::default(); @@ -1622,7 +1915,7 @@ mod tests { assert!(storage_result.is_some()); assert_eq!(storage_result.unwrap().disks[0].state, "ok"); - let server_result = handle_server_info_failure(Some(&server_cache), "peer-1", &endpoints); + let server_result = handle_server_info_failure(Some(&server_cache), "peer-1", &endpoints, None); assert_eq!(server_result.state, "online"); } } diff --git a/crates/madmin/src/info_commands.rs b/crates/madmin/src/info_commands.rs index 1e35c7cea..d0759aea8 100644 --- a/crates/madmin/src/info_commands.rs +++ b/crates/madmin/src/info_commands.rs @@ -24,6 +24,18 @@ pub enum ItemState { Offline, Initializing, Online, + /// The peer could not be probed this cycle but is not confirmed down. + /// Distinct from `Offline` (confirmed unreachable past the failure + /// threshold) and `Initializing` (genuinely still starting up): a + /// transient probe miss on an otherwise-healthy member must not read as + /// an ejection. See rustfs/backlog#1049. + Unknown, + /// The peer's admin/properties RPC keeps failing, but its drives are still + /// answering the independent disk-health heartbeat — the node is alive and + /// serving data, only its admin surface is unreachable. Reported instead of + /// `Offline` so a stuck admin path (e.g. a saturated scanner) is not + /// mistaken for an ejected node. See rustfs/backlog#1049 (P0-B). + Degraded, } impl ItemState { @@ -32,6 +44,8 @@ impl ItemState { ItemState::Offline => "offline", ItemState::Initializing => "initializing", ItemState::Online => "online", + ItemState::Unknown => "unknown", + ItemState::Degraded => "degraded", } } @@ -40,6 +54,8 @@ impl ItemState { "offline" => Some(ItemState::Offline), "initializing" => Some(ItemState::Initializing), "online" => Some(ItemState::Online), + "unknown" => Some(ItemState::Unknown), + "degraded" => Some(ItemState::Degraded), _ => None, } } @@ -186,6 +202,8 @@ pub struct BackendInfo { pub const ITEM_OFFLINE: &str = "offline"; pub const ITEM_INITIALIZING: &str = "initializing"; pub const ITEM_ONLINE: &str = "online"; +pub const ITEM_UNKNOWN: &str = "unknown"; +pub const ITEM_DEGRADED: &str = "degraded"; #[derive(Clone, Debug, Default, Serialize, Deserialize)] pub struct MemStats { @@ -326,6 +344,14 @@ pub struct ErasureBackend { pub online_disks: usize, #[serde(rename = "offlineDisks")] pub offline_disks: usize, + /// Drives whose owning member could not be probed this cycle and is not + /// confirmed offline. Kept as a separate bucket so + /// `online + offline + unknown == totalDrivesPerSet` holds even while a + /// member is unreachable, instead of the fourth drive silently vanishing + /// from the totals (rustfs/backlog#1049). `#[serde(default)]` keeps older + /// payloads that predate this field deserializable. + #[serde(rename = "unknownDisks", default)] + pub unknown_disks: usize, #[serde(rename = "standardSCParity")] pub standard_sc_parity: Option, #[serde(rename = "rrSCParity")] @@ -409,6 +435,8 @@ mod tests { assert_eq!(ItemState::Offline.to_string(), ITEM_OFFLINE); assert_eq!(ItemState::Initializing.to_string(), ITEM_INITIALIZING); assert_eq!(ItemState::Online.to_string(), ITEM_ONLINE); + assert_eq!(ItemState::Unknown.to_string(), ITEM_UNKNOWN); + assert_eq!(ItemState::Degraded.to_string(), ITEM_DEGRADED); } #[test] @@ -416,6 +444,8 @@ mod tests { assert_eq!(ItemState::from_string(ITEM_OFFLINE), Some(ItemState::Offline)); assert_eq!(ItemState::from_string(ITEM_INITIALIZING), Some(ItemState::Initializing)); assert_eq!(ItemState::from_string(ITEM_ONLINE), Some(ItemState::Online)); + assert_eq!(ItemState::from_string(ITEM_UNKNOWN), Some(ItemState::Unknown)); + assert_eq!(ItemState::from_string(ITEM_DEGRADED), Some(ItemState::Degraded)); } #[test] @@ -1145,6 +1175,7 @@ mod tests { backend_type: BackendType::ErasureType, online_disks: 8, offline_disks: 0, + unknown_disks: 0, standard_sc_parity: Some(2), rr_sc_parity: Some(1), total_sets: vec![2], @@ -1154,6 +1185,7 @@ mod tests { assert!(matches!(erasure_backend.backend_type, BackendType::ErasureType)); assert_eq!(erasure_backend.online_disks, 8); assert_eq!(erasure_backend.offline_disks, 0); + assert_eq!(erasure_backend.unknown_disks, 0); assert_eq!(erasure_backend.standard_sc_parity.unwrap(), 2); assert_eq!(erasure_backend.rr_sc_parity.unwrap(), 1); assert_eq!(erasure_backend.total_sets.len(), 1); @@ -1271,5 +1303,7 @@ mod tests { assert_eq!(ITEM_OFFLINE, "offline"); assert_eq!(ITEM_INITIALIZING, "initializing"); assert_eq!(ITEM_ONLINE, "online"); + assert_eq!(ITEM_UNKNOWN, "unknown"); + assert_eq!(ITEM_DEGRADED, "degraded"); } }