Skip to main content

helix_driver_native/
engine_loop.rs

1//! engine_loop — native 薄壳(pump 核已下沉 helix-driver-host::engine)。
2//!
3//! ## 下沉后的分工
4//!
5//! pump 核(`BoundedSpawner` / `OwnedEffect`+`consume_effects` / `dispatch_effects` /
6//! `TimerRegistry` / `run_http_envelope` / 双通道 biased select + graceful drain)统一在
7//! `helix-driver-host::engine`(泛型 over Storage/HttpRequester/EventSink/FrameSender/Clock)。
8//! native 这层只保留:
9//! - 具体 ports 装配(`NativeStorage` / `NativeHttp` / `NativeEventSink` / `NativeClock`)
10//! - 与现网兼容的公开 API 签名([`EngineConfig`] / [`TransportTable`] / [`run_engine_loop`])
11//! - tick 源接线(`SharedWsClient::with_inbound_tick`,在 `transport.rs`)
12//!
13//! NativeEventSink(broadcast egress)特化留 native;(下一轮)FFI 传批处理 vtable sink,
14//! 共用同一 host 泛型壳,只差 egress + tick 源。
15//!
16//! ## 五不变量真源
17//!
18//! 五不变量(回灌必达 / 主循环不等 I/O / 同 key 写序 / 回灌优先 / graceful drain)
19//! 的实现与文档真源现在在 `helix-driver-host::engine` 头注 + 本 crate AGENTS.md。
20
21use std::sync::Arc;
22
23use helix_driver_host::engine::{self, EngineDeps, TransportLifecycleEvent, TransportTraceSink};
24use helix_driver_host::{AsyncMetricSink, NoopMetricSink, SharedFileUploader, TraceHooks};
25use helix_driver_host::{TickIngressReceiver, TickIngressSender};
26// rows↔bytes 编解码:host 单一真源,native re-export(helix-im 反序列化端对齐 host codec)。
27pub use helix_driver_host::{rows_from_reply_bytes, rows_to_reply_bytes};
28
29use helix_core::effect::TransportId;
30use helix_core::{ExecutionShell, Tick};
31use std::collections::HashMap;
32use tokio::sync::{mpsc, oneshot};
33
34use crate::clock::NativeClock;
35use crate::event_sink::NativeEventSink;
36use crate::http::NativeHttp;
37use crate::storage::NativeStorage;
38use crate::transport::NativeTransport;
39
40mod reconnect;
41use reconnect::{start_native_reconnect, NativeReconnectRuntime};
42
43/// TransportId → 实际 WS 连接 的路由表(native 特化:值为 `NativeTransport`)。
44///
45/// host 装配阶段填入(`host.add_transport(...)` → `TransportId`);host 泵的 `Effect::Send`
46/// 分支按 id 查表调 `NativeTransport::send`(`&self`)。多连接形态从一开始做成 `HashMap`。
47/// `Arc<_>`(无外层 `Mutex`):FrameSender::send 只需 `&self`(host `SharedWsClient` 内部 state
48/// 锁串行化),表里只需共享句柄;connect 在装配端独占 `&mut` 持有时完成(见 host engine 头注)。
49pub type TransportTable = HashMap<TransportId, Arc<NativeTransport>>;
50
51/// 事件泵配置(native 具体 ports 聚合)。
52///
53/// 装配后由 [`run_engine_loop`] / [`run_engine_loop_with_transports`] 转成 host 的
54/// [`EngineDeps`](注入 `NativeClock`),调泛型壳 `helix_driver_host::engine::run_engine_loop`。
55pub struct EngineConfig {
56    pub storage: NativeStorage,
57    pub http: NativeHttp,
58    pub uploader: SharedFileUploader,
59    pub event_sink: NativeEventSink,
60    /// Engine 热路径指标 sink;composition root 应与 HTTP/WS/Storage/Event 共用同一 Arc。
61    pub metrics: Arc<dyn AsyncMetricSink>,
62    /// Http effect 并发上限(Http BoundedSpawner 的 worker 数 N)。
63    ///
64    /// 冷启动时 core 一步可吐 N 条 `Effect::Http`;若全部齐发会瞬间打满底层连接
65    /// (reqwest pool **不**限制新连接并发)。此字段 = Http/HttpFire BoundedSpawner 的常驻
66    /// worker 数:第 N+1 条请求在有界队列里等 worker,不并发齐发。host(ambient API 合法落点)
67    /// 读 env 注入;建议默认 8。`0` 会被 `.max(1)` 兜底为 1。
68    pub max_http_inflight: usize,
69}
70
71impl EngineConfig {
72    pub fn new(
73        storage: NativeStorage,
74        http: NativeHttp,
75        uploader: SharedFileUploader,
76        event_sink: NativeEventSink,
77        max_http_inflight: usize,
78    ) -> Self {
79        Self {
80            storage,
81            http,
82            uploader,
83            event_sink,
84            metrics: Arc::new(NoopMetricSink),
85            max_http_inflight,
86        }
87    }
88
89    pub fn with_metric_sink(mut self, metrics: Arc<dyn AsyncMetricSink>) -> Self {
90        self.metrics = metrics;
91        self
92    }
93
94    /// 转成 host 泛型壳所需的 [`EngineDeps`](注入 `NativeClock`)。
95    fn into_deps(
96        self,
97        trace: TraceHooks,
98        transport_lifecycle_tx: Option<mpsc::UnboundedSender<TransportLifecycleEvent>>,
99        transport_trace_tx: Option<TransportTraceSink>,
100    ) -> EngineDeps<NativeStorage, NativeHttp, SharedFileUploader, NativeEventSink, NativeClock>
101    {
102        EngineDeps {
103            storage: Arc::new(self.storage),
104            http: Arc::new(self.http),
105            uploader: Arc::new(self.uploader),
106            event_sink: Arc::new(self.event_sink),
107            clock: NativeClock,
108            trace,
109            metrics: self.metrics,
110            max_http_inflight: self.max_http_inflight,
111            transport_lifecycle_tx,
112            transport_trace_tx,
113        }
114    }
115}
116
117/// 运行 ExecutionShell 的 tokio 事件泵(native 入口,无 transport 路由)。
118///
119/// 无 transport 的兼容入口:Send 表为空 → `Effect::Send` 落 warn(回退既有行为)。
120///
121/// ## 使用方式
122///
123/// ```rust,ignore
124/// let (tick_tx, tick_rx) = mpsc::channel::<Tick>(256);
125/// let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
126/// let config = EngineConfig {
127///     storage,
128///     http,
129///     uploader: helix_driver_host::SharedFileUploader::default(),
130///     event_sink,
131///     metrics: std::sync::Arc::new(helix_driver_host::NoopMetricSink),
132///     max_http_inflight: 8,
133/// };
134/// tokio::spawn(run_engine_loop(shell, tick_rx, tick_tx.clone(), config, shutdown_rx));
135/// ```
136///
137/// ## Shutdown 语义
138///
139/// `shutdown_tx.send(())` 后进入 graceful drain,返回 = 全部已接收 Persist + 必达 Http 已落地。
140pub async fn run_engine_loop(
141    shell: ExecutionShell,
142    tick_rx: mpsc::Receiver<Tick>,
143    tick_tx: mpsc::Sender<Tick>,
144    config: EngineConfig,
145    shutdown_rx: oneshot::Receiver<()>,
146) {
147    run_engine_loop_with_transports(
148        shell,
149        tick_rx,
150        tick_tx,
151        config,
152        shutdown_rx,
153        TransportTable::new(),
154    )
155    .await
156}
157
158/// 生产指标入口:保持单一有界队列,在元素内携带真实入队时间戳。
159pub async fn run_engine_loop_stamped(
160    shell: ExecutionShell,
161    tick_rx: TickIngressReceiver,
162    tick_tx: TickIngressSender,
163    config: EngineConfig,
164    shutdown_rx: oneshot::Receiver<()>,
165) {
166    run_engine_loop_with_transports_and_trace_stamped(
167        shell,
168        tick_rx,
169        tick_tx,
170        config,
171        TraceHooks::noop(),
172        shutdown_rx,
173        TransportTable::new(),
174    )
175    .await;
176}
177
178/// 带 transport 路由表的事件泵入口(A1)。
179///
180/// 与 [`run_engine_loop`] 唯一区别:`transports` 把 `TransportId → NativeTransport`
181/// 接进泵,使 `Effect::Send` 真正发到 WS。host 装配时 `connect()` 各 transport 后填表。
182///
183/// 装配 native ports → host 泛型壳 `engine::run_engine_loop`(单一 pump 核)。
184pub async fn run_engine_loop_with_transports(
185    shell: ExecutionShell,
186    tick_rx: mpsc::Receiver<Tick>,
187    tick_tx: mpsc::Sender<Tick>,
188    config: EngineConfig,
189    shutdown_rx: oneshot::Receiver<()>,
190    transports: TransportTable,
191) {
192    run_engine_loop_with_transports_and_trace(
193        shell,
194        tick_rx,
195        tick_tx,
196        config,
197        TraceHooks::noop(),
198        shutdown_rx,
199        transports,
200    )
201    .await;
202}
203
204/// 显式注入 trace hooks 的 native composition-root 入口;旧入口保持真 no-op。
205pub async fn run_engine_loop_with_transports_and_trace(
206    shell: ExecutionShell,
207    tick_rx: mpsc::Receiver<Tick>,
208    tick_tx: mpsc::Sender<Tick>,
209    config: EngineConfig,
210    trace: TraceHooks,
211    shutdown_rx: oneshot::Receiver<()>,
212    transports: TransportTable,
213) {
214    let NativeReconnectRuntime {
215        transport_rx,
216        lifecycle_tx,
217        trace_sink,
218        trace_stats,
219        tasks,
220        shutdown_tx: reconnect_shutdown_tx,
221    } = start_native_reconnect(&transports);
222    engine::run_engine_loop(
223        shell,
224        tick_rx,
225        tick_tx,
226        config.into_deps(trace, lifecycle_tx, trace_sink),
227        shutdown_rx,
228        transports,
229        transport_rx,
230    )
231    .await;
232    if let Some(stats) = trace_stats {
233        let dropped = stats.dropped_count();
234        if dropped > 0 {
235            tracing::warn!(dropped, "native transport trace 有界队列发生丢弃");
236        }
237    }
238    reconnect_shutdown_tx.send_replace(true);
239    for task in tasks {
240        task.abort();
241    }
242}
243
244pub async fn run_engine_loop_with_transports_and_trace_stamped(
245    shell: ExecutionShell,
246    tick_rx: TickIngressReceiver,
247    tick_tx: TickIngressSender,
248    config: EngineConfig,
249    trace: TraceHooks,
250    shutdown_rx: oneshot::Receiver<()>,
251    transports: TransportTable,
252) {
253    let NativeReconnectRuntime {
254        transport_rx,
255        lifecycle_tx,
256        trace_sink,
257        trace_stats,
258        tasks,
259        shutdown_tx: reconnect_shutdown_tx,
260    } = start_native_reconnect(&transports);
261    engine::run_engine_loop_stamped(
262        shell,
263        tick_rx,
264        tick_tx,
265        config.into_deps(trace, lifecycle_tx, trace_sink),
266        shutdown_rx,
267        transports,
268        transport_rx,
269    )
270    .await;
271    if let Some(stats) = trace_stats {
272        let dropped = stats.dropped_count();
273        if dropped > 0 {
274            tracing::warn!(dropped, "native transport trace 有界队列发生丢弃");
275        }
276    }
277    reconnect_shutdown_tx.send_replace(true);
278    for task in tasks {
279        task.abort();
280    }
281}
282
283#[cfg(test)]
284#[path = "engine_loop_tests.rs"]
285mod tests;