小类随手记

Rust 双模式实时交互:一套事件总线同时服务 Web 与 Tauri

记录一种 Rust 应用的实时数据交互设计:核心层提供传输无关的事件总线,Web 模式封装为 SSE,Tauri 模式封装为原生事件,前端用统一抽象消费。

本文记录一种 Rust 应用的实时数据交互设计:一个 Rust 后端同时服务两种前端形态(浏览器 Web 客户端与 Tauri 桌面客户端 WebView)时,如何用"传输无关的事件总线 + 两端薄适配 + 前端统一订阅抽象"实现实时数据推送,避免为两种传输各写一套逻辑。

背景与挑战

一个 Rust 应用常常有两种形态:Web 版(Rust 起 HTTP 服务,浏览器访问)和桌面版(Tauri 壳 + 内置 WebView)。业务逻辑在核心层共享,但前后端通信通道完全不同:

  • Web 模式:浏览器只能走 HTTP,实时推送用 SSE 或 WebSocket
  • Tauri 模式:前端运行在 Rust 进程内的 WebView 里,可以直接走 Tauri 的原生事件通道(emit/listen),绕开 HTTP

如果为两种模式分别实现实时交互,会产生两个问题:

  • 业务代码里混入传输细节,核心层被 axum::Ssetauri::Emitter 污染
  • 前端要维护两套订阅逻辑,行为不一致(比如鉴权、重连策略不同)

目标设计:核心层只发布领域事件,不感知传输;服务壳与 Tauri 壳各做一个薄适配层;前端提供与 API 调用对称的统一订阅抽象。

总体架构

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
写操作(核心层业务逻辑)
        │ 发布
┌──────────────────────┐
│  EventBus(事件总线) │  传输无关,基于 tokio::sync::broadcast
└───────┬──────────────┘
        │ 订阅
   ┌────┴─────┐
   ▼          ▼
服务壳(SSE)   Tauri壳(原生Event)
GET /api/events  app.emit("app://event")
   │              │
   ▼              ▼
浏览器 fetch 流   WebView listen
   └──────┬───────┘
前端统一 event-client(按运行环境分流)
   TanStack Query invalidate / 状态 store

核心思路:核心层发布事件时不知道事件会以什么方式到达前端;传输适配只在两个壳里完成;前端用一个客户端抽象屏蔽差异。

事件总线设计

事件总线放在核心层,与业务逻辑同进程,基于 tokio::sync::broadcast

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
pub struct Event {
    pub topic: &'static str,   // 事件主题,如 "task.progress"
    pub seq: u64,              // 单调递增序号,用于断线补漏
    pub payload: serde_json::Value,
}

pub struct EventBus {
    tx: broadcast::Sender<Event>,
}

impl EventBus {
    pub fn publish(&self, topic: &'static str, payload: serde_json::Value) {
        let seq = self.next_seq();
        let _ = self.tx.send(Event { topic, seq, payload });
    }

    pub fn subscribe(&self) -> broadcast::Receiver<Event> {
        self.tx.subscribe()
    }
}

选择 broadcast 而不是 mpsc 的原因:

  • 天然支持多订阅者:多个浏览器 Tab、多个设备、Tauri 转发器可以同时订阅同一份事件流
  • 慢订阅者自动落后,配合 lagged 计数可以检测断线期间丢了多少事件
  • 容量按最坏情况配置(如 1024),事件发布是内存操作,成本极低

事件主题统一为 &'static str 常量,payload 一律 serde_json::Value。这样无论走 SSE 还是 Tauri 原生通道,payload 都是同一个结构,序列化零转换。

发布侧:谁写谁发布

实时交互的关键约定:所有写路径必须收敛到核心层的方法上,在成功完成后发布对应主题的事件。如果写路径分散在壳层,就会漏发事件导致各端状态不一致。

两个典型场景:

  • 长任务进度:worker 循环里拿到进度后发布 task.progress,但要节流(如每 500ms 一条),避免高频事件压垮订阅者和前端;终态(完成/失败/取消)立即发布不节流
  • 配置变更:增删改配置成功后发布 config.changed,前端收到后使相关查询失效,重新拉取

节流放哪一层?放在发布侧(核心层)而不是消费侧,因为一个发布者可能对应多个订阅者,源头节流才能保证总流量可控。payload 保持小体积,只带最小必要字段。

服务壳:SSE 适配

Web 模式用 SSE(Server-Sent Events)而不是 WebSocket:

  • SSE 基于普通 HTTP,天然支持鉴权头、自动重连、与现有 REST 体系共用中间件
  • 单向推送足够(前端不需要通过事件通道回传数据,回传走普通 API)
  • axum 有现成的 Sse + KeepAlive 支持
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
pub async fn events_subscribe(
    State(app): State<AppContext>,
    Query(params): Query<AuthParams>,   // ?token= 鉴权
) -> Sse<impl Stream<Item = Result<Event, Infallible>>> {
    // 1. 手动校验 JWT(EventSource / fetch 无法自定义 Authorization 头时的备选)
    // 2. 从 EventBus 订阅一个 receiver
    let rx = app.events.subscribe();
    let stream = BroadcastStream::new(rx).map(|item| {
        let event = item.unwrap_or_else(|lagged| /* 断线补偿事件 */);
        Ok(Event::json(event))
    });
    Sse::new(stream).keep_alive(KeepAlive::new().interval(Duration::from_secs(15)))
}

设计要点:

  • 每个连接一个 broadcast 订阅者,连接断开时 receiver 自动丢弃,不泄漏
  • 15s KeepAlive ping 保持代理/负载均衡器不断连,也用于检测死连接
  • SSE 事件名直接用 topic,前端 addEventListener(topic) 按主题分发
  • 鉴权不走全局中间件,路由挂在受保护路由组之外,handler 内自校验,支持 ?token= 查询参数——因为浏览器 EventSource 无法自定义请求头

Tauri 壳:原生事件适配

Tauri 模式没有 HTTP,前端 WebView 与 Rust 进程通过原生事件通道通信。适配层只是一个转发器:订阅 EventBus,把每个事件 app.emit 出去。

1
2
3
4
5
6
7
8
// setup 中 spawn 一个转发任务
let bus = app.state::<AppContext>().events.clone();
tauri::async_runtime::spawn(async move {
    let mut rx = bus.subscribe();
    while let Ok(event) = rx.recv().await {
        let _ = app.emit("app://event", &event);
    }
});

要点:

  • payload 是 serde_json::ValueSerialize + Clone 天然满足,直接传引用即可
  • 桌面模式是单进程内通信,无需鉴权、无需重连、无需 KeepAlive,比 SSE 简单一个数量级
  • 事件通道名固定一个(如 app://event),内部再按 topic 分发,与 SSE 的"事件名 = topic"在语义上对齐

前端统一订阅抽象

前端与 api-client 对称,提供一个 event-client 单例:启动时检测运行环境(浏览器 fetch 还是 Tauri invoke),选择对应实现,向上暴露同一套 API:

1
2
3
4
5
type EventHandler = (payload: unknown) => void;

// 统一接口
subscribe(topic: string, handler: EventHandler): Unsubscribe;
onStatus(cb: (connected: boolean) => void): void;

Web 模式:fetch + ReadableStream 解析 SSE

浏览器原生 EventSource 有两个硬伤:无法自定义 Authorization 头(只能 ?token=,token 会进日志),且无法精细控制重连。所以用 fetch + ReadableStream 手动解析 SSE 帧:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
async function connectSSE(token: string) {
  const res = await fetch('/api/events', {
    headers: { Authorization: `Bearer ${token}` },
  });
  const reader = res.body!.getReader();
  const decoder = new TextDecoder();
  let buffer = '';
  while (true) {
    const { done, value } = await reader.read();
    if (done) break;
    buffer += decoder.decode(value, { stream: true });
    // 按空行切分 SSE 帧,解析 event: / data: 字段
    // 按 topic 分发到 handler
  }
}

Tauri 模式:原生 listen

动态 import @tauri-apps/api/event,避免 Web 包打进多余的 Tauri 依赖:

1
2
3
4
5
6
7
async function connectTauri() {
  const { listen } = await import('@tauri-apps/api/event');
  await listen('app://event', (e) => {
    const event = e.payload as { topic: string; payload: unknown };
    dispatch(event.topic, event.payload);
  });
}

统一的重连与状态管理

  • 连接状态(connected/disconnected)通过 onStatus 通知 hooks
  • Web 模式断线后指数退避重连(1s → 2s → 4s → 上限),401 时先刷新 token 再重连
  • Tauri 模式进程内通道几乎不会断,但代码路径保持同一套接口

消费侧:事件驱动的查询失效与状态合并

前端拿到事件后怎么用,分两种场景:

查询失效(配置类事件)

配置变更事件不携带完整数据,只带变更动作。前端收到后让相关 TanStack Query 失效(invalidate),下次读取自动重新拉取。这样事件通道只负责"通知变了",数据一致性仍由查询层保证,不用在事件里塞大 payload。

1
2
3
4
// 全局订阅,收到配置变更就失效对应查询
subscribe('config.changed', () => {
  queryClient.invalidateQueries({ queryKey: ['config'] });
});

状态合并(进度类事件)

高频进度事件不适合进查询层(会触发整页 re-render),用 Zustand store 单独管理:

1
2
3
4
// 进度 store:task_id -> 最新进度 payload
subscribe('task.progress', (p) => {
  useProgressStore.setState((s) => ({ [p.taskId]: p.payload }));
});

UI 渲染时把查询数据(任务的静态信息)与 store 数据(实时进度)合并。终态事件到达时清掉 store 对应条目并失效查询,恢复为查询层的数据源。

轮询降级

订阅可用时停掉轮询,断开时自动恢复轮询兜底:

1
refetchInterval: connected ? false : 3000,

这条规则保证:实时通道正常时零轮询、零延迟;断线期间退化到轮询,数据不中断;重连成功后自动切回事件模式。

可靠性设计

  • 断线补漏:事件带 seq 单调递增,订阅者可以检测丢失。落地时先做"连接建立即快照补齐"——重连成功后把相关查询全部 invalidate 一遍,用一次全量拉取弥补断线窗口丢失的事件;seq 留作后续精确补漏
  • 事件丢失兜底:broadcast 满容量时慢订阅者会 lagged,转发器要处理 RecvError::Lagged,把 lagged 计数当作一次补偿信号(触发快照补齐)
  • 高频事件节流在发布侧,保证多订阅者场景下总量可控
  • 长连接与普通 API 的超时策略分开:SSE 连接不能套用 API 的 30s 超时,前端要为事件通道单独实现

总结

这套设计最终形成三个对称:

  • 核心层不感知传输,发布侧与传输侧解耦
  • 服务壳与 Tauri 壳各只做 50 行左右的适配代码,没有业务逻辑
  • 前端 api-client(请求)与 event-client(订阅)对称,按运行环境分流

对于"同一个 Rust 核心要同时喂浏览器和桌面端"的应用,这是成本最低、行为最一致的实时交互方案。事件总线、SSE、Tauri 原生事件三者各司其职,前端永远只面对一个订阅接口。

相关文章

comments powered by Disqus
Theme Stack