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}