Skip to main content

ecat_data_iotdb/
lib.rs

1// Copyright (c) 2026 erik <erik@erik.xyz> — https://erik.xyz
2use async_trait::async_trait;
3use ecat_circuit_breaker::{Breaker, BreakerConfig};
4use ecat_data::{
5    BackendKind, DataPoint, FieldValue, TsdbClient, breaker_error_to_backend_error,
6    run_with_timeout,
7};
8use ecat_errors::{Error, ErrorCode};
9use ecat_tls::TlsClientConfig;
10use serde::Deserialize;
11use std::sync::Arc;
12use std::time::Duration;
13use tokio::sync::{Semaphore, SemaphorePermit};
14
15#[cfg(feature = "metrics")]
16mod metrics;
17#[cfg(feature = "metrics")]
18pub use metrics::register_outbound_metrics;
19
20#[derive(Debug, Clone, Deserialize)]
21pub struct IotdbConfig {
22    pub base_url: String,
23    pub username: String,
24    pub password: String,
25    #[serde(default)]
26    pub tls: Option<TlsClientConfig>,
27    /// 单次调用超时秒数。`0` = 禁用;未配置 = 30 秒。
28    ///
29    /// 这是**外层**预算,与 reqwest 自带的总超时(`from_config` 建的 client 有
30    /// 30 秒、`new` 没有)取先到者。
31    #[serde(default)]
32    pub query_timeout_secs: Option<u64>,
33    /// 熔断配置;省略则用保守默认(失败率 0.5、窗口 30 秒、打开 10 秒)。
34    ///
35    /// 熔断**默认开启** —— 保守阈值下只在持续失败时打开。**当前没有总开关**:
36    /// `BreakerConfig` 只有阈值字段,没有 `enabled`(不许写 `{"enabled": false}`:
37    /// 那是反序列化错误,或被 `#[serde(default)]` 静默吞掉后以为关掉了)。
38    /// 真要停用,只能把阈值调到不可能触发(如 `failure_ratio: 1.1`)。
39    #[serde(default)]
40    pub breaker: Option<BreakerConfig>,
41    /// 并发上限。`0` = **不限并发**(与 `query_timeout_secs: 0` = 禁用同构);
42    /// 未配置 = 32。
43    ///
44    /// reqwest **只有** `pool_max_idle_per_host`(空闲保留数),没有「最大总连接数」
45    /// —— 默认无上限意味着并发无背压。上限由本 crate 的信号量实现,不是 reqwest 的旋钮。
46    #[serde(default)]
47    pub max_concurrency: Option<usize>,
48}
49
50pub struct IotdbClient {
51    client: reqwest::Client,
52    base_url: String,
53    username: String,
54    password: String,
55    query_timeout: Option<Duration>,
56    /// 逐 client 一个 —— 熔断器要挂在**后端实例**上,不是进程上。
57    breaker: Arc<Breaker>,
58    /// `None` = 不限并发(`max_concurrency: 0`)。
59    semaphore: Option<Arc<Semaphore>>,
60}
61
62impl IotdbClient {
63    pub fn new(
64        base_url: impl Into<String>,
65        username: impl Into<String>,
66        password: impl Into<String>,
67    ) -> Self {
68        Self {
69            client: reqwest::Client::new(),
70            base_url: base_url.into(),
71            username: username.into(),
72            password: password.into(),
73            query_timeout: query_timeout(None),
74            breaker: Arc::new(Breaker::new(BreakerConfig::default())),
75            semaphore: Some(Arc::new(Semaphore::new(32))),
76        }
77    }
78
79    pub fn from_config(cfg: IotdbConfig) -> Result<Self, Error> {
80        let client = ecat_tls::build_reqwest_client(&cfg.tls)
81            .map_err(|e| Error::new(ErrorCode::Internal, "iotdb", format!("TLS: {e}")))?;
82        let breaker = Arc::new(Breaker::new(cfg.breaker.unwrap_or_default()));
83        // 自动接线(lead 裁决 2026-10-08):`metrics` feature 下**构造即注册**,
84        // 用户代码零变化。`ecat_metrics::register_outbound_metrics` 是幂等的
85        // (同一 backend 重复注册是覆盖闭包),所以多 client 不会炸。
86        // 探针:注释掉下面两行 ⇒ from_config_registers_outbound_metrics 红。
87        #[cfg(feature = "metrics")]
88        crate::register_outbound_metrics(Arc::clone(&breaker));
89        Ok(Self {
90            client,
91            base_url: cfg.base_url,
92            username: cfg.username,
93            password: cfg.password,
94            query_timeout: query_timeout(cfg.query_timeout_secs),
95            breaker,
96            semaphore: match cfg.max_concurrency {
97                // `0` = 不限并发(与 `query_timeout_secs: 0` = 禁用同构):
98                // 不建信号量。建 `Semaphore::new(0)` 会让每次调用静默无限挂起
99                // —— `guarded` 的第一句就是 `permit().await`,超时层在它里面。
100                Some(0) => None,
101                Some(n) => Some(Arc::new(Semaphore::new(n))),
102                None => Some(Arc::new(Semaphore::new(32))),
103            },
104        })
105    }
106
107    /// 本 client 的熔断器。`metrics` feature 注册指标时要读它的状态与打开次数。
108    pub fn breaker(&self) -> Arc<Breaker> {
109        Arc::clone(&self.breaker)
110    }
111
112    /// 取一个并发许可;不限并发(`max_concurrency: 0`)时返回 `None`。
113    /// 信号量从不 `close()`,`AcquireError` 不可达。
114    async fn permit(&self) -> Option<SemaphorePermit<'_>> {
115        match &self.semaphore {
116            Some(sem) => Some(sem.acquire().await.expect("semaphore is never closed")),
117            None => None,
118        }
119    }
120
121    /// 一次出站调用的公共外壳:**许可 → 熔断 → 超时**。
122    ///
123    /// 这个顺序不能改:**超时若在外层,熔断器会对卡死的后端永久失明** ——
124    /// 超时触发时 `tokio::time::timeout` 会 drop 内层 future,而熔断器记失败的那句
125    /// 在 `f().await` **之后**,于是每次都只留下一次 drop、窗口里什么都不记,
126    /// 熔断器永远不打开(这正是本设计要防的头号场景)。详见批次 5a 计划的「出入 4」。
127    ///
128    /// 许可在最外:还在排队的请求**还没碰后端**,不该计入熔断失败、也不该被超时掐断。
129    ///
130    /// `kind` **写死**不收参数:本 crate 的两个 I/O 方法同属 `TsdbClient` 一个家族
131    /// (ClickHouse 收参数是因为它有 `SqlExecutor` / `TsdbClient` 两条**不同家族**的路径
132    /// 共用外壳)。多一个永不变化的入参就多一个填错的机会,
133    /// 而填错只是静默少数(`ecat-data/src/timeout.rs:15-35`),没有编译期保护。
134    ///
135    /// **壳的边界 = 公开方法**:本 crate 的 `write` **循环里每次迭代都发 HTTP**,
136    /// 所以整个循环在这一个壳内 —— 谁把这段挪到循环里面(每个点一个预算),
137    /// `whole_call_budget_covers_every_request_in_write` 会红。判据:问「这一步发 HTTP 吗?」
138    async fn guarded<F, T: 'static>(&self, fut: F) -> Result<T, Error>
139    where
140        F: std::future::Future<Output = Result<T, Error>> + Send,
141    {
142        let _permit = self.permit().await;
143        self.breaker
144            .call(|| run_with_timeout(BackendKind::Tsdb, self.query_timeout, fut))
145            .await
146            .map_err(|e| breaker_error_to_backend_error(e, "iotdb"))
147    }
148}
149
150/// `0` 表示显式禁用超时;未配置时为 30 秒。
151fn query_timeout(secs: Option<u64>) -> Option<Duration> {
152    match secs {
153        None => Some(Duration::from_secs(30)),
154        Some(0) => None,
155        Some(s) => Some(Duration::from_secs(s)),
156    }
157}
158
159#[async_trait]
160impl TsdbClient for IotdbClient {
161    async fn write(&self, points: &[DataPoint]) -> Result<(), Error> {
162        // 整个循环**一个预算**:本方法每个点发一次 POST(`rest/v2/insertTablet`),
163        // 预算必须罩住整次调用 —— 包在循环里面就会变成「每个点一个预算」,
164        // 一次 write 的墙钟上限随点数线性放大,熔断窗口也被拆成 N 份记账
165        // (测试 whole_call_budget_covers_every_request_in_write 盯着这条)。
166        self.guarded(async {
167            for p in points {
168                // Apache IoTDB REST v2 insertTablet body:
169                // {"device": "...", "is_aligned": false, "timestamps": [...],
170                //  "measurements": [...], "data_types": [...], "values": [[...]]}
171                // `device` = measurement; tags are not representable in this API.
172                let mut measurements = Vec::with_capacity(p.fields.len());
173                let mut data_types = Vec::with_capacity(p.fields.len());
174                let mut values: Vec<serde_json::Value> = Vec::with_capacity(p.fields.len());
175                for (k, v) in &p.fields {
176                    measurements.push(k.clone());
177                    let (dt, val) = match v {
178                        FieldValue::Float(f) => (
179                            "DOUBLE",
180                            serde_json::Value::Number(
181                                serde_json::Number::from_f64(*f).unwrap_or(0.into()),
182                            ),
183                        ),
184                        FieldValue::Int(i) => ("INT64", serde_json::Value::Number((*i).into())),
185                        FieldValue::String(s) => ("TEXT", serde_json::Value::String(s.clone())),
186                        FieldValue::Bool(b) => ("BOOLEAN", serde_json::Value::Bool(*b)),
187                    };
188                    data_types.push(dt);
189                    values.push(val);
190                }
191                let body = serde_json::json!({
192                    "device": p.measurement,
193                    "is_aligned": false,
194                    "timestamps": [p.timestamp.unwrap_or(0)],
195                    "measurements": measurements,
196                    "data_types": data_types,
197                    "values": [values],
198                });
199                let resp = self
200                    .client
201                    .post(format!("{}/rest/v2/insertTablet", self.base_url))
202                    .basic_auth(&self.username, Some(&self.password))
203                    .header("Content-Type", "application/json")
204                    .json(&body)
205                    .send()
206                    .await
207                    .map_err(|e| {
208                        Error::new(ErrorCode::Internal, "iotdb", format!("iotdb write: {e}"))
209                    })?;
210                if !resp.status().is_success() {
211                    return Err(Error::new(
212                        ErrorCode::Internal,
213                        "iotdb",
214                        resp.text().await.unwrap_or_default(),
215                    ));
216                }
217                // IoTDB REST v2 may return HTTP 200 with a body `code` != 200 on
218                // some failures; surface those too.
219                if let Ok(v) = resp.json::<serde_json::Value>().await
220                    && let Some(code) = v.get("code").and_then(|c| c.as_i64())
221                    && code != 200
222                {
223                    return Err(Error::new(
224                        ErrorCode::Internal,
225                        "iotdb",
226                        format!(
227                            "iotdb write failed: code {code}: {}",
228                            v.get("message")
229                                .and_then(|m| m.as_str())
230                                .unwrap_or("no message")
231                        ),
232                    ));
233                }
234            }
235            Ok(())
236        })
237        .await
238    }
239
240    async fn query(&self, sql: &str) -> Result<serde_json::Value, Error> {
241        self.guarded(async {
242            let resp = self
243                .client
244                .post(format!("{}/rest/v2/query", self.base_url))
245                .basic_auth(&self.username, Some(&self.password))
246                .header("Content-Type", "text/plain; charset=utf-8")
247                .body(sql.to_string())
248                .send()
249                .await
250                .map_err(|e| {
251                    Error::new(ErrorCode::Internal, "iotdb", format!("iotdb query: {e}"))
252                })?;
253            if !resp.status().is_success() {
254                return Err(Error::new(
255                    ErrorCode::Internal,
256                    "iotdb",
257                    resp.text().await.unwrap_or_default(),
258                ));
259            }
260            resp.json()
261                .await
262                .map_err(|e| Error::new(ErrorCode::Internal, "iotdb", format!("iotdb parse: {e}")))
263        })
264        .await
265    }
266
267    // `delete` 走 `TsdbClient` 的 trait 默认实现(`ecat-data/src/tsdb.rs:55`),
268    // **不包 `guarded`**:默认实现的「不支持」是**调用方的用法错**,不是后端故障。
269    // 包了之后 8 次「不支持」就会打开熔断器,之后**正常写入/查询全被拒绝**。守测试见
270    // `tests/resilience.rs::delete_default_does_not_trip_the_breaker`。
271}
272
273#[cfg(test)]
274mod tests;