diff --git a/crates/df-tunnel/src/tunnel.rs b/crates/df-tunnel/src/tunnel.rs index 52fd30b..2111adf 100644 --- a/crates/df-tunnel/src/tunnel.rs +++ b/crates/df-tunnel/src/tunnel.rs @@ -81,6 +81,13 @@ pub trait TunnelClient: Send + Sync { /// /// 非阻塞入队,实际 socket 写入异步完成。连接断开时返回 NotConnected。 async fn send_event(&self, event: TunnelEvent) -> Result<()>; + + /// 发送原始事件(桌面端 → 云后端 → 小程序),payload 为 AiChatEvent 序列化 JSON Value + /// + /// Phase3 阶段2(D2 全 19 变体透传):EventBus subscriber 把 AiChatEvent Value 经此方法 + /// 透传(不经 TunnelEvent 强类型子集)。与 [`send_event`](TunnelClient::send_event) 伴行, + /// 后者保留作高频子集快捷方式(D5 强类型保留)。非阻塞入队,连接断开返 NotConnected。 + async fn send_raw_event(&self, payload: serde_json::Value) -> Result<()>; } /// 握手探测窗口(relay 校验通过后静默进入收发循环,无 ack 帧; @@ -105,8 +112,11 @@ struct ConnState { /// 收发循环出站帧(SendLoop 内部消费) enum OutFrame { - /// 业务事件(序列化后写 socket) - Event(TunnelEvent), + /// 业务事件(已序列化 JSON 字符串,run_loop 直接写 socket) + /// + /// send_event(TunnelEvent) 与 send_raw_event(Value) 各自预序列化为 String 入队, + /// 统一经此变体写 socket(DRY:序列化在调用侧,run_loop 只负责写)。 + Event(String), /// 优雅关闭 Close, } @@ -260,13 +270,29 @@ impl TunnelClient for WsTunnelClient { } async fn send_event(&self, event: TunnelEvent) -> Result<()> { + // 序列化提前(OutFrame::Event 持 String),失败返 Serde 错给调用方 + let json = serde_json::to_string(&event)?; let conn = self.conn.lock().await; let Some(state) = conn.as_ref() else { return Err(TunnelError::NotConnected); }; state .tx - .send(OutFrame::Event(event)) + .send(OutFrame::Event(json)) + .map_err(|_| TunnelError::NotConnected)?; + Ok(()) + } + + async fn send_raw_event(&self, payload: serde_json::Value) -> Result<()> { + // AiChatEvent Value 序列化为 JSON 字符串入队(与 send_event 共用 OutFrame::Event(String)) + let json = serde_json::to_string(&payload)?; + let conn = self.conn.lock().await; + let Some(state) = conn.as_ref() else { + return Err(TunnelError::NotConnected); + }; + state + .tx + .send(OutFrame::Event(json)) .map_err(|_| TunnelError::NotConnected)?; Ok(()) } @@ -358,14 +384,8 @@ async fn run_loop( // 出站队列 maybe_out = rx.recv() => { match maybe_out { - Some(OutFrame::Event(ev)) => { - let json = match serde_json::to_string(&ev) { - Ok(j) => j, - Err(e) => { - tracing::warn!(error = %e, "事件序列化失败,跳过"); - continue; - } - }; + // OutFrame::Event 持已序列化 JSON 字符串(send_event/send_raw_event 预序列化) + Some(OutFrame::Event(json)) => { if let Err(e) = sink.send(WsMessage::Text(json)).await { tracing::warn!(error = %e, "socket 写入失败,收发循环退出"); break;