mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 08:18:18 +00:00
fix(admin): report unreachable info members as unknown with balanced drive counts (#4607)
fix(admin): report unreachable info members as `unknown` with balanced drive counts `GET /rustfs/admin/v3/info` synthesizes an entry for every peer, but a peer whose properties RPC failed below the offline threshold with no cached snapshot was reported as `initializing` with an empty drive list. That drive list being empty meant its drives disappeared from the backend summary entirely, so `onlineDisks + offlineDisks` no longer equaled `totalDrivesPerSet` (an operator sees `onlineDisks: 3, offlineDisks: 0` for a 4-drive set), and `initializing` is a misleading state for a member that has been running for hours — it reads like a node was ejected when nothing is wrong (upstream rustfs/rustfs#4566). Changes: - madmin: add `ItemState::Unknown` ("unknown") and an `unknownDisks` bucket on `ErasureBackend` (serde `default`, so older payloads still deserialize). - ecstore diagnostics: `get_online_offline_disks_stats` now returns a third `unknown` bucket instead of folding those drives into `offline`, keeping `online + offline + unknown == totalDrivesPerSet`. - ecstore notification_sys: a probe miss below the failure threshold with no cache is now reported as `unknown` (not `initializing`) and carries the host's drives from the pool topology (tagged `unknown`) so the summary stays balanced; drive synthesis is factored into `synthesized_disks(host, endpoints, state)`, shared by the offline and unknown paths. Tracked in rustfs/backlog#1049 (P1-A + P0-A). Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -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<String, usize> = HashMap::new();
|
||||
let mut offline_disks: HashMap<String, usize> = HashMap::new();
|
||||
let mut unknown_disks: HashMap<String, usize> = 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::<usize>()) {
|
||||
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::<usize>() + unknown_disks.values().sum::<usize>()) {
|
||||
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<HashMap<i32, HashMap<i32, ErasureSetInfo>>> {
|
||||
@@ -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]
|
||||
|
||||
@@ -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<ServerProperties>,
|
||||
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<SystemTime>,
|
||||
}
|
||||
|
||||
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<NotificationSys> = 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<PeerAdminCache>>, 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<rustfs_madmin::Disk>,
|
||||
}
|
||||
|
||||
/// 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<PeerDiskHealth> {
|
||||
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<PeerAdminCache>>,
|
||||
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<PeerAdminCache>>, host: &str, info: &ServerProperties) {
|
||||
@@ -1163,13 +1314,31 @@ fn update_server_info_cache(cache: Option<&Mutex<PeerAdminCache>>, 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<SystemTime>) -> 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<rustfs_madmin::Disk> {
|
||||
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<rustfs_madmin::Disk> {
|
||||
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<String>) -> 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");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<usize>,
|
||||
#[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");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user