1use async_trait::async_trait;
3use ecat_errors::{Error, ErrorCode};
4use std::collections::HashMap;
5
6#[derive(Debug, Clone)]
7pub enum FieldValue {
8 Float(f64),
9 Int(i64),
10 String(String),
11 Bool(bool),
12}
13
14#[derive(Debug, Clone)]
15pub struct DataPoint {
16 pub measurement: String,
17 pub tags: HashMap<String, String>,
18 pub fields: HashMap<String, FieldValue>,
19 pub timestamp: Option<i64>,
20}
21
22impl DataPoint {
23 pub fn new(measurement: impl Into<String>) -> Self {
24 Self {
25 measurement: measurement.into(),
26 tags: HashMap::new(),
27 fields: HashMap::new(),
28 timestamp: None,
29 }
30 }
31
32 pub fn with_tag(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
33 self.tags.insert(key.into(), value.into());
34 self
35 }
36
37 pub fn with_field(mut self, key: impl Into<String>, value: FieldValue) -> Self {
38 self.fields.insert(key.into(), value);
39 self
40 }
41
42 pub fn with_timestamp(mut self, ts: i64) -> Self {
43 self.timestamp = Some(ts);
44 self
45 }
46}
47
48#[async_trait]
49pub trait TsdbClient: Send + Sync {
50 async fn write(&self, points: &[DataPoint]) -> Result<(), Error>;
51 async fn query(&self, query: &str) -> Result<serde_json::Value, Error>;
52
53 async fn delete(&self, _query: &str) -> Result<(), Error> {
56 Err(Error::new(
57 ErrorCode::Internal,
58 "tsdb",
59 "delete not supported by this backend",
60 ))
61 }
62}
63
64#[cfg(test)]
65mod tests {
66 use super::*;
67
68 #[test]
69 fn datapoint_builder_accumulates_metadata() {
70 let dp = DataPoint::new("cpu")
71 .with_tag("host", "h1")
72 .with_tag("env", "prod")
73 .with_field("usage", FieldValue::Float(0.85))
74 .with_field("count", FieldValue::Int(3))
75 .with_field("active", FieldValue::Bool(true))
76 .with_field("name", FieldValue::String("web".into()))
77 .with_timestamp(1_700_000_000_000);
78 assert_eq!(dp.measurement, "cpu");
79 assert_eq!(dp.tags.len(), 2);
80 assert_eq!(dp.tags.get("host").map(String::as_str), Some("h1"));
81 assert_eq!(dp.fields.len(), 4);
82 assert!(
83 matches!(dp.fields.get("usage"), Some(FieldValue::Float(v)) if (v - 0.85).abs() < 1e-9)
84 );
85 assert!(matches!(dp.fields.get("count"), Some(FieldValue::Int(3))));
86 assert!(matches!(
87 dp.fields.get("active"),
88 Some(FieldValue::Bool(true))
89 ));
90 assert!(matches!(
91 dp.fields.get("name"),
92 Some(FieldValue::String(s)) if s == "web"
93 ));
94 assert_eq!(dp.timestamp, Some(1_700_000_000_000));
95 }
96
97 #[test]
98 fn datapoint_defaults_are_empty() {
99 let dp = DataPoint::new("mem");
100 assert!(dp.tags.is_empty());
101 assert!(dp.fields.is_empty());
102 assert_eq!(dp.timestamp, None);
103 }
104
105 #[test]
106 fn datapoint_overwrites_existing_tag() {
107 let dp = DataPoint::new("cpu")
108 .with_tag("host", "a")
109 .with_tag("host", "b");
110 assert_eq!(dp.tags.get("host").map(String::as_str), Some("b"));
111 assert_eq!(dp.tags.len(), 1);
112 }
113
114 struct NoDeleteClient;
115
116 #[async_trait]
117 impl TsdbClient for NoDeleteClient {
118 async fn write(&self, _points: &[DataPoint]) -> Result<(), Error> {
119 Ok(())
120 }
121 async fn query(&self, _query: &str) -> Result<serde_json::Value, Error> {
122 Ok(serde_json::Value::Null)
123 }
124 }
125
126 #[tokio::test]
127 async fn delete_defaults_to_not_supported_error() {
128 let client = NoDeleteClient;
129 let err = client.delete("DROP TABLE t").await.unwrap_err();
130 assert!(
131 err.to_string().contains("delete not supported"),
132 "unexpected error: {err}"
133 );
134 }
135}