1use std::collections::HashMap;
16
17use chrono::{DateTime, Utc};
18use serde::{Deserialize, Serialize};
19
20use crate::health::MemInfo;
21
22#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
23pub struct TimedAction {
24 #[serde(rename = "count")]
25 pub count: u64,
26 #[serde(rename = "acc_time_ns")]
27 pub acc_time: u64,
28 #[serde(rename = "bytes")]
29 pub bytes: u64,
30}
31
32impl TimedAction {
33 pub fn merge(&mut self, other: &TimedAction) {
34 self.count += other.count;
35 self.acc_time += other.acc_time;
36 self.bytes += other.bytes;
37 }
38}
39
40#[derive(Clone, Debug, Default, Serialize, Deserialize)]
41pub struct DiskIOStats {
42 #[serde(rename = "read_ios")]
43 pub read_ios: u64,
44 #[serde(rename = "read_merges")]
45 pub read_merges: u64,
46 #[serde(rename = "read_sectors")]
47 pub read_sectors: u64,
48 #[serde(rename = "read_ticks")]
49 pub read_ticks: u64,
50 #[serde(rename = "write_ios")]
51 pub write_ios: u64,
52 #[serde(rename = "write_merges")]
53 pub write_merges: u64,
54 #[serde(rename = "write_sectors")]
55 pub write_sectors: u64,
56 #[serde(rename = "write_ticks")]
57 pub write_ticks: u64,
58 #[serde(rename = "current_ios")]
59 pub current_ios: u64,
60 #[serde(rename = "total_ticks")]
61 pub total_ticks: u64,
62 #[serde(rename = "req_ticks")]
63 pub req_ticks: u64,
64 #[serde(rename = "discard_ios")]
65 pub discard_ios: u64,
66 #[serde(rename = "discard_merges")]
67 pub discard_merges: u64,
68 #[serde(rename = "discard_secotrs")]
69 pub discard_sectors: u64,
70 #[serde(rename = "discard_ticks")]
71 pub discard_ticks: u64,
72 #[serde(rename = "flush_ios")]
73 pub flush_ios: u64,
74 #[serde(rename = "flush_ticks")]
75 pub flush_ticks: u64,
76}
77
78#[derive(Clone, Debug, Default, Serialize, Deserialize)]
79pub struct DiskMetric {
80 #[serde(rename = "collected")]
81 pub collected_at: DateTime<Utc>,
82 #[serde(rename = "n_disks")]
83 pub n_disks: usize,
84 #[serde(rename = "offline")]
85 pub offline: usize,
86 #[serde(rename = "healing")]
87 pub healing: usize,
88 #[serde(rename = "life_time_ops")]
89 pub life_time_ops: HashMap<String, u64>,
90 #[serde(rename = "last_minute")]
91 pub last_minute: Operations,
92 #[serde(rename = "iostats")]
93 pub io_stats: DiskIOStats,
94}
95
96impl DiskMetric {
97 pub fn merge(&mut self, other: &DiskMetric) {
98 if self.collected_at < other.collected_at {
99 self.collected_at = other.collected_at;
100 }
101 self.n_disks += other.n_disks;
102 self.offline += other.offline;
103 self.healing += other.healing;
104
105 for (k, v) in other.life_time_ops.iter() {
106 *self.life_time_ops.entry(k.clone()).or_insert(0) += v;
107 }
108
109 for (k, v) in other.last_minute.operations.iter() {
110 self.last_minute.operations.entry(k.clone()).or_default().merge(v);
111 }
112 }
113}
114
115#[derive(Clone, Debug, Default, Serialize, Deserialize)]
116pub struct LastMinute {
117 #[serde(rename = "actions")]
118 pub actions: HashMap<String, TimedAction>,
119 #[serde(rename = "ilm")]
120 pub ilm: HashMap<String, TimedAction>,
121}
122
123#[derive(Clone, Debug, Default, Serialize, Deserialize)]
124pub struct ScannerMetrics {
125 #[serde(rename = "collected")]
126 pub collected_at: DateTime<Utc>,
127 #[serde(rename = "current_cycle")]
128 pub current_cycle: u64,
129 #[serde(rename = "current_started")]
130 pub current_started: DateTime<Utc>,
131 #[serde(rename = "cycle_complete_times")]
132 pub cycles_completed_at: Vec<DateTime<Utc>>,
133 #[serde(rename = "ongoing_buckets")]
134 pub ongoing_buckets: usize,
135 #[serde(rename = "life_time_ops")]
136 pub life_time_ops: HashMap<String, u64>,
137 #[serde(rename = "ilm_ops")]
138 pub life_time_ilm: HashMap<String, u64>,
139 #[serde(rename = "last_minute")]
140 pub last_minute: LastMinute,
141 #[serde(rename = "active")]
142 pub active_paths: Vec<String>,
143}
144
145impl ScannerMetrics {
146 pub fn merge(&mut self, other: &Self) {
147 if self.collected_at < other.collected_at {
148 self.collected_at = other.collected_at;
149 }
150
151 if self.ongoing_buckets < other.ongoing_buckets {
152 self.ongoing_buckets = other.ongoing_buckets;
153 }
154
155 if self.current_cycle < other.current_cycle {
156 self.current_cycle = other.current_cycle;
157 self.cycles_completed_at = other.cycles_completed_at.clone();
158 self.current_started = other.current_started;
159 }
160
161 if other.cycles_completed_at.len() > self.cycles_completed_at.len() {
162 self.cycles_completed_at = other.cycles_completed_at.clone();
163 }
164
165 if !other.life_time_ops.is_empty() && self.life_time_ops.is_empty() {
166 self.life_time_ops = other.life_time_ops.clone();
167 }
168
169 for (k, v) in other.life_time_ops.iter() {
170 *self.life_time_ops.entry(k.clone()).or_default() += v;
171 }
172
173 for (k, v) in other.last_minute.actions.iter() {
174 self.last_minute.actions.entry(k.clone()).or_default().merge(v);
175 }
176
177 for (k, v) in other.life_time_ilm.iter() {
178 *self.life_time_ilm.entry(k.clone()).or_default() += v;
179 }
180
181 for (k, v) in other.last_minute.ilm.iter() {
182 self.last_minute.ilm.entry(k.clone()).or_default().merge(v);
183 }
184
185 self.active_paths.extend(other.active_paths.clone());
186
187 self.active_paths.sort();
188 }
189}
190
191#[derive(Clone, Debug, Default, Serialize, Deserialize)]
192pub struct Metrics {
193 #[serde(rename = "scanner", skip_serializing_if = "Option::is_none")]
194 pub scanner: Option<ScannerMetrics>,
195 #[serde(rename = "disk", skip_serializing_if = "Option::is_none")]
196 pub disk: Option<DiskMetric>,
197 #[serde(rename = "os", skip_serializing_if = "Option::is_none")]
198 pub os: Option<OsMetrics>,
199 #[serde(rename = "batchJobs", skip_serializing_if = "Option::is_none")]
200 pub batch_jobs: Option<BatchJobMetrics>,
201 #[serde(rename = "siteResync", skip_serializing_if = "Option::is_none")]
202 pub site_resync: Option<SiteResyncMetrics>,
203 #[serde(rename = "net", skip_serializing_if = "Option::is_none")]
204 pub net: Option<NetMetrics>,
205 #[serde(rename = "mem", skip_serializing_if = "Option::is_none")]
206 pub mem: Option<MemMetrics>,
207 #[serde(rename = "cpu", skip_serializing_if = "Option::is_none")]
208 pub cpu: Option<CPUMetrics>,
209 #[serde(rename = "rpc", skip_serializing_if = "Option::is_none")]
210 pub rpc: Option<RPCMetrics>,
211}
212
213impl Metrics {
214 pub fn merge(&mut self, other: &Self) {
215 if let Some(scanner) = other.scanner.as_ref() {
216 match self.scanner {
217 Some(ref mut s_scanner) => s_scanner.merge(scanner),
218 None => self.scanner = Some(scanner.clone()),
219 }
220 }
221
222 if let Some(disk) = other.disk.as_ref() {
223 match self.disk {
224 Some(ref mut s_disk) => s_disk.merge(disk),
225 None => self.disk = Some(disk.clone()),
226 }
227 }
228
229 if let Some(os) = other.os.as_ref() {
230 match self.os {
231 Some(ref mut s_os) => s_os.merge(os),
232 None => self.os = Some(os.clone()),
233 }
234 }
235
236 if let Some(batch_jobs) = other.batch_jobs.as_ref() {
237 match self.batch_jobs {
238 Some(ref mut s_batch_jobs) => s_batch_jobs.merge(batch_jobs),
239 None => self.batch_jobs = Some(batch_jobs.clone()),
240 }
241 }
242
243 if let Some(site_resync) = other.site_resync.as_ref() {
244 match self.site_resync {
245 Some(ref mut s_site_resync) => s_site_resync.merge(site_resync),
246 None => self.site_resync = Some(site_resync.clone()),
247 }
248 }
249
250 if let Some(net) = other.net.as_ref() {
251 match self.net {
252 Some(ref mut s_net) => s_net.merge(net),
253 None => self.net = Some(net.clone()),
254 }
255 }
256
257 if let Some(rpc) = other.rpc.as_ref() {
258 match self.rpc {
259 Some(ref mut s_rpc) => s_rpc.merge(rpc),
260 None => self.rpc = Some(rpc.clone()),
261 }
262 }
263 }
264}
265
266#[derive(Clone, Debug, Default, Serialize, Deserialize)]
267pub struct RPCMetrics {
268 #[serde(rename = "collectedAt")]
269 pub collected_at: DateTime<Utc>,
270
271 pub connected: i32,
272
273 #[serde(rename = "reconnectCount")]
274 pub reconnect_count: i32,
275
276 pub disconnected: i32,
277
278 #[serde(rename = "outgoingStreams")]
279 pub outgoing_streams: i32,
280
281 #[serde(rename = "incomingStreams")]
282 pub incoming_streams: i32,
283
284 #[serde(rename = "outgoingBytes")]
285 pub outgoing_bytes: i64,
286
287 #[serde(rename = "incomingBytes")]
288 pub incoming_bytes: i64,
289
290 #[serde(rename = "outgoingMessages")]
291 pub outgoing_messages: i64,
292
293 #[serde(rename = "incomingMessages")]
294 pub incoming_messages: i64,
295
296 pub out_queue: i32,
297
298 #[serde(rename = "lastPongTime")]
299 pub last_pong_time: DateTime<Utc>,
300
301 #[serde(rename = "lastPingMS")]
302 pub last_ping_ms: f64,
303
304 #[serde(rename = "maxPingDurMS")]
305 pub max_ping_dur_ms: f64, #[serde(rename = "lastConnectTime")]
308 pub last_connect_time: DateTime<Utc>,
309
310 #[serde(rename = "byDestination", skip_serializing_if = "Option::is_none")]
311 pub by_destination: Option<HashMap<String, RPCMetrics>>,
312
313 #[serde(rename = "byCaller", skip_serializing_if = "Option::is_none")]
314 pub by_caller: Option<HashMap<String, RPCMetrics>>,
315}
316
317impl RPCMetrics {
318 pub fn merge(&mut self, other: &Self) {
319 if self.collected_at < other.collected_at {
320 self.collected_at = other.collected_at;
321 }
322
323 if self.last_connect_time < other.last_connect_time {
324 self.last_connect_time = other.last_connect_time;
325 }
326
327 self.connected += other.connected;
328 self.disconnected += other.disconnected;
329 self.reconnect_count += other.reconnect_count;
330 self.outgoing_streams += other.outgoing_streams;
331 self.incoming_streams += other.incoming_streams;
332 self.outgoing_bytes += other.outgoing_bytes;
333 self.incoming_bytes += other.incoming_bytes;
334 self.outgoing_messages += other.outgoing_messages;
335 self.incoming_messages += other.incoming_messages;
336 self.out_queue += other.out_queue;
337
338 if self.last_pong_time < other.last_pong_time {
339 self.last_pong_time = other.last_pong_time;
340 self.last_ping_ms = other.last_ping_ms;
341 }
342
343 if self.max_ping_dur_ms < other.max_ping_dur_ms {
344 self.max_ping_dur_ms = other.max_ping_dur_ms;
345 }
346
347 if let Some(by_destination) = other.by_destination.as_ref() {
348 match self.by_destination.as_mut() {
349 Some(s_by_de) => {
350 for (key, value) in by_destination {
351 s_by_de
352 .entry(key.to_string())
353 .and_modify(|v| v.merge(value))
354 .or_insert(value.clone());
355 }
356 }
357 None => self.by_destination = Some(by_destination.clone()),
358 }
359 }
360
361 if let Some(by_caller) = other.by_caller.as_ref() {
362 match self.by_caller.as_mut() {
363 Some(s_by_caller) => {
364 for (key, value) in by_caller {
365 s_by_caller
366 .entry(key.to_string())
367 .and_modify(|v| v.merge(value))
368 .or_insert(value.clone());
369 }
370 }
371 None => self.by_caller = Some(by_caller.clone()),
372 }
373 }
374 }
375}
376
377#[derive(Clone, Debug, Default, Serialize, Deserialize)]
378pub struct CPUMetrics {}
379
380#[derive(Clone, Debug, Default, Serialize, Deserialize)]
381pub struct NetMetrics {
382 #[serde(rename = "collected")]
383 pub collected_at: DateTime<Utc>,
384 #[serde(rename = "interfaceName")]
385 pub interface_name: String,
386 #[serde(rename = "netstats")]
387 pub net_stats: NetDevLine,
388}
389
390impl NetMetrics {
391 pub fn merge(&mut self, other: &Self) {
392 if self.collected_at < other.collected_at {
393 self.collected_at = other.collected_at;
394 }
395
396 self.net_stats.rx_bytes += other.net_stats.rx_bytes;
397 self.net_stats.rx_packets += other.net_stats.rx_packets;
398 self.net_stats.rx_errors += other.net_stats.rx_errors;
399 self.net_stats.rx_dropped += other.net_stats.rx_dropped;
400 self.net_stats.rx_fifo += other.net_stats.rx_fifo;
401 self.net_stats.rx_frame += other.net_stats.rx_frame;
402 self.net_stats.rx_compressed += other.net_stats.rx_compressed;
403 self.net_stats.rx_multicast += other.net_stats.rx_multicast;
404 self.net_stats.tx_bytes += other.net_stats.tx_bytes;
405 self.net_stats.tx_packets += other.net_stats.tx_packets;
406 self.net_stats.tx_errors += other.net_stats.tx_errors;
407 self.net_stats.tx_dropped += other.net_stats.tx_dropped;
408 self.net_stats.tx_fifo += other.net_stats.tx_fifo;
409 self.net_stats.tx_collisions += other.net_stats.tx_collisions;
410 self.net_stats.tx_carrier += other.net_stats.tx_carrier;
411 self.net_stats.tx_compressed += other.net_stats.tx_compressed;
412 }
413}
414
415#[derive(Clone, Debug, Default, Serialize, Deserialize)]
416pub struct NetDevLine {
417 #[serde(rename = "name")]
418 pub name: String, #[serde(rename = "rx_bytes")]
421 pub rx_bytes: u64, #[serde(rename = "rx_packets")]
424 pub rx_packets: u64, #[serde(rename = "rx_errors")]
427 pub rx_errors: u64, #[serde(rename = "rx_dropped")]
430 pub rx_dropped: u64, #[serde(rename = "rx_fifo")]
433 pub rx_fifo: u64, #[serde(rename = "rx_frame")]
436 pub rx_frame: u64, #[serde(rename = "rx_compressed")]
439 pub rx_compressed: u64, #[serde(rename = "rx_multicast")]
442 pub rx_multicast: u64, #[serde(rename = "tx_bytes")]
445 pub tx_bytes: u64, #[serde(rename = "tx_packets")]
448 pub tx_packets: u64, #[serde(rename = "tx_errors")]
451 pub tx_errors: u64, #[serde(rename = "tx_dropped")]
454 pub tx_dropped: u64, #[serde(rename = "tx_fifo")]
457 pub tx_fifo: u64, #[serde(rename = "tx_collisions")]
460 pub tx_collisions: u64, #[serde(rename = "tx_carrier")]
463 pub tx_carrier: u64, #[serde(rename = "tx_compressed")]
466 pub tx_compressed: u64, }
468
469#[derive(Clone, Debug, Default, Serialize, Deserialize)]
470pub struct MemMetrics {
471 #[serde(rename = "collected")]
472 pub collected_at: DateTime<Utc>,
473 #[serde(rename = "memInfo")]
474 pub info: MemInfo,
475}
476
477#[derive(Clone, Debug, Default, Serialize, Deserialize)]
478pub struct SiteResyncMetrics {
479 #[serde(rename = "collected")]
480 pub collected_at: DateTime<Utc>,
481 #[serde(rename = "resyncStatus", skip_serializing_if = "Option::is_none")]
482 pub resync_status: Option<String>,
483 #[serde(rename = "startTime")]
484 pub start_time: DateTime<Utc>,
485 #[serde(rename = "lastUpdate")]
486 pub last_update: DateTime<Utc>,
487 #[serde(rename = "numBuckets")]
488 pub num_buckets: i64,
489 #[serde(rename = "resyncID")]
490 pub resync_id: String,
491 #[serde(rename = "deplID")]
492 pub depl_id: String,
493 #[serde(rename = "completedReplicationSize")]
494 pub replicated_size: i64,
495 #[serde(rename = "replicationCount")]
496 pub replicated_count: i64,
497 #[serde(rename = "failedReplicationSize")]
498 pub failed_size: i64,
499 #[serde(rename = "failedReplicationCount")]
500 pub failed_count: i64,
501 #[serde(rename = "failedBuckets")]
502 pub failed_buckets: Vec<String>,
503 #[serde(rename = "bucket", skip_serializing_if = "Option::is_none")]
504 pub bucket: Option<String>,
505 #[serde(rename = "object", skip_serializing_if = "Option::is_none")]
506 pub object: Option<String>,
507}
508
509impl SiteResyncMetrics {
510 pub fn merge(&mut self, other: &Self) {
511 if self.collected_at < other.collected_at {
512 *self = other.clone();
513 }
514 }
515}
516
517#[derive(Clone, Debug, Default, Serialize, Deserialize)]
518pub struct BatchJobMetrics {
519 #[serde(rename = "collected")]
520 pub collected_at: DateTime<Utc>,
521 #[serde(rename = "Jobs")]
522 pub jobs: HashMap<String, JobMetric>,
523}
524
525impl BatchJobMetrics {
526 pub fn merge(&mut self, other: &BatchJobMetrics) {
527 if other.jobs.is_empty() {
528 return;
529 }
530
531 if self.collected_at < other.collected_at {
532 self.collected_at = other.collected_at;
533 }
534
535 for (k, v) in other.jobs.clone().into_iter() {
536 self.jobs.insert(k, v);
537 }
538 }
539}
540
541#[derive(Clone, Debug, Default, Serialize, Deserialize)]
542pub struct JobMetric {
543 #[serde(rename = "jobID")]
544 pub job_id: String,
545 #[serde(rename = "jobType")]
546 pub job_type: String,
547 #[serde(rename = "startTime")]
548 pub start_time: DateTime<Utc>,
549 #[serde(rename = "lastUpdate")]
550 pub last_update: DateTime<Utc>,
551 #[serde(rename = "retryAttempts")]
552 pub retry_attempts: i32,
553 pub complete: bool,
554 pub failed: bool,
555 #[serde(skip_serializing_if = "Option::is_none")]
557 pub replicate: Option<ReplicateInfo>,
558 #[serde(skip_serializing_if = "Option::is_none")]
559 pub key_rotate: Option<KeyRotationInfo>,
560 #[serde(skip_serializing_if = "Option::is_none")]
561 pub expired: Option<ExpirationInfo>,
562}
563
564#[derive(Clone, Debug, Default, Serialize, Deserialize)]
565pub struct ReplicateInfo {
566 #[serde(rename = "lastBucket")]
567 pub bucket: String,
568 #[serde(rename = "lastObject")]
569 pub object: String,
570 #[serde(rename = "objects")]
571 pub objects: i64,
572 #[serde(rename = "objectsFailed")]
573 pub objects_failed: i64,
574 #[serde(rename = "bytesTransferred")]
575 pub bytes_transferred: i64,
576 #[serde(rename = "bytesFailed")]
577 pub bytes_failed: i64,
578}
579
580#[derive(Clone, Debug, Default, Serialize, Deserialize)]
581pub struct ExpirationInfo {
582 #[serde(rename = "lastBucket")]
583 pub bucket: String,
584 #[serde(rename = "lastObject")]
585 pub object: String,
586 #[serde(rename = "objects")]
587 pub objects: i64,
588 #[serde(rename = "objectsFailed")]
589 pub objects_failed: i64,
590}
591
592#[derive(Clone, Debug, Default, Serialize, Deserialize)]
593pub struct KeyRotationInfo {
594 #[serde(rename = "lastBucket")]
595 pub bucket: String,
596 #[serde(rename = "lastObject")]
597 pub object: String,
598 #[serde(rename = "objects")]
599 pub objects: i64,
600 #[serde(rename = "objectsFailed")]
601 pub objects_failed: i64,
602}
603
604#[derive(Debug, Default, Serialize, Deserialize)]
605pub struct RealtimeMetrics {
606 #[serde(rename = "errors")]
607 pub errors: Vec<String>,
608 #[serde(rename = "hosts")]
609 pub hosts: Vec<String>,
610 #[serde(rename = "aggregated")]
611 pub aggregated: Metrics,
612 #[serde(rename = "by_host")]
613 pub by_host: HashMap<String, Metrics>,
614 #[serde(rename = "by_disk")]
615 pub by_disk: HashMap<String, DiskMetric>,
616 #[serde(rename = "final")]
617 pub finally: bool,
618}
619
620impl RealtimeMetrics {
621 pub fn merge(&mut self, other: Self) {
622 if !other.errors.is_empty() {
623 self.errors.extend(other.errors);
624 }
625
626 for (k, v) in other.by_host.into_iter() {
627 *self.by_host.entry(k).or_default() = v;
628 }
629
630 self.hosts.extend(other.hosts);
631 self.aggregated.merge(&other.aggregated);
632 self.hosts.sort();
633
634 for (k, v) in other.by_disk.into_iter() {
635 self.by_disk.entry(k.to_string()).and_modify(|h| *h = v.clone()).or_insert(v);
636 }
637 }
638}
639
640#[derive(Clone, Debug, Default, Serialize, Deserialize)]
641pub struct OsMetrics {
642 #[serde(rename = "collected")]
643 pub collected_at: DateTime<Utc>,
644 #[serde(rename = "life_time_ops")]
645 pub life_time_ops: HashMap<String, u64>,
646 #[serde(rename = "last_minute")]
647 pub last_minute: Operations,
648}
649
650impl OsMetrics {
651 pub fn merge(&mut self, other: &Self) {
652 if self.collected_at < other.collected_at {
653 self.collected_at = other.collected_at;
654 }
655
656 for (k, v) in other.life_time_ops.iter() {
657 *self.life_time_ops.entry(k.clone()).or_default() += v;
658 }
659
660 for (k, v) in other.last_minute.operations.iter() {
661 self.last_minute.operations.entry(k.clone()).or_default().merge(v);
662 }
663 }
664}
665
666#[derive(Clone, Debug, Default, Serialize, Deserialize)]
667pub struct Operations {
668 #[serde(rename = "operations")]
669 pub operations: HashMap<String, TimedAction>,
670}