Skip to main content

rc_core/admin/
observability.rs

1//! Typed contracts for RustFS scanner, storage, and realtime metrics APIs.
2
3use std::collections::BTreeMap;
4use std::fmt;
5
6use async_trait::async_trait;
7use serde::{Deserialize, Serialize};
8
9use crate::Result;
10
11/// Maximum number of snapshots accepted from one metrics request.
12pub const MAX_METRICS_SAMPLES: u16 = 120;
13/// Maximum encoded size of one NDJSON metrics record.
14pub const MAX_METRICS_LINE_BYTES: usize = 1024 * 1024;
15/// Maximum encoded size accepted for an entire metrics response.
16pub const MAX_METRICS_RESPONSE_BYTES: usize = 16 * 1024 * 1024;
17
18/// RustFS realtime metrics selector bits from beta.10.
19#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
20#[serde(rename_all = "kebab-case")]
21pub enum MetricsScope {
22    Scanner,
23    Disk,
24    Os,
25    BatchJobs,
26    SiteResync,
27    Network,
28    Memory,
29    Cpu,
30    Rpc,
31    All,
32}
33
34impl MetricsScope {
35    /// Return the exact `types` bit accepted by RustFS Admin API v3.
36    pub const fn bit(self) -> u32 {
37        match self {
38            Self::Scanner => 1 << 0,
39            Self::Disk => 1 << 1,
40            Self::Os => 1 << 2,
41            Self::BatchJobs => 1 << 3,
42            Self::SiteResync => 1 << 4,
43            Self::Network => 1 << 5,
44            Self::Memory => 1 << 6,
45            Self::Cpu => 1 << 7,
46            Self::Rpc => 1 << 8,
47            Self::All => (1 << 9) - 1,
48        }
49    }
50}
51
52impl fmt::Display for MetricsScope {
53    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
54        let value = match self {
55            Self::Scanner => "scanner",
56            Self::Disk => "disk",
57            Self::Os => "os",
58            Self::BatchJobs => "batch-jobs",
59            Self::SiteResync => "site-resync",
60            Self::Network => "network",
61            Self::Memory => "memory",
62            Self::Cpu => "cpu",
63            Self::Rpc => "rpc",
64            Self::All => "all",
65        };
66        formatter.write_str(value)
67    }
68}
69
70/// Bounded realtime metrics request.
71#[derive(Debug, Clone, PartialEq, Eq)]
72pub struct MetricsQuery {
73    pub scopes: Vec<MetricsScope>,
74    pub hosts: Vec<String>,
75    pub disks: Vec<String>,
76    pub interval: Option<String>,
77    pub samples: u16,
78    pub by_host: bool,
79    pub by_disk: bool,
80    pub job_id: Option<String>,
81    pub deployment_id: Option<String>,
82}
83
84impl Default for MetricsQuery {
85    fn default() -> Self {
86        Self {
87            scopes: vec![MetricsScope::All],
88            hosts: Vec::new(),
89            disks: Vec::new(),
90            interval: None,
91            samples: 1,
92            by_host: false,
93            by_disk: false,
94            job_id: None,
95            deployment_id: None,
96        }
97    }
98}
99
100impl MetricsQuery {
101    /// Combine requested selector bits without allowing duplicate scopes to alter the mask.
102    pub fn types_mask(&self) -> u32 {
103        if self.scopes.is_empty() || self.scopes.contains(&MetricsScope::All) {
104            return MetricsScope::All.bit();
105        }
106        self.scopes.iter().fold(0, |mask, scope| mask | scope.bit())
107    }
108
109    /// Stable label used by human and machine output.
110    pub fn scope_label(&self) -> String {
111        if self.scopes.is_empty() || self.scopes.contains(&MetricsScope::All) {
112            return "all".to_string();
113        }
114        let mut scopes = self
115            .scopes
116            .iter()
117            .map(ToString::to_string)
118            .collect::<Vec<_>>();
119        scopes.sort();
120        scopes.dedup();
121        scopes.join(",")
122    }
123}
124
125/// One scope-specific metrics object kept in its native JSON shape.
126#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
127#[serde(transparent)]
128pub struct MetricGroup(pub BTreeMap<String, serde_json::Value>);
129
130/// Scope groups carried by a RustFS realtime metrics snapshot.
131#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
132pub struct MetricGroups {
133    #[serde(default, skip_serializing_if = "Option::is_none")]
134    pub scanner: Option<MetricGroup>,
135    #[serde(default, skip_serializing_if = "Option::is_none")]
136    pub disk: Option<MetricGroup>,
137    #[serde(default, skip_serializing_if = "Option::is_none")]
138    pub os: Option<MetricGroup>,
139    #[serde(rename = "batchJobs", default, skip_serializing_if = "Option::is_none")]
140    pub batch_jobs: Option<MetricGroup>,
141    #[serde(
142        rename = "siteResync",
143        default,
144        skip_serializing_if = "Option::is_none"
145    )]
146    pub site_resync: Option<MetricGroup>,
147    #[serde(default, skip_serializing_if = "Option::is_none")]
148    pub net: Option<MetricGroup>,
149    #[serde(default, skip_serializing_if = "Option::is_none")]
150    pub mem: Option<MetricGroup>,
151    #[serde(default, skip_serializing_if = "Option::is_none")]
152    pub cpu: Option<MetricGroup>,
153    #[serde(default, skip_serializing_if = "Option::is_none")]
154    pub rpc: Option<MetricGroup>,
155    #[serde(flatten, default)]
156    pub extra: BTreeMap<String, serde_json::Value>,
157}
158
159impl MetricGroups {
160    /// Iterate current and future metric groups in deterministic name order.
161    pub fn groups(&self) -> Vec<(&str, &MetricGroup)> {
162        let mut groups = Vec::new();
163        for (name, group) in [
164            ("batch-jobs", self.batch_jobs.as_ref()),
165            ("cpu", self.cpu.as_ref()),
166            ("disk", self.disk.as_ref()),
167            ("memory", self.mem.as_ref()),
168            ("network", self.net.as_ref()),
169            ("os", self.os.as_ref()),
170            ("rpc", self.rpc.as_ref()),
171            ("scanner", self.scanner.as_ref()),
172            ("site-resync", self.site_resync.as_ref()),
173        ] {
174            if let Some(group) = group {
175                groups.push((name, group));
176            }
177        }
178        groups
179    }
180}
181
182/// One JSON Lines record from `/rustfs/admin/v3/metrics`.
183#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
184pub struct RealtimeMetrics {
185    #[serde(default)]
186    pub errors: Vec<String>,
187    #[serde(default)]
188    pub hosts: Vec<String>,
189    #[serde(default)]
190    pub aggregated: MetricGroups,
191    #[serde(rename = "by_host", default)]
192    pub by_host: BTreeMap<String, MetricGroups>,
193    #[serde(rename = "by_disk", default)]
194    pub by_disk: BTreeMap<String, MetricGroup>,
195    #[serde(rename = "final", default)]
196    pub final_sample: bool,
197    #[serde(flatten, default)]
198    pub extra: BTreeMap<String, serde_json::Value>,
199}
200
201/// Fully bounded result of a RustFS metrics query.
202#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
203pub struct MetricsBatch {
204    pub snapshots: Vec<RealtimeMetrics>,
205    pub encoded_bytes: usize,
206}
207
208impl MetricsBatch {
209    /// Whether RustFS reported incomplete data or omitted its final marker.
210    pub fn is_partial(&self) -> bool {
211        self.snapshots
212            .iter()
213            .any(|snapshot| !snapshot.errors.is_empty())
214            || self
215                .snapshots
216                .last()
217                .is_some_and(|snapshot| !snapshot.final_sample)
218    }
219}
220
221/// Scanner freshness metadata calculated by RustFS.
222#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
223pub struct ScannerFreshness {
224    #[serde(default)]
225    pub state: String,
226    #[serde(default)]
227    pub last_cycle_end_unix_secs: u64,
228    #[serde(default)]
229    pub max_expected_age_seconds: u64,
230    #[serde(default)]
231    pub reason: Option<String>,
232    #[serde(flatten, default)]
233    pub extra: BTreeMap<String, serde_json::Value>,
234}
235
236/// Typed operational fields from the scanner report, with future fields retained.
237#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
238pub struct ScannerMetrics {
239    #[serde(default)]
240    pub collected_at: String,
241    #[serde(default)]
242    pub current_cycle: u64,
243    #[serde(default)]
244    pub current_started: String,
245    #[serde(default)]
246    pub cycles_completed_at: Vec<String>,
247    #[serde(default)]
248    pub ongoing_buckets: usize,
249    #[serde(default)]
250    pub active_scan_paths: usize,
251    #[serde(default)]
252    pub oldest_active_path_age_seconds: u64,
253    #[serde(default)]
254    pub active_paths: Vec<String>,
255    #[serde(default)]
256    pub current_scan_mode: String,
257    #[serde(default)]
258    pub leader_lock_state: String,
259    #[serde(default)]
260    pub leader_lock_held_by_this_process: bool,
261    #[serde(default)]
262    pub leader_lock_last_error: String,
263    #[serde(default)]
264    pub last_cycle_end_unix_secs: u64,
265    #[serde(default)]
266    pub current_cycle_objects_scanned: u64,
267    #[serde(default)]
268    pub current_cycle_directories_scanned: u64,
269    #[serde(default)]
270    pub current_cycle_bucket_drive_failures: u64,
271    #[serde(default)]
272    pub last_cycle_result: String,
273    #[serde(default)]
274    pub last_cycle_result_code: u64,
275    #[serde(default)]
276    pub last_cycle_partial_reason: String,
277    #[serde(default)]
278    pub last_cycle_partial_source: String,
279    #[serde(default)]
280    pub last_cycle_duration_seconds: f64,
281    #[serde(default)]
282    pub last_cycle_objects_scanned: u64,
283    #[serde(default)]
284    pub last_cycle_directories_scanned: u64,
285    #[serde(default)]
286    pub last_cycle_bucket_drive_failures: u64,
287    #[serde(default)]
288    pub failed_cycles: u64,
289    #[serde(default)]
290    pub partial_cycles: u64,
291    #[serde(flatten, default)]
292    pub extra: BTreeMap<String, serde_json::Value>,
293}
294
295/// Effective scanner condition presented by the CLI.
296#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
297#[serde(rename_all = "kebab-case")]
298pub enum ScannerHealth {
299    Healthy,
300    Stale,
301    Empty,
302    Partial,
303    Disabled,
304}
305
306impl fmt::Display for ScannerHealth {
307    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
308        let value = match self {
309            Self::Healthy => "healthy",
310            Self::Stale => "stale",
311            Self::Empty => "empty",
312            Self::Partial => "partial",
313            Self::Disabled => "disabled",
314        };
315        formatter.write_str(value)
316    }
317}
318
319/// Effective scanner cycle scheduling information.
320#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
321pub struct ScannerCycleSchedule {
322    #[serde(default)]
323    pub effective_interval_seconds: u64,
324    #[serde(default)]
325    pub clean_idle_backoff_enabled: bool,
326    #[serde(default = "default_backoff_multiplier")]
327    pub clean_idle_backoff_multiplier: u64,
328    #[serde(flatten, default)]
329    pub extra: BTreeMap<String, serde_json::Value>,
330}
331
332const fn default_backoff_multiplier() -> u64 {
333    1
334}
335
336/// One scanner runtime setting and the source selected by RustFS.
337#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
338pub struct ScannerRuntimeConfigValue<T> {
339    pub value: T,
340    #[serde(default)]
341    pub source: String,
342}
343
344/// Typed subset of scanner runtime settings.
345#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
346pub struct ScannerRuntimeConfig {
347    #[serde(default)]
348    pub speed: Option<ScannerRuntimeConfigValue<String>>,
349    #[serde(default)]
350    pub delay: Option<ScannerRuntimeConfigValue<f64>>,
351    #[serde(default)]
352    pub max_wait_seconds: Option<ScannerRuntimeConfigValue<f64>>,
353    #[serde(default)]
354    pub idle_mode: Option<ScannerRuntimeConfigValue<bool>>,
355    #[serde(default)]
356    pub start_delay_seconds: Option<ScannerRuntimeConfigValue<Option<u64>>>,
357    #[serde(default)]
358    pub cycle_interval_seconds: Option<ScannerRuntimeConfigValue<u64>>,
359    #[serde(default)]
360    pub bitrot_cycle_seconds: Option<ScannerRuntimeConfigValue<Option<u64>>>,
361    #[serde(default)]
362    pub cycle_max_duration_seconds: Option<ScannerRuntimeConfigValue<Option<u64>>>,
363    #[serde(default)]
364    pub cycle_max_objects: Option<ScannerRuntimeConfigValue<Option<u64>>>,
365    #[serde(default)]
366    pub cycle_max_directories: Option<ScannerRuntimeConfigValue<Option<u64>>>,
367    #[serde(flatten, default)]
368    pub extra: BTreeMap<String, serde_json::Value>,
369}
370
371/// RustFS scanner status response.
372#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
373pub struct ScannerStatus {
374    #[serde(default)]
375    pub enabled: bool,
376    #[serde(default)]
377    pub disabled_reason: Option<String>,
378    #[serde(default)]
379    pub freshness: ScannerFreshness,
380    #[serde(default)]
381    pub metrics: ScannerMetrics,
382    #[serde(default)]
383    pub cycle_schedule: ScannerCycleSchedule,
384    #[serde(default)]
385    pub runtime_config: ScannerRuntimeConfig,
386    #[serde(flatten, default)]
387    pub extra: BTreeMap<String, serde_json::Value>,
388}
389
390impl ScannerStatus {
391    /// Classify operational state without treating stale or disabled data as transport failure.
392    pub fn health(&self) -> ScannerHealth {
393        if !self.enabled {
394            return ScannerHealth::Disabled;
395        }
396        if self.freshness.state.eq_ignore_ascii_case("stale") {
397            return ScannerHealth::Stale;
398        }
399        if self
400            .metrics
401            .last_cycle_result
402            .eq_ignore_ascii_case("partial")
403            || (!self.metrics.last_cycle_partial_reason.is_empty()
404                && !self
405                    .metrics
406                    .last_cycle_partial_reason
407                    .eq_ignore_ascii_case("unknown"))
408        {
409            return ScannerHealth::Partial;
410        }
411        if self.metrics.current_cycle == 0
412            && self.metrics.last_cycle_end_unix_secs == 0
413            && self.freshness.last_cycle_end_unix_secs == 0
414        {
415            return ScannerHealth::Empty;
416        }
417        ScannerHealth::Healthy
418    }
419}
420
421/// Per-operation disk counters from storage information.
422#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
423pub struct StorageDiskMetrics {
424    #[serde(default)]
425    pub last_minute: BTreeMap<String, serde_json::Value>,
426    #[serde(default)]
427    pub api_calls: BTreeMap<String, u64>,
428    #[serde(default)]
429    pub total_waiting: u32,
430    #[serde(default)]
431    pub total_errors_availability: u64,
432    #[serde(default)]
433    pub total_errors_timeout: u64,
434    #[serde(default)]
435    pub total_writes: u64,
436    #[serde(default)]
437    pub total_deletes: u64,
438}
439
440/// One disk returned by RustFS storage information.
441#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
442pub struct StorageDisk {
443    #[serde(default)]
444    pub endpoint: String,
445    #[serde(rename = "rootDisk", default)]
446    pub root_disk: bool,
447    #[serde(rename = "path", default)]
448    pub drive_path: String,
449    #[serde(default)]
450    pub healing: bool,
451    #[serde(default)]
452    pub scanning: bool,
453    #[serde(default)]
454    pub state: String,
455    #[serde(default)]
456    pub uuid: String,
457    #[serde(default)]
458    pub model: Option<String>,
459    #[serde(rename = "totalspace", default)]
460    pub total_space: u64,
461    #[serde(rename = "usedspace", default)]
462    pub used_space: u64,
463    #[serde(rename = "availspace", default)]
464    pub available_space: u64,
465    #[serde(rename = "readthroughput", default)]
466    pub read_throughput: Option<f64>,
467    #[serde(rename = "writethroughput", default)]
468    pub write_throughput: Option<f64>,
469    #[serde(rename = "readlatency", default)]
470    pub read_latency: Option<f64>,
471    #[serde(rename = "writelatency", default)]
472    pub write_latency: Option<f64>,
473    #[serde(default)]
474    pub utilization: Option<f64>,
475    #[serde(default)]
476    pub metrics: Option<StorageDiskMetrics>,
477    #[serde(default)]
478    pub heal_info: Option<serde_json::Value>,
479    #[serde(default)]
480    pub used_inodes: u64,
481    #[serde(default)]
482    pub free_inodes: u64,
483    #[serde(default)]
484    pub local: bool,
485    #[serde(default)]
486    pub pool_index: i32,
487    #[serde(default)]
488    pub set_index: i32,
489    #[serde(default)]
490    pub disk_index: i32,
491    #[serde(rename = "runtimeState", default)]
492    pub runtime_state: Option<String>,
493    #[serde(rename = "offlineDurationSeconds", default)]
494    pub offline_duration_seconds: Option<u64>,
495    #[serde(rename = "capacityObservationSource", default)]
496    pub capacity_observation_source: Option<String>,
497    #[serde(rename = "capacityObservationAgeSeconds", default)]
498    pub capacity_observation_age_seconds: Option<u64>,
499    #[serde(rename = "physicalDeviceIds", default)]
500    pub physical_device_ids: Option<Vec<String>>,
501    #[serde(flatten, default)]
502    pub extra: BTreeMap<String, serde_json::Value>,
503}
504
505/// RustFS storage backend kind.
506#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
507pub enum StorageBackendKind {
508    #[default]
509    Unknown,
510    #[serde(rename = "FS")]
511    Fs,
512    Erasure,
513    #[serde(other)]
514    Other,
515}
516
517impl fmt::Display for StorageBackendKind {
518    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
519        let value = match self {
520            Self::Unknown => "unknown",
521            Self::Fs => "filesystem",
522            Self::Erasure => "erasure",
523            Self::Other => "other",
524        };
525        formatter.write_str(value)
526    }
527}
528
529/// Storage backend topology and erasure coding configuration.
530#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
531pub struct StorageBackend {
532    #[serde(rename = "BackendType", alias = "backend_type", default)]
533    pub kind: StorageBackendKind,
534    #[serde(rename = "OnlineDisks", alias = "online_disks", default)]
535    pub online_disks: BTreeMap<String, usize>,
536    #[serde(rename = "OfflineDisks", alias = "offline_disks", default)]
537    pub offline_disks: BTreeMap<String, usize>,
538    #[serde(rename = "StandardSCData", default)]
539    pub standard_sc_data: Vec<usize>,
540    #[serde(rename = "StandardSCParities", default)]
541    pub standard_sc_parities: Vec<usize>,
542    #[serde(rename = "StandardSCParity", default)]
543    pub standard_sc_parity: Option<usize>,
544    #[serde(rename = "RRSCData", default)]
545    pub rr_sc_data: Vec<usize>,
546    #[serde(rename = "RRSCParities", default)]
547    pub rr_sc_parities: Vec<usize>,
548    #[serde(rename = "RRSCParity", default)]
549    pub rr_sc_parity: Option<usize>,
550    #[serde(rename = "TotalSets", alias = "total_sets", default)]
551    pub total_sets: Vec<usize>,
552    #[serde(rename = "DrivesPerSet", alias = "drives_per_set", default)]
553    pub drives_per_set: Vec<usize>,
554    #[serde(flatten, default)]
555    pub extra: BTreeMap<String, serde_json::Value>,
556}
557
558/// RustFS storage information response body.
559#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
560pub struct StorageInfo {
561    #[serde(default)]
562    pub disks: Vec<StorageDisk>,
563    #[serde(default)]
564    pub backend: StorageBackend,
565}
566
567impl StorageInfo {
568    pub fn total_capacity(&self) -> u64 {
569        self.disks.iter().map(|disk| disk.total_space).sum()
570    }
571
572    pub fn used_capacity(&self) -> u64 {
573        self.disks.iter().map(|disk| disk.used_space).sum()
574    }
575
576    pub fn online_disks(&self) -> usize {
577        self.disks
578            .iter()
579            .filter(|disk| {
580                disk.runtime_state
581                    .as_deref()
582                    .unwrap_or(&disk.state)
583                    .eq_ignore_ascii_case("online")
584                    || disk.state.eq_ignore_ascii_case("ok")
585            })
586            .count()
587    }
588}
589
590/// Read-only scanner, storage, and metrics Admin API boundary.
591#[async_trait]
592pub trait ObservabilityApi: Send + Sync {
593    async fn scanner_status(&self) -> Result<ScannerStatus>;
594    async fn storage_info(&self) -> Result<StorageInfo>;
595    async fn realtime_metrics(&self, query: &MetricsQuery) -> Result<MetricsBatch>;
596}
597
598#[cfg(test)]
599mod tests {
600    use super::*;
601
602    #[test]
603    fn metrics_scope_masks_match_rustfs_beta_10() {
604        assert_eq!(MetricsScope::Scanner.bit(), 1);
605        assert_eq!(MetricsScope::Disk.bit(), 2);
606        assert_eq!(MetricsScope::Os.bit(), 4);
607        assert_eq!(MetricsScope::BatchJobs.bit(), 8);
608        assert_eq!(MetricsScope::SiteResync.bit(), 16);
609        assert_eq!(MetricsScope::Network.bit(), 32);
610        assert_eq!(MetricsScope::Memory.bit(), 64);
611        assert_eq!(MetricsScope::Cpu.bit(), 128);
612        assert_eq!(MetricsScope::Rpc.bit(), 256);
613        assert_eq!(MetricsScope::All.bit(), 511);
614    }
615
616    #[test]
617    fn metrics_query_deduplicates_scopes_without_losing_bits() {
618        let query = MetricsQuery {
619            scopes: vec![
620                MetricsScope::Scanner,
621                MetricsScope::Disk,
622                MetricsScope::Scanner,
623            ],
624            ..Default::default()
625        };
626
627        assert_eq!(query.types_mask(), 3);
628    }
629
630    #[test]
631    fn scanner_health_distinguishes_healthy_stale_empty_partial_and_disabled() {
632        let healthy = scanner_fixture("fresh", "success", 7, true);
633        assert_eq!(healthy.health(), ScannerHealth::Healthy);
634
635        let stale = scanner_fixture("stale", "success", 7, true);
636        assert_eq!(stale.health(), ScannerHealth::Stale);
637
638        let empty = scanner_fixture("unknown", "unknown", 0, true);
639        assert_eq!(empty.health(), ScannerHealth::Empty);
640
641        let partial = scanner_fixture("fresh", "partial", 7, true);
642        assert_eq!(partial.health(), ScannerHealth::Partial);
643
644        let disabled = scanner_fixture("unknown", "unknown", 0, false);
645        assert_eq!(disabled.health(), ScannerHealth::Disabled);
646    }
647
648    #[test]
649    fn scanner_status_preserves_unknown_server_fields() {
650        let status: ScannerStatus = serde_json::from_str(
651            r#"{
652                "enabled":true,
653                "freshness":{"state":"fresh","last_cycle_end_unix_secs":10,"max_expected_age_seconds":120,"reason":null,"server_hint":"retained"},
654                "metrics":{"collected_at":"2026-07-21T04:00:00Z","current_cycle":2,"last_cycle_result":"success","future_counter":9},
655                "cycle_schedule":{"effective_interval_seconds":60,"clean_idle_backoff_enabled":false,"clean_idle_backoff_multiplier":1},
656                "runtime_config":{"speed":{"value":"fast","source":"default"},"future_setting":{"value":1,"source":"config"}}
657            }"#,
658        )
659        .expect("scanner status should deserialize");
660
661        assert_eq!(status.metrics.extra["future_counter"], 9);
662        assert_eq!(status.freshness.extra["server_hint"], "retained");
663        assert!(status.runtime_config.extra.contains_key("future_setting"));
664    }
665
666    #[test]
667    fn realtime_metrics_preserve_numeric_values_labels_and_timestamps() {
668        let snapshot: RealtimeMetrics = serde_json::from_str(
669            r#"{
670                "errors":[],
671                "hosts":["node-1"],
672                "aggregated":{"net":{"collected":"2026-07-21T04:00:00Z","netstats":{"rx_bytes":42}}},
673                "by_host":{"node-1":{"rpc":{"collectedAt":"2026-07-21T04:00:01Z","incomingBytes":7}}},
674                "by_disk":{"/data1":{"collected":"2026-07-21T04:00:02Z","n_disks":1}},
675                "final":true
676            }"#,
677        )
678        .expect("realtime metrics should deserialize");
679
680        assert_eq!(
681            snapshot.aggregated.net.expect("net group").0["netstats"]["rx_bytes"],
682            42
683        );
684        assert_eq!(
685            snapshot.by_host["node-1"]
686                .rpc
687                .as_ref()
688                .expect("rpc group")
689                .0["collectedAt"],
690            "2026-07-21T04:00:01Z"
691        );
692        assert_eq!(snapshot.by_disk["/data1"].0["n_disks"], 1);
693        assert!(snapshot.final_sample);
694    }
695
696    #[test]
697    fn storage_info_reads_current_rustfs_field_names() {
698        let info: StorageInfo = serde_json::from_str(
699            r#"{
700                "disks":[{
701                    "endpoint":"http://node1:9000",
702                    "path":"/data1",
703                    "state":"online",
704                    "totalspace":100,
705                    "usedspace":40,
706                    "availspace":60,
707                    "runtimeState":"online",
708                    "pool_index":0,
709                    "set_index":1,
710                    "disk_index":2
711                }],
712                "backend":{"BackendType":"Erasure","OnlineDisks":{"set-1":1},"OfflineDisks":{}}
713            }"#,
714        )
715        .expect("storage info should deserialize");
716
717        assert_eq!(info.disks[0].drive_path, "/data1");
718        assert_eq!(info.disks[0].used_space, 40);
719        assert_eq!(info.backend.kind, StorageBackendKind::Erasure);
720        assert_eq!(info.backend.online_disks.values().sum::<usize>(), 1);
721    }
722
723    #[test]
724    fn storage_info_preserves_missing_observations_as_unavailable() {
725        let info: StorageInfo = serde_json::from_str(
726            r#"{
727                "disks":[{
728                    "endpoint":"http://node1:9000",
729                    "path":"/data1",
730                    "state":"online",
731                    "totalspace":100,
732                    "usedspace":40,
733                    "availspace":60
734                }],
735                "backend":{"BackendType":"Erasure"}
736            }"#,
737        )
738        .expect("storage info should deserialize");
739
740        assert!(info.disks[0].read_throughput.is_none());
741        assert!(info.disks[0].write_throughput.is_none());
742        assert!(info.disks[0].read_latency.is_none());
743        assert!(info.disks[0].write_latency.is_none());
744        assert!(info.disks[0].utilization.is_none());
745        assert!(info.disks[0].metrics.is_none());
746    }
747
748    fn scanner_fixture(
749        freshness: &str,
750        last_result: &str,
751        current_cycle: u64,
752        enabled: bool,
753    ) -> ScannerStatus {
754        ScannerStatus {
755            enabled,
756            disabled_reason: (!enabled).then(|| "disabled by configuration".to_string()),
757            freshness: ScannerFreshness {
758                state: freshness.to_string(),
759                ..Default::default()
760            },
761            metrics: ScannerMetrics {
762                current_cycle,
763                last_cycle_result: last_result.to_string(),
764                ..Default::default()
765            },
766            ..Default::default()
767        }
768    }
769}