nodedb_types/timeseries/
ingest.rs1use std::collections::HashMap;
6
7use serde::{Deserialize, Serialize};
8
9#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
11pub struct MetricSample {
12 pub timestamp_ms: i64,
13 pub value: f64,
14}
15
16#[derive(Debug, Clone, Serialize, Deserialize)]
18pub struct LogEntry {
19 pub timestamp_ms: i64,
20 pub data: Vec<u8>,
21}
22
23#[derive(Debug, Clone, Copy, PartialEq, Eq)]
25#[non_exhaustive]
26pub enum IngestResult {
27 Ok,
29 FlushNeeded,
31 Rejected,
33}
34
35impl IngestResult {
36 pub fn is_flush_needed(&self) -> bool {
37 matches!(self, Self::FlushNeeded)
38 }
39
40 pub fn is_rejected(&self) -> bool {
41 matches!(self, Self::Rejected)
42 }
43}
44
45#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
47pub struct TimeRange {
48 pub start_ms: i64,
49 pub end_ms: i64,
50}
51
52impl TimeRange {
53 pub fn new(start_ms: i64, end_ms: i64) -> Self {
54 Self { start_ms, end_ms }
55 }
56
57 pub fn contains(&self, ts: i64) -> bool {
58 ts >= self.start_ms && ts <= self.end_ms
59 }
60
61 pub fn overlaps(&self, other: &TimeRange) -> bool {
63 self.start_ms <= other.end_ms && other.start_ms <= self.end_ms
64 }
65}
66
67#[derive(
72 Debug,
73 Clone,
74 Default,
75 PartialEq,
76 Serialize,
77 Deserialize,
78 zerompk::ToMessagePack,
79 zerompk::FromMessagePack,
80)]
81pub struct SymbolDictionary {
82 forward: HashMap<String, u32>,
84 reverse: Vec<String>,
86}
87
88impl SymbolDictionary {
89 pub fn new() -> Self {
90 Self::default()
91 }
92
93 pub fn resolve(&mut self, value: &str, max_cardinality: u32) -> Option<u32> {
97 if let Some(&id) = self.forward.get(value) {
98 return Some(id);
99 }
100 if self.reverse.len() as u32 >= max_cardinality {
101 return None;
102 }
103 let id = self.reverse.len() as u32;
104 self.forward.insert(value.to_string(), id);
105 self.reverse.push(value.to_string());
106 Some(id)
107 }
108
109 pub fn get(&self, id: u32) -> Option<&str> {
111 self.reverse.get(id as usize).map(|s| s.as_str())
112 }
113
114 pub fn get_id(&self, value: &str) -> Option<u32> {
116 self.forward.get(value).copied()
117 }
118
119 pub fn len(&self) -> usize {
121 self.reverse.len()
122 }
123
124 pub fn is_empty(&self) -> bool {
125 self.reverse.is_empty()
126 }
127
128 pub fn merge(&mut self, other: &SymbolDictionary, max_cardinality: u32) -> Vec<u32> {
132 let mut remap = Vec::with_capacity(other.reverse.len());
133 for symbol in &other.reverse {
134 match self.resolve(symbol, max_cardinality) {
135 Some(new_id) => remap.push(new_id),
136 None => remap.push(u32::MAX),
137 }
138 }
139 remap
140 }
141}