Skip to main content

ecat_data/
tsdb.rs

1// Copyright (c) 2026 erik <erik@erik.xyz> — https://erik.xyz
2use 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    /// Delete data using a backend-specific query (e.g. `DELETE FROM ...`).
54    /// Backends that cannot delete return an error.
55    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}