Skip to main content

helix_driver_host/
timer.rs

1//! TimerRegistry — ScheduleTimer / CancelTimer 兑现器
2//!
3//! 维护 TimerId → JoinHandle 映射。
4//! 到期时通过 mpsc channel 把 Tick::Timer{id} 发给 engine_loop。
5//!
6//! ## 设计要点
7//!
8//! - `schedule()` 先 cancel 同 id 的旧 timer(ScheduleTimer 可覆盖,幂等)
9//! - `cancel()` abort JoinHandle + 移除(幂等,id 不存在时静默)
10//! - JoinHandle 不 await:fire-and-forget spawn,timer 到期自行发送 Tick
11//! - **one-shot 回收契约**:fire-once 后既不 re-arm 也不发 `CancelTimer` 的一次性
12//!   timer(如 IM send-timeout),其 JoinHandle 会自然到期退出但不会自我移除——
13//!   长运行会话每触发一次就漏一条已完成条目。修法:spawn 闭包在 send 完成后经
14//!   `done_tx` 回灌自己的 id,`reap()` 在 schedule/cancel/pending_count 三个低频点
15//!   惰性把已完成 id 从 handles 移除(不引入后台 task / 不阻塞热路径)。
16//!   语义上与 core 侧 one-shot 概念对齐(core 改 `timer_map` remove-on-fire,本文件
17//!   只管 host 侧 `handles` map 回收)——二者**不共享文件、不互相依赖、仅语义对齐**。
18
19use helix_core::Tick;
20use helix_core::TimerId;
21use std::collections::HashMap;
22use std::sync::atomic::{AtomicUsize, Ordering};
23use std::sync::Arc;
24use std::time::Instant;
25use tokio::task::JoinHandle;
26
27use crate::metrics::{
28    AsyncMetricSink, LabelKey, MetricEvent, MetricId, MetricLabels, NoopMetricSink,
29};
30use crate::tick_ingress::EngineTickSender;
31
32pub struct TimerRegistry {
33    handles: HashMap<TimerId, JoinHandle<()>>,
34    /// one-shot 自清理回灌通道:spawn 闭包 fire 完毕后把自身 id 送回,`reap()` 据此惰性回收。
35    /// unbounded 对齐本 crate「回灌必达走 unbounded」不变量(避免回灌阻塞产生背压死锁)。
36    done_tx: tokio::sync::mpsc::UnboundedSender<TimerId>,
37    done_rx: tokio::sync::mpsc::UnboundedReceiver<TimerId>,
38    metrics: Arc<dyn AsyncMetricSink>,
39    active: Arc<AtomicUsize>,
40}
41
42impl Default for TimerRegistry {
43    fn default() -> Self {
44        Self::new()
45    }
46}
47
48impl TimerRegistry {
49    pub fn new() -> Self {
50        Self::with_metrics(Arc::new(NoopMetricSink))
51    }
52
53    /// 构造带 host 指标的 timer 注册表;sink 只允许非阻塞 try_record。
54    pub(crate) fn with_metrics(metrics: Arc<dyn AsyncMetricSink>) -> Self {
55        let (done_tx, done_rx) = tokio::sync::mpsc::unbounded_channel();
56        let registry = Self {
57            handles: HashMap::new(),
58            done_tx,
59            done_rx,
60            metrics,
61            active: Arc::new(AtomicUsize::new(0)),
62        };
63        record_timer_gauge(registry.metrics.as_ref(), 0);
64        registry
65    }
66
67    /// 惰性回收已自然到期的 one-shot timer 条目:非阻塞排空 done 通道,把已完成 id
68    /// 从 `handles` 移除。`try_recv` 空即停(O(已完成数),不阻塞、不在热路径)。
69    ///
70    /// **reschedule 同 id 竞态护栏**:done 通道里可能滞留「旧 id 已 fire」的回灌,而此后
71    /// schedule 用同 id 重新 arm 插入了新 handle。本函数仅做 `handles.remove(&id)` 排空,
72    /// 调用约定保证 **reap 永远先于本次操作对 handles 的写入**(schedule/cancel 内 reap
73    /// 在 cancel/insert 之前跑)——这样新 insert 的活 handle 不会被同批 reap 误删。
74    fn reap(&mut self) {
75        while let Ok(id) = self.done_rx.try_recv() {
76            self.handles.remove(&id);
77        }
78    }
79
80    /// 调度 timer:after_ms 毫秒后把 `Tick::Timer{id}` 发到 tick_tx。
81    ///
82    /// 若同 id 的 timer 已存在,先取消(ScheduleTimer 覆盖语义)。
83    ///
84    /// `tick_tx` 是 engine_loop 的入站通道,timer 触发与其他 Tick 源(WS/HTTP/Command)
85    /// 统一汇入同一 mpsc,由串行事件泵顺序处理。
86    pub(crate) fn schedule(&mut self, id: TimerId, after_ms: u64, tick_tx: EngineTickSender) {
87        // 先惰性回收已自然到期的 one-shot 条目(必须先于本次对 handles 的写入,
88        // 防止下方 insert 的新 handle 被同批旧 done 信号误删——见 reap 护栏)。
89        self.reap();
90        // 先取消已有同 id timer(ScheduleTimer 可覆盖)
91        self.cancel_with_reason(id, "replaced");
92
93        let done_tx = self.done_tx.clone();
94        let metrics = Arc::clone(&self.metrics);
95        let active = Arc::clone(&self.active);
96        let due_at = Instant::now() + tokio::time::Duration::from_millis(after_ms);
97        let current = active.fetch_add(1, Ordering::Relaxed) + 1;
98        record_timer_gauge(metrics.as_ref(), current);
99        record_timer_counter(metrics.as_ref(), MetricId::TimerScheduledTotal, "scheduled");
100        let handle = tokio::spawn(async move {
101            tokio::time::sleep(tokio::time::Duration::from_millis(after_ms)).await;
102            record_timer_histogram(
103                metrics.as_ref(),
104                MetricId::TimerLatenessSeconds,
105                Instant::now().saturating_duration_since(due_at),
106            );
107            let delivery_started = Instant::now();
108            // send 失败(接收端已关闭)静默忽略(正常关机场景)
109            let delivered = tick_tx.send(Tick::Timer(id)).await.is_ok();
110            record_timer_histogram(
111                metrics.as_ref(),
112                MetricId::TimerDeliveryWaitSeconds,
113                delivery_started.elapsed(),
114            );
115            if delivered {
116                record_timer_counter(metrics.as_ref(), MetricId::TimerFiredTotal, "delivered");
117            } else {
118                record_timer_counter(
119                    metrics.as_ref(),
120                    MetricId::TimerDeliveryFailedTotal,
121                    "closed",
122                );
123            }
124            let current = decrement_saturating(active.as_ref());
125            record_timer_gauge(metrics.as_ref(), current);
126            // fire 完毕,回灌自身 id 供 registry 惰性回收;registry 已 drop 时静默忽略。
127            done_tx.send(id).ok();
128        });
129        self.handles.insert(id, handle);
130    }
131
132    /// 取消 timer(幂等,id 不存在时静默忽略)。
133    pub fn cancel(&mut self, id: TimerId) {
134        self.cancel_with_reason(id, "explicit");
135    }
136
137    /// 取消指定 timer,并把覆盖与显式取消区分为稳定 status 标签。
138    fn cancel_with_reason(&mut self, id: TimerId, reason: &'static str) {
139        // 先惰性回收(必须先于本次 remove 对 handles 的读写,与 schedule 同护栏)。
140        self.reap();
141        if let Some(h) = self.handles.remove(&id) {
142            if h.is_finished() {
143                return;
144            }
145            h.abort();
146            let current = decrement_saturating(self.active.as_ref());
147            record_timer_gauge(self.metrics.as_ref(), current);
148            record_timer_counter(self.metrics.as_ref(), MetricId::TimerCancelledTotal, reason);
149        }
150    }
151
152    /// 当前挂起的 timer 数量(测试用)
153    #[cfg(test)]
154    pub fn pending_count(&mut self) -> usize {
155        // 先回收已自然到期的 one-shot 条目再报数,否则已 fire 的条目仍滞留计数。
156        self.reap();
157        self.handles.len()
158    }
159}
160
161/// Drop 时 abort 所有挂起 timer 任务:把「`run_engine_loop` 返回 = 完全 drain」的
162/// 保证扩展到 timer。`JoinHandle::drop` 仅 detach 不 abort——否则长 timer 任务会泄漏到
163/// runtime 关停,且持 tick_tx clone 延迟 channel 关闭。与 `cancel` 的 abort 语义一致。
164impl Drop for TimerRegistry {
165    fn drop(&mut self) {
166        for (_, h) in self.handles.drain() {
167            h.abort();
168        }
169    }
170}
171
172/// 记录 timer 低基数 Counter,不携带 TimerId。
173fn record_timer_counter(metrics: &dyn AsyncMetricSink, id: MetricId, status: &'static str) {
174    if metrics.is_enabled() {
175        let _ = metrics.try_record(MetricEvent::counter(
176            id,
177            1.0,
178            MetricLabels::one(LabelKey::Stage, "timer").with(LabelKey::Status, status),
179        ));
180    }
181}
182
183/// 记录 timer 耗时直方图,禁止把 timer payload 带入标签。
184fn record_timer_histogram(metrics: &dyn AsyncMetricSink, id: MetricId, value: std::time::Duration) {
185    if metrics.is_enabled() {
186        let _ = metrics.try_record(MetricEvent::histogram(
187            id,
188            value.as_secs_f64(),
189            MetricLabels::one(LabelKey::Stage, "timer"),
190        ));
191    }
192}
193
194/// 发布当前活跃 timer 数量。
195fn record_timer_gauge(metrics: &dyn AsyncMetricSink, active: usize) {
196    if metrics.is_enabled() {
197        let _ = metrics.try_record(MetricEvent::gauge(
198            MetricId::TimerActive,
199            active as f64,
200            MetricLabels::one(LabelKey::Stage, "timer"),
201        ));
202    }
203}
204
205/// 原子饱和递减,处理 cancel 与自然 fire 的竞争边界。
206fn decrement_saturating(value: &AtomicUsize) -> usize {
207    value
208        .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
209            Some(current.saturating_sub(1))
210        })
211        .unwrap_or_default()
212        .saturating_sub(1)
213}
214
215// ─── 单元测试 ────────────────────────────────────────────────────────────────
216
217#[cfg(test)]
218mod tests {
219    use super::*;
220    use crate::metrics::RecordOutcome;
221    use helix_core::effect::TimerId;
222    use std::sync::Mutex;
223    use tokio::sync::mpsc;
224
225    #[derive(Default)]
226    struct RecordingMetricSink(Mutex<Vec<MetricEvent>>);
227
228    impl AsyncMetricSink for RecordingMetricSink {
229        /// 保存 timer 外部指标事件,不启动真实 exporter。
230        fn try_record(&self, event: MetricEvent) -> RecordOutcome {
231            self.0.lock().unwrap().push(event);
232            RecordOutcome::Accepted
233        }
234    }
235
236    #[tokio::test]
237    async fn test_timer_fires_after_delay() {
238        let (tx, mut rx) = mpsc::channel::<Tick>(8);
239        let mut registry = TimerRegistry::new();
240
241        let id = TimerId::from_raw(42);
242        registry.schedule(id, 10, EngineTickSender::Raw(tx)); // 10ms 后触发
243
244        let tick = tokio::time::timeout(tokio::time::Duration::from_millis(200), rx.recv())
245            .await
246            .expect("timer should fire within 200ms")
247            .expect("channel should not be closed");
248
249        assert!(matches!(tick, Tick::Timer(tid) if tid == id));
250    }
251
252    #[tokio::test]
253    async fn test_cancel_prevents_fire() {
254        let (tx, mut rx) = mpsc::channel::<Tick>(8);
255        let mut registry = TimerRegistry::new();
256
257        let id = TimerId::from_raw(99);
258        // 10s timer — 立即 cancel 后,在 50ms 窗口内不应收到 Tick::Timer
259        registry.schedule(id, 10_000, EngineTickSender::Raw(tx.clone())); // clone tx 保持通道开放
260        registry.cancel(id);
261        // tx 仍然 alive(clone 持有),rx.recv() 会阻塞(而非立即返回 None)
262        // 因此 timeout 50ms 后应超时(Err),而非收到消息(Ok)
263
264        let result = tokio::time::timeout(tokio::time::Duration::from_millis(50), rx.recv()).await;
265
266        // timeout elapsed = Err(Elapsed),表示 timer 未触发(符合预期)
267        assert!(
268            result.is_err(),
269            "cancelled timer should not fire within 50ms, result was ready"
270        );
271
272        // 显式 drop tx(避免 rx 挂起)
273        drop(tx);
274    }
275
276    #[tokio::test]
277    async fn test_reschedule_overwrites() {
278        let (tx, mut rx) = mpsc::channel::<Tick>(8);
279        let mut registry = TimerRegistry::new();
280
281        let id = TimerId::from_raw(7);
282        registry.schedule(id, 10_000, EngineTickSender::Raw(tx.clone())); // 10s timer
283        registry.schedule(id, 10, EngineTickSender::Raw(tx)); // 10ms 覆盖
284
285        let tick = tokio::time::timeout(tokio::time::Duration::from_millis(200), rx.recv())
286            .await
287            .expect("reschedule should fire fast")
288            .expect("channel open");
289
290        assert!(matches!(tick, Tick::Timer(tid) if tid == id));
291    }
292
293    #[tokio::test]
294    async fn test_cancel_idempotent() {
295        let mut registry = TimerRegistry::new();
296        let id = TimerId::from_raw(1);
297        // 取消不存在的 id 不 panic
298        registry.cancel(id);
299        registry.cancel(id);
300    }
301
302    /// one-shot timer 自然到期(收到 Tick::Timer)后,handles 不残留该 id(防缓慢内存泄漏)。
303    #[tokio::test]
304    async fn test_oneshot_timer_reaped_after_fire() {
305        let (tx, mut rx) = mpsc::channel::<Tick>(8);
306        let mut registry = TimerRegistry::new();
307
308        let id = TimerId::from_raw(123);
309        registry.schedule(id, 10, EngineTickSender::Raw(tx)); // 10ms 后触发的一次性 timer
310
311        // 等待 timer 自然触发
312        let tick = tokio::time::timeout(tokio::time::Duration::from_millis(200), rx.recv())
313            .await
314            .expect("timer should fire within 200ms")
315            .expect("channel should not be closed");
316        assert!(matches!(tick, Tick::Timer(tid) if tid == id));
317
318        // 闭包在 send 之后才回灌 done_tx;让出执行权确保回灌已入通道,
319        // 随后 reap(pending_count 内部触发)把已完成 id 从 handles 移除。
320        tokio::task::yield_now().await;
321
322        // pending_count() 内部调 reap(),应观察到该 one-shot 已被回收 → 归零。
323        assert_eq!(
324            registry.pending_count(),
325            0,
326            "one-shot timer should be reaped from handles after natural fire"
327        );
328    }
329
330    /// timer 自然触发必须闭合 scheduled、fired、lateness、delivery 与 active Gauge。
331    #[tokio::test]
332    async fn observed_timer_records_lifecycle_without_timer_id_label() {
333        let metrics = Arc::new(RecordingMetricSink::default());
334        let (tx, mut rx) = mpsc::channel::<Tick>(8);
335        let mut registry = TimerRegistry::with_metrics(metrics.clone());
336
337        registry.schedule(TimerId::from_raw(321), 1, EngineTickSender::Raw(tx));
338        tokio::time::timeout(tokio::time::Duration::from_millis(200), rx.recv())
339            .await
340            .expect("timer 应在窗口内触发")
341            .expect("timer channel 应保持打开");
342        tokio::task::yield_now().await;
343
344        let events = metrics.0.lock().unwrap();
345        for expected in [
346            MetricId::TimerScheduledTotal,
347            MetricId::TimerFiredTotal,
348            MetricId::TimerLatenessSeconds,
349            MetricId::TimerDeliveryWaitSeconds,
350            MetricId::TimerActive,
351        ] {
352            assert!(
353                events.iter().any(|event| event.id == expected),
354                "timer 缺少 {:?}",
355                expected
356            );
357        }
358        assert!(events
359            .iter()
360            .filter(|event| event.id == MetricId::TimerActive)
361            .any(|event| event.value == 0.0));
362    }
363}