helix-driver-host 0.1.13

Helix Native 与 FFI 共用的存储、网络和执行驱动
Documentation
//! TimerRegistry — ScheduleTimer / CancelTimer 兑现器
//!
//! 维护 TimerId → JoinHandle 映射。
//! 到期时通过 mpsc channel 把 Tick::Timer{id} 发给 engine_loop。
//!
//! ## 设计要点
//!
//! - `schedule()` 先 cancel 同 id 的旧 timer(ScheduleTimer 可覆盖,幂等)
//! - `cancel()` abort JoinHandle + 移除(幂等,id 不存在时静默)
//! - JoinHandle 不 await:fire-and-forget spawn,timer 到期自行发送 Tick
//! - **one-shot 回收契约**:fire-once 后既不 re-arm 也不发 `CancelTimer` 的一次性
//!   timer(如 IM send-timeout),其 JoinHandle 会自然到期退出但不会自我移除——
//!   长运行会话每触发一次就漏一条已完成条目。修法:spawn 闭包在 send 完成后经
//!   `done_tx` 回灌自己的 id,`reap()` 在 schedule/cancel/pending_count 三个低频点
//!   惰性把已完成 id 从 handles 移除(不引入后台 task / 不阻塞热路径)。
//!   语义上与 core 侧 one-shot 概念对齐(core 改 `timer_map` remove-on-fire,本文件
//!   只管 host 侧 `handles` map 回收)——二者**不共享文件、不互相依赖、仅语义对齐**。

use helix_core::Tick;
use helix_core::TimerId;
use std::collections::HashMap;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::Instant;
use tokio::task::JoinHandle;

use crate::metrics::{
    AsyncMetricSink, LabelKey, MetricEvent, MetricId, MetricLabels, NoopMetricSink,
};
use crate::tick_ingress::EngineTickSender;

pub struct TimerRegistry {
    handles: HashMap<TimerId, JoinHandle<()>>,
    /// one-shot 自清理回灌通道:spawn 闭包 fire 完毕后把自身 id 送回,`reap()` 据此惰性回收。
    /// unbounded 对齐本 crate「回灌必达走 unbounded」不变量(避免回灌阻塞产生背压死锁)。
    done_tx: tokio::sync::mpsc::UnboundedSender<TimerId>,
    done_rx: tokio::sync::mpsc::UnboundedReceiver<TimerId>,
    metrics: Arc<dyn AsyncMetricSink>,
    active: Arc<AtomicUsize>,
}

impl Default for TimerRegistry {
    fn default() -> Self {
        Self::new()
    }
}

impl TimerRegistry {
    pub fn new() -> Self {
        Self::with_metrics(Arc::new(NoopMetricSink))
    }

    /// 构造带 host 指标的 timer 注册表;sink 只允许非阻塞 try_record。
    pub(crate) fn with_metrics(metrics: Arc<dyn AsyncMetricSink>) -> Self {
        let (done_tx, done_rx) = tokio::sync::mpsc::unbounded_channel();
        let registry = Self {
            handles: HashMap::new(),
            done_tx,
            done_rx,
            metrics,
            active: Arc::new(AtomicUsize::new(0)),
        };
        record_timer_gauge(registry.metrics.as_ref(), 0);
        registry
    }

    /// 惰性回收已自然到期的 one-shot timer 条目:非阻塞排空 done 通道,把已完成 id
    /// 从 `handles` 移除。`try_recv` 空即停(O(已完成数),不阻塞、不在热路径)。
    ///
    /// **reschedule 同 id 竞态护栏**:done 通道里可能滞留「旧 id 已 fire」的回灌,而此后
    /// schedule 用同 id 重新 arm 插入了新 handle。本函数仅做 `handles.remove(&id)` 排空,
    /// 调用约定保证 **reap 永远先于本次操作对 handles 的写入**(schedule/cancel 内 reap
    /// 在 cancel/insert 之前跑)——这样新 insert 的活 handle 不会被同批 reap 误删。
    fn reap(&mut self) {
        while let Ok(id) = self.done_rx.try_recv() {
            self.handles.remove(&id);
        }
    }

    /// 调度 timer:after_ms 毫秒后把 `Tick::Timer{id}` 发到 tick_tx。
    ///
    /// 若同 id 的 timer 已存在,先取消(ScheduleTimer 覆盖语义)。
    ///
    /// `tick_tx` 是 engine_loop 的入站通道,timer 触发与其他 Tick 源(WS/HTTP/Command)
    /// 统一汇入同一 mpsc,由串行事件泵顺序处理。
    pub(crate) fn schedule(&mut self, id: TimerId, after_ms: u64, tick_tx: EngineTickSender) {
        // 先惰性回收已自然到期的 one-shot 条目(必须先于本次对 handles 的写入,
        // 防止下方 insert 的新 handle 被同批旧 done 信号误删——见 reap 护栏)。
        self.reap();
        // 先取消已有同 id timer(ScheduleTimer 可覆盖)
        self.cancel_with_reason(id, "replaced");

        let done_tx = self.done_tx.clone();
        let metrics = Arc::clone(&self.metrics);
        let active = Arc::clone(&self.active);
        let due_at = Instant::now() + tokio::time::Duration::from_millis(after_ms);
        let current = active.fetch_add(1, Ordering::Relaxed) + 1;
        record_timer_gauge(metrics.as_ref(), current);
        record_timer_counter(metrics.as_ref(), MetricId::TimerScheduledTotal, "scheduled");
        let handle = tokio::spawn(async move {
            tokio::time::sleep(tokio::time::Duration::from_millis(after_ms)).await;
            record_timer_histogram(
                metrics.as_ref(),
                MetricId::TimerLatenessSeconds,
                Instant::now().saturating_duration_since(due_at),
            );
            let delivery_started = Instant::now();
            // send 失败(接收端已关闭)静默忽略(正常关机场景)
            let delivered = tick_tx.send(Tick::Timer(id)).await.is_ok();
            record_timer_histogram(
                metrics.as_ref(),
                MetricId::TimerDeliveryWaitSeconds,
                delivery_started.elapsed(),
            );
            if delivered {
                record_timer_counter(metrics.as_ref(), MetricId::TimerFiredTotal, "delivered");
            } else {
                record_timer_counter(
                    metrics.as_ref(),
                    MetricId::TimerDeliveryFailedTotal,
                    "closed",
                );
            }
            let current = decrement_saturating(active.as_ref());
            record_timer_gauge(metrics.as_ref(), current);
            // fire 完毕,回灌自身 id 供 registry 惰性回收;registry 已 drop 时静默忽略。
            done_tx.send(id).ok();
        });
        self.handles.insert(id, handle);
    }

    /// 取消 timer(幂等,id 不存在时静默忽略)。
    pub fn cancel(&mut self, id: TimerId) {
        self.cancel_with_reason(id, "explicit");
    }

    /// 取消指定 timer,并把覆盖与显式取消区分为稳定 status 标签。
    fn cancel_with_reason(&mut self, id: TimerId, reason: &'static str) {
        // 先惰性回收(必须先于本次 remove 对 handles 的读写,与 schedule 同护栏)。
        self.reap();
        if let Some(h) = self.handles.remove(&id) {
            if h.is_finished() {
                return;
            }
            h.abort();
            let current = decrement_saturating(self.active.as_ref());
            record_timer_gauge(self.metrics.as_ref(), current);
            record_timer_counter(self.metrics.as_ref(), MetricId::TimerCancelledTotal, reason);
        }
    }

    /// 当前挂起的 timer 数量(测试用)
    #[cfg(test)]
    pub fn pending_count(&mut self) -> usize {
        // 先回收已自然到期的 one-shot 条目再报数,否则已 fire 的条目仍滞留计数。
        self.reap();
        self.handles.len()
    }
}

/// Drop 时 abort 所有挂起 timer 任务:把「`run_engine_loop` 返回 = 完全 drain」的
/// 保证扩展到 timer。`JoinHandle::drop` 仅 detach 不 abort——否则长 timer 任务会泄漏到
/// runtime 关停,且持 tick_tx clone 延迟 channel 关闭。与 `cancel` 的 abort 语义一致。
impl Drop for TimerRegistry {
    fn drop(&mut self) {
        for (_, h) in self.handles.drain() {
            h.abort();
        }
    }
}

/// 记录 timer 低基数 Counter,不携带 TimerId。
fn record_timer_counter(metrics: &dyn AsyncMetricSink, id: MetricId, status: &'static str) {
    if metrics.is_enabled() {
        let _ = metrics.try_record(MetricEvent::counter(
            id,
            1.0,
            MetricLabels::one(LabelKey::Stage, "timer").with(LabelKey::Status, status),
        ));
    }
}

/// 记录 timer 耗时直方图,禁止把 timer payload 带入标签。
fn record_timer_histogram(metrics: &dyn AsyncMetricSink, id: MetricId, value: std::time::Duration) {
    if metrics.is_enabled() {
        let _ = metrics.try_record(MetricEvent::histogram(
            id,
            value.as_secs_f64(),
            MetricLabels::one(LabelKey::Stage, "timer"),
        ));
    }
}

/// 发布当前活跃 timer 数量。
fn record_timer_gauge(metrics: &dyn AsyncMetricSink, active: usize) {
    if metrics.is_enabled() {
        let _ = metrics.try_record(MetricEvent::gauge(
            MetricId::TimerActive,
            active as f64,
            MetricLabels::one(LabelKey::Stage, "timer"),
        ));
    }
}

/// 原子饱和递减,处理 cancel 与自然 fire 的竞争边界。
fn decrement_saturating(value: &AtomicUsize) -> usize {
    value
        .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
            Some(current.saturating_sub(1))
        })
        .unwrap_or_default()
        .saturating_sub(1)
}

// ─── 单元测试 ────────────────────────────────────────────────────────────────

#[cfg(test)]
mod tests {
    use super::*;
    use crate::metrics::RecordOutcome;
    use helix_core::effect::TimerId;
    use std::sync::Mutex;
    use tokio::sync::mpsc;

    #[derive(Default)]
    struct RecordingMetricSink(Mutex<Vec<MetricEvent>>);

    impl AsyncMetricSink for RecordingMetricSink {
        /// 保存 timer 外部指标事件,不启动真实 exporter。
        fn try_record(&self, event: MetricEvent) -> RecordOutcome {
            self.0.lock().unwrap().push(event);
            RecordOutcome::Accepted
        }
    }

    #[tokio::test]
    async fn test_timer_fires_after_delay() {
        let (tx, mut rx) = mpsc::channel::<Tick>(8);
        let mut registry = TimerRegistry::new();

        let id = TimerId::from_raw(42);
        registry.schedule(id, 10, EngineTickSender::Raw(tx)); // 10ms 后触发

        let tick = tokio::time::timeout(tokio::time::Duration::from_millis(200), rx.recv())
            .await
            .expect("timer should fire within 200ms")
            .expect("channel should not be closed");

        assert!(matches!(tick, Tick::Timer(tid) if tid == id));
    }

    #[tokio::test]
    async fn test_cancel_prevents_fire() {
        let (tx, mut rx) = mpsc::channel::<Tick>(8);
        let mut registry = TimerRegistry::new();

        let id = TimerId::from_raw(99);
        // 10s timer — 立即 cancel 后,在 50ms 窗口内不应收到 Tick::Timer
        registry.schedule(id, 10_000, EngineTickSender::Raw(tx.clone())); // clone tx 保持通道开放
        registry.cancel(id);
        // tx 仍然 alive(clone 持有),rx.recv() 会阻塞(而非立即返回 None)
        // 因此 timeout 50ms 后应超时(Err),而非收到消息(Ok)

        let result = tokio::time::timeout(tokio::time::Duration::from_millis(50), rx.recv()).await;

        // timeout elapsed = Err(Elapsed),表示 timer 未触发(符合预期)
        assert!(
            result.is_err(),
            "cancelled timer should not fire within 50ms, result was ready"
        );

        // 显式 drop tx(避免 rx 挂起)
        drop(tx);
    }

    #[tokio::test]
    async fn test_reschedule_overwrites() {
        let (tx, mut rx) = mpsc::channel::<Tick>(8);
        let mut registry = TimerRegistry::new();

        let id = TimerId::from_raw(7);
        registry.schedule(id, 10_000, EngineTickSender::Raw(tx.clone())); // 10s timer
        registry.schedule(id, 10, EngineTickSender::Raw(tx)); // 10ms 覆盖

        let tick = tokio::time::timeout(tokio::time::Duration::from_millis(200), rx.recv())
            .await
            .expect("reschedule should fire fast")
            .expect("channel open");

        assert!(matches!(tick, Tick::Timer(tid) if tid == id));
    }

    #[tokio::test]
    async fn test_cancel_idempotent() {
        let mut registry = TimerRegistry::new();
        let id = TimerId::from_raw(1);
        // 取消不存在的 id 不 panic
        registry.cancel(id);
        registry.cancel(id);
    }

    /// one-shot timer 自然到期(收到 Tick::Timer)后,handles 不残留该 id(防缓慢内存泄漏)。
    #[tokio::test]
    async fn test_oneshot_timer_reaped_after_fire() {
        let (tx, mut rx) = mpsc::channel::<Tick>(8);
        let mut registry = TimerRegistry::new();

        let id = TimerId::from_raw(123);
        registry.schedule(id, 10, EngineTickSender::Raw(tx)); // 10ms 后触发的一次性 timer

        // 等待 timer 自然触发
        let tick = tokio::time::timeout(tokio::time::Duration::from_millis(200), rx.recv())
            .await
            .expect("timer should fire within 200ms")
            .expect("channel should not be closed");
        assert!(matches!(tick, Tick::Timer(tid) if tid == id));

        // 闭包在 send 之后才回灌 done_tx;让出执行权确保回灌已入通道,
        // 随后 reap(pending_count 内部触发)把已完成 id 从 handles 移除。
        tokio::task::yield_now().await;

        // pending_count() 内部调 reap(),应观察到该 one-shot 已被回收 → 归零。
        assert_eq!(
            registry.pending_count(),
            0,
            "one-shot timer should be reaped from handles after natural fire"
        );
    }

    /// timer 自然触发必须闭合 scheduled、fired、lateness、delivery 与 active Gauge。
    #[tokio::test]
    async fn observed_timer_records_lifecycle_without_timer_id_label() {
        let metrics = Arc::new(RecordingMetricSink::default());
        let (tx, mut rx) = mpsc::channel::<Tick>(8);
        let mut registry = TimerRegistry::with_metrics(metrics.clone());

        registry.schedule(TimerId::from_raw(321), 1, EngineTickSender::Raw(tx));
        tokio::time::timeout(tokio::time::Duration::from_millis(200), rx.recv())
            .await
            .expect("timer 应在窗口内触发")
            .expect("timer channel 应保持打开");
        tokio::task::yield_now().await;

        let events = metrics.0.lock().unwrap();
        for expected in [
            MetricId::TimerScheduledTotal,
            MetricId::TimerFiredTotal,
            MetricId::TimerLatenessSeconds,
            MetricId::TimerDeliveryWaitSeconds,
            MetricId::TimerActive,
        ] {
            assert!(
                events.iter().any(|event| event.id == expected),
                "timer 缺少 {:?}",
                expected
            );
        }
        assert!(events
            .iter()
            .filter(|event| event.id == MetricId::TimerActive)
            .any(|event| event.value == 0.0));
    }
}