1use std::collections::BTreeMap;
4use std::fmt;
5
6use async_trait::async_trait;
7use serde::{Deserialize, Serialize};
8
9use crate::Result;
10
11pub const MAX_METRICS_SAMPLES: u16 = 120;
13pub const MAX_METRICS_LINE_BYTES: usize = 1024 * 1024;
15pub const MAX_METRICS_RESPONSE_BYTES: usize = 16 * 1024 * 1024;
17
18#[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 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#[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 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 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#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
127#[serde(transparent)]
128pub struct MetricGroup(pub BTreeMap<String, serde_json::Value>);
129
130#[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 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#[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#[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 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#[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#[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#[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#[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#[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#[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#[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 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#[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#[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#[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#[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#[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#[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}