Skip to main content

ecat_data/
timeout.rs

1// Copyright (c) 2026 erik <erik@erik.xyz> — https://erik.xyz
2use crate::rdbms::RdbmsError;
3use ecat_errors::{Error, ErrorCode};
4use std::sync::atomic::{AtomicU64, Ordering};
5use std::time::Duration;
6
7/// 未提交即 Drop 的事务累计数(`Transaction` 的 Drop guard 递增)。
8pub static TRANSACTIONS_LEAKED: AtomicU64 = AtomicU64::new(0);
9
10/// 出站调用的后端类别,用于把超时计数分维度 ——
11/// 否则只有一个总数,看不出是哪个后端在出事。
12///
13/// 判别值即 [`TIMEOUTS`] 的下标,**顺序不可改**。
14///
15/// # 选哪个维度:跟 trait 家族走,不跟产品品类走
16///
17/// 一个产品可以同时属于两个家族:ClickHouse 既实现 `SqlExecutor` / `RdbmsClient`
18/// (错误类型 [`RdbmsError`])又实现 `TsdbClient`(错误类型 [`Error`])。
19/// 判据是**调用点包着哪个 trait**:
20///
21/// | 调用点包的 trait | `kind` |
22/// |---|---|
23/// | `SqlExecutor` / `RdbmsClient` | [`Rdbms`](BackendKind::Rdbms) |
24/// | `Cache` | [`Cache`](BackendKind::Cache) |
25/// | `SearchClient` | [`Search`](BackendKind::Search) |
26/// | `GraphClient` | [`Graph`](BackendKind::Graph) |
27/// | `DocumentClient` | [`Document`](BackendKind::Document) |
28/// | `StorageClient` | [`Storage`](BackendKind::Storage) |
29/// | `TsdbClient` | [`Tsdb`](BackendKind::Tsdb) |
30///
31/// 只有 `Rdbms` 维度被 `ecat-data-sqlx` / `ecat-data-mssql` 的 `metrics` 读成
32/// `ecat_rdbms_query_timeout_total`:走 `SqlExecutor` 的调用点若错填 `Tsdb`,
33/// 该指标会**静默少数**。`E: TimeoutError` 与 `kind` 是互相独立的两个轴
34/// (`Tsdb` + [`RdbmsError`] 照样编译),所以**没有编译期保护** —— 只能靠这条
35/// 判据加评审。
36#[derive(Debug, Clone, Copy, PartialEq, Eq)]
37pub enum BackendKind {
38    Rdbms = 0,
39    Cache = 1,
40    Search = 2,
41    Graph = 3,
42    Document = 4,
43    Storage = 5,
44    Tsdb = 6,
45}
46
47impl BackendKind {
48    /// `Error::reason` 用的后端短标识,与 [`TIMEOUTS`] 的维度一一对应。
49    ///
50    /// 粒度取家族(7 维)而非产品名:计数只到家族一级,reason 也到这一级,
51    /// 否则同一个槽会散出 `"redis"` / `"cache"` 两种 reason。
52    pub fn slug(&self) -> &'static str {
53        match self {
54            BackendKind::Rdbms => "rdbms",
55            BackendKind::Cache => "cache",
56            BackendKind::Search => "search",
57            BackendKind::Graph => "graph",
58            BackendKind::Document => "document",
59            BackendKind::Storage => "storage",
60            BackendKind::Tsdb => "tsdb",
61        }
62    }
63}
64
65/// 各类后端的累计超时次数,下标即 [`BackendKind`] 的判别值。
66///
67/// 用进程级静态量而非依赖 `ecat-metrics`:本 crate 保持零外部依赖,
68/// 指标的读取方按需接入(各后端的 `metrics` feature)。
69///
70/// 本量取代了 5.0.0 的 `QUERY_TIMEOUTS`(只覆盖 RDBMS)。
71///
72/// 取槽用 [`timeout_counter`],别自己写 `TIMEOUTS[kind as usize]`。
73pub static TIMEOUTS: [AtomicU64; 7] = [const { AtomicU64::new(0) }; 7];
74
75/// 取某类后端的超时计数器。
76///
77/// 内部是**无 `_` 分支**的 `match`,不用 `TIMEOUTS[kind as usize]`:
78/// 加第 8 个变体时这里直接编译失败,而不是静默编译通过、运行期数错槽。
79pub fn timeout_counter(kind: BackendKind) -> &'static AtomicU64 {
80    match kind {
81        BackendKind::Rdbms => &TIMEOUTS[0],
82        BackendKind::Cache => &TIMEOUTS[1],
83        BackendKind::Search => &TIMEOUTS[2],
84        BackendKind::Graph => &TIMEOUTS[3],
85        BackendKind::Document => &TIMEOUTS[4],
86        BackendKind::Storage => &TIMEOUTS[5],
87        BackendKind::Tsdb => &TIMEOUTS[6],
88    }
89}
90
91/// 可被超时包装的错误类型。
92///
93/// 六个非 RDBMS trait(`Cache` / `SearchClient` / `GraphClient` /
94/// `DocumentClient` / `StorageClient` / `TsdbClient`)共用 `ecat_errors::Error`,
95/// RDBMS 路径用 [`RdbmsError`] —— 两个实现就够,不需要六套。
96pub trait TimeoutError: Sized {
97    /// `kind` 只被 `Error` 那个实现用来填 `reason`([`BackendKind::slug`]);
98    /// [`RdbmsError::Timeout`] 没有 reason 字段,忽略它。
99    fn from_timeout(kind: BackendKind, d: Duration) -> Self;
100}
101
102impl TimeoutError for RdbmsError {
103    fn from_timeout(_kind: BackendKind, d: Duration) -> Self {
104        RdbmsError::Timeout(format!("query exceeded {d:?}"))
105    }
106}
107
108impl TimeoutError for Error {
109    fn from_timeout(kind: BackendKind, d: Duration) -> Self {
110        // `reason` 按全仓约定放**组件标识**;「超时」已由 `DeadlineExceeded`
111        // 表达,再写 `"timeout"` 会让按 reason 过滤/告警的人漏掉全部超时。
112        Error::new(
113            ErrorCode::DeadlineExceeded,
114            kind.slug(),
115            format!("call exceeded {d:?}"),
116        )
117    }
118}
119
120/// 给一次出站调用套一层超时。
121///
122/// `None` 表示禁用超时,直接透传结果。超时发生时递增 [`timeout_counter`] 的对应
123/// 维度,并返回 `E::from_timeout`。
124///
125/// **`None` 才是禁用;`Some(Duration::ZERO)` 是「立刻超时」**,不是禁用 ——
126/// 本仓别处「`0` = 禁用」的惯例(如 `ecat-data-sqlx/src/config.rs` 的
127/// `query_timeout_secs`)在这里不适用,别照着填。tokio 先 poll 内层:已就绪的
128/// future 仍会成功,只有挂着的才当场判超时。
129///
130/// 这是池耗尽的头号防线:`acquire_timeout` 只约束「等连接」,
131/// 拿到连接后卡死的查询会一直占着它。
132pub async fn run_with_timeout<F, T, E>(
133    kind: BackendKind,
134    timeout: Option<Duration>,
135    fut: F,
136) -> Result<T, E>
137where
138    F: std::future::Future<Output = Result<T, E>>,
139    E: TimeoutError,
140{
141    match timeout {
142        None => fut.await,
143        Some(d) => match tokio::time::timeout(d, fut).await {
144            Ok(result) => result,
145            Err(_) => {
146                timeout_counter(kind).fetch_add(1, Ordering::Relaxed);
147                Err(E::from_timeout(kind, d))
148            }
149        },
150    }
151}
152
153#[cfg(test)]
154mod tests {
155    use super::*;
156    use crate::rdbms::RdbmsError;
157    use std::sync::atomic::Ordering;
158
159    /// 该维度的当前值。测试用 `SeqCst`:别的测试也在并发地加同一个槽。
160    fn count(kind: BackendKind) -> u64 {
161        timeout_counter(kind).load(Ordering::SeqCst)
162    }
163
164    /// 钉住 7 个下标与判别值一致。
165    ///
166    /// [`timeout_counter`] 的无 `_` 分支 `match` 只保证**编译期穷尽**:把臂写反
167    /// (如 `Cache` 接到 `TIMEOUTS[2]`)照样编译通过,而且是这里唯一可能的编辑
168    /// 错误。别的用例只覆盖 Tsdb 与 Storage 两个槽,Cache↔Search 写反了没人会
169    /// 发现 —— 这个映射要被后续各后端的 `metrics` 依赖,按槽钉死。
170    #[test]
171    fn timeout_counter_indexes_match_discriminants() {
172        for k in [
173            BackendKind::Rdbms,
174            BackendKind::Cache,
175            BackendKind::Search,
176            BackendKind::Graph,
177            BackendKind::Document,
178            BackendKind::Storage,
179            BackendKind::Tsdb,
180        ] {
181            assert!(
182                std::ptr::eq(timeout_counter(k), &TIMEOUTS[k as usize]),
183                "timeout_counter({k:?}) 必须与 TIMEOUTS[{}] 是同一个槽",
184                k as usize
185            );
186        }
187    }
188
189    #[tokio::test]
190    async fn none_timeout_passes_result_through() {
191        let r: Result<u64, RdbmsError> =
192            run_with_timeout(BackendKind::Rdbms, None, async { Ok(42) }).await;
193        assert_eq!(r.unwrap(), 42);
194    }
195
196    #[tokio::test]
197    async fn fast_future_completes_within_timeout() {
198        let r: Result<u64, RdbmsError> =
199            run_with_timeout(BackendKind::Rdbms, Some(Duration::from_secs(5)), async {
200                Ok(1)
201            })
202            .await;
203        assert_eq!(r.unwrap(), 1);
204    }
205
206    /// **前提:同二进制里只有本用例写 `Rdbms` 槽。** 精确断言(`== before + 1`)
207    /// 依赖它 —— 别处若有用例写 `Rdbms`,它的 `+1` 会插进 `before` 与断言之间
208    /// (libtest 并行跑用例)。新增写 `Rdbms` 的用例时:改成方向断言(`> before`),
209    /// 或与本案串行。
210    #[tokio::test]
211    async fn slow_future_times_out_and_counts() {
212        let before = count(BackendKind::Rdbms);
213        let r: Result<(), RdbmsError> =
214            run_with_timeout(BackendKind::Rdbms, Some(Duration::from_millis(10)), async {
215                tokio::time::sleep(Duration::from_millis(200)).await;
216                Ok(())
217            })
218            .await;
219        let err = r.unwrap_err();
220        assert!(matches!(err, RdbmsError::Timeout(_)), "got: {err:?}");
221        assert_eq!(count(BackendKind::Rdbms), before + 1);
222    }
223
224    /// `ecat_errors::Error` 路径:必须是 `DeadlineExceeded` 而不是笼统的 `Internal`,
225    /// 否则调用方没法把「超时」与「后端报错」分开处理。`reason` 还要是**组件标识**
226    /// (`BackendKind::slug`)—— 全仓约定 reason 放组件名,写 `"timeout"` 会让按
227    /// reason 过滤/告警的人漏掉全部超时。
228    #[tokio::test]
229    async fn error_path_maps_to_deadline_exceeded() {
230        let r: Result<(), Error> =
231            run_with_timeout(BackendKind::Cache, Some(Duration::from_millis(10)), async {
232                tokio::time::sleep(Duration::from_millis(200)).await;
233                Ok(())
234            })
235            .await;
236        let err = r.unwrap_err();
237        assert_eq!(err.code, ErrorCode::DeadlineExceeded, "got: {err:?}");
238        assert_eq!(err.reason, "cache", "reason 应为组件标识,got: {err:?}");
239    }
240
241    /// 分维度计数**互不串**(spec §8 判据 4)。
242    /// 两个维度都真开火,而不是只断言「某维度 += 1」—— 后者在没有分维度时也会过。
243    ///
244    /// 观测槽取 `Storage` 而非 `Rdbms`:`slow_future_times_out_and_counts`
245    /// 同样在等 10ms 后推进 `Rdbms` 槽,两条用例被 libtest 并发调度时窗口重叠,
246    /// 拿 `Rdbms` 当观测槽是竞态(实测 12/15 次误红)。`Storage` 无其它写者,
247    /// 断言才确定;若实现把超时计到别的槽,这一条仍会红。
248    ///
249    /// **约束:证人槽必须是同二进制内没有任何其它用例会写的槽**,否则同一个竞态
250    /// 复发(写它的用例会把 `+1` 插进 `witness_before` 与断言之间)。加用例前
251    /// `grep -rn 'BackendKind::Storage'` 确认无人写它;要写就先换证人槽。
252    #[tokio::test]
253    async fn timeout_counters_are_per_backend_kind() {
254        let witness_before = count(BackendKind::Storage);
255        let tsdb_before = count(BackendKind::Tsdb);
256        let slow = || async {
257            tokio::time::sleep(Duration::from_millis(200)).await;
258            Ok::<(), Error>(())
259        };
260        let _: Result<(), Error> =
261            run_with_timeout(BackendKind::Tsdb, Some(Duration::from_millis(10)), slow()).await;
262        assert_eq!(count(BackendKind::Tsdb), tsdb_before + 1);
263        assert_eq!(
264            count(BackendKind::Storage),
265            witness_before,
266            "Tsdb 的超时不得落到别的槽"
267        );
268    }
269}