首页
/ airi 中的 @moeru/eventa:跨进程事件通信与类型安全 RPC 实战指南

airi 中的 @moeru/eventa:跨进程事件通信与类型安全 RPC 实战指南

2026-09-05 13:58:33作者:乔或婵

本篇技术指南围绕 airi 仓库中的 @moeru/eventa 技能文档(.agents/skills/eventa/SKILL.md)展开,系统讲解这个“传输感知的事件库”如何支撑 Electron IPC、WebSocket、Web Worker、BroadcastChannel 等多种进程间通信,并演示一元 RPC、三种流式调用模式、取消机制与多跳频道路由。读完后,你将能够复用 airi 中的真实代码模式,为任意 Electron/Node/浏览器应用定义类型安全的事件与 RPC 接口。

核心设计思想

Eventa 的整个 API 围绕三个基本理念构建,理解它们是掌握这个库的前提:

  1. 事件是一等公民 —— 定义一次带类型的事件,到处复用;
  2. 传输层可替换 —— 同一份事件定义可以在 Electron IPC、WebSocket、Web Workers、BroadcastChannel、EventEmitter、EventTarget 和 Worker Threads 之间无缝切换;
  3. RPC 只是事件的组合 —— invoke(一元调用)和 stream(流式调用)模式都从同一套事件原语组合而来,而不是独立的一套机制。

这种设计在 airi 中得到了大量验证:当前仓库通过 pnpm-workspace.yaml@moeru/eventa 锁定在 catalog 版本 1.0.0,并被 packages/electron-eventapackages/better-wsserver/packages/server-sdk-sharedpackages/plugin-sdkpackages/pipelines-audio 等众多包依赖,是贯穿桌面端、Web 端与后端的统一事件基础设施。

API 速览

事件定义与上下文(Event Definition & Context)

defineEventa 的泛型参数即事件负载类型。createContext() 创建一个内存基础上下文,适用于同进程通信;ctx.emit 发送事件,ctx.on 注册监听:

import { createContext, defineEventa } from '@moeru/eventa'

// Define a typed event (the generic is the payload type)
const move = defineEventa<{ x: number, y: number }>()

// Create a base context (in-memory, useful for same-process communication)
const ctx = createContext()

// Emit and listen
ctx.emit(move, { x: 100, y: 200 })
ctx.on(move, ({ body }) => console.log(body.x, body.y))

airi 中的真实用例印证了这一模式。packages/pipelines-audio/src/eventa.ts 为音频流水线集中定义了一组事件,如 speechSegmentEventspeechTtsRequestEventspeechPlaybackStartEvent,并将它们聚合成 speechPipelineEventMap 供流水线各环节按名订阅;packages/better-ws/src/client/index.ts 也用 defineEventa<ClientStateChange>('better-ws:client:state-change') 定义了 WebSocket 客户端的状态变更事件。可以看到项目遵循的命名约定是 领域:路径:事件名 的点分风格。

一元 RPC(Invoke)

defineInvokeEventa<ResponseType, RequestType>(optionalName) 的关键点是 Response 类型在前,Request 类型在后

import { createContext, defineInvokeEventa, defineInvokeHandler, defineInvoke } from '@moeru/eventa'

const ctx = createContext()

// defineInvokeEventa<ResponseType, RequestType>(optionalName)
const echo = defineInvokeEventa<{ output: string }, { input: string }>('rpc:echo')

// Register handler (server side)
defineInvokeHandler(ctx, echo, ({ input }) => ({ output: input.toUpperCase() }))

// Create invoke function (client side)
const invokeEcho = defineInvoke(ctx, echo)
const result = await invokeEcho({ input: 'hello' }) // { output: 'HELLO' }

airi 的桌面端就是这样组织的:packages/electron-eventa 包(README)被定位为“跨 AIRI 各应用共享的 Electron IPC Eventa 契约定义层”。例如 packages/electron-eventa/src/electron/window.ts 中定义了窗口控制契约:

import { defineEventa, defineInvokeEventa } from '@moeru/eventa'

export const bounds = defineEventa<Rectangle>('eventa:event:electron:window:bounds')
const getBounds = defineInvokeEventa<ReturnType<BrowserWindow['getBounds']>>('eventa:invoke:electron:window:get-bounds')
const setBounds = defineInvokeEventa<void, Parameters<BrowserWindow['setBounds']>>('eventa:invoke:electron:window:set-bounds')

这里通过 ReturnType / Parameters 直接从 Electron 类型推导请求与响应类型,保证 RPC 契约与 Electron API 签名严格同步——这正是“事件定义放在共享模块”规则的实际收益。

服务端流式 RPC(Server-Streaming)

服务端流用于“一次请求、多次产出”的场景,例如进度上报。事件定义中 Response 类型是一个可区分的更新联合体:

import { createContext, defineInvokeEventa, defineStreamInvoke, defineStreamInvokeHandler, toStreamHandler } from '@moeru/eventa'

const ctx = createContext()
const sync = defineInvokeEventa<
  { type: 'progress' | 'result', value: number },
  { jobId: string }
>('rpc:sync')

// Generator-style handler
defineStreamInvokeHandler(ctx, sync, async function* ({ jobId }) {
  for (let i = 1; i <= 5; i++) {
    yield { type: 'progress' as const, value: i * 20 }
  }
  yield { type: 'result' as const, value: 100 }
})

// Or imperative style with toStreamHandler
defineStreamInvokeHandler(ctx, sync, toStreamHandler(async ({ payload, emit }) => {
  emit({ type: 'progress', value: 0 })
  emit({ type: 'result', value: 100 })
}))

// Consume as async iterator
const stream = defineStreamInvoke(ctx, sync)
for await (const update of stream({ jobId: 'import' })) {
  console.log(update.type, update.value)
}

服务端处理器支持两种写法:生成器风格(async function* 逐帧 yield)和命令式风格(toStreamHandler 包装的 emit 回调)。客户端统一通过 defineStreamInvoke 得到函数,返回值可作 async iterable 消费。

客户端流(Stream Input, Unary Output)

反过来,“多次输入、单次结果”的场景中,Request 类型是一个 ReadableStream,handler 以异步迭代器消费它后返回单个结果:

const recordRoute = defineInvokeEventa<
  { distance: number, points: number },
  ReadableStream<{ lat: number, lng: number }>
>('rpc:record-route')

defineInvokeHandler(ctx, recordRoute, async (stream) => {
  let points = 0
  for await (const _ of stream) points += 1
  return { distance: points * 10, points }
})

const invoke = defineInvoke(ctx, recordRoute)
const input = new ReadableStream({
  start(c) { c.enqueue({ lat: 0, lng: 0 }); c.enqueue({ lat: 1, lng: 1 }); c.close() },
})
await invoke(input)

注意客户端流的 handler 走 defineInvokeHandler 而非 defineStreamInvokeHandler:因为响应是单值的,只有请求侧是流。

双向流(Bidirectional Streaming)

Request 与 Response 都是 ReadableStream 时构成双向流:handler 的生成器参数即为入站流,yield 即出站帧:

const routeChat = defineInvokeEventa<
  { message: string },
  ReadableStream<{ message: string }>
>('rpc:route-chat')

defineStreamInvokeHandler(ctx, routeChat, async function* (incoming) {
  for await (const note of incoming) {
    yield { message: `echo: ${note.message}` }
  }
})

const stream = defineStreamInvoke(ctx, routeChat)
for await (const note of stream(outgoing)) {
  console.log(note.message)
}

中止与取消(Abort/Cancel)

取消分客户端与服务端两侧。客户端通过 AbortController 传入 signal;服务端 handler 的第二个参数 options 携带 abortController,可检查 signal.aborted 或监听 abort 事件做清理:

// Client-side cancellation
const controller = new AbortController()
const promise = invokeMethod({ input: 'work' }, { signal: controller.signal })
controller.abort('user cancelled')

// Server-side abort awareness
defineInvokeHandler(ctx, event, async ({ input }, options) => {
  const signal = options?.abortController?.signal
  if (signal?.aborted) return { output: 'aborted' }
  signal?.addEventListener('abort', () => { /* cleanup */ }, { once: true })
  return { output: `done: ${input}` }
})

多跳频道(Multi-hop Channels)

频道(channel)构成有序的路由链,事件、一元 invoke、所有流帧以及调用取消都会经由中间上下文转发:

import { linkChannel, pipeChannel } from '@moeru/eventa'

pipeChannel(a, b, c) // a -> b -> c
linkChannel(a, b, c) // a <-> b <-> c

使用多跳频道时有几条必须掌握的规则(直接来自技能文档):

  • 不存在 ac 的直连边,需要扇出时应使用多条显式 pipe;
  • 销毁(dispose)某个频道只移除它自己的边,context 的 abort 绝不会沿 link 级联传播;
  • 一个连通图中,每个 invoke 定义只能有一个有效 handler

底层机制上,每次本地 emit() 都会创建 EventaInner,其 deliveryId 能够跨越频道跳数和传输层序列化存活;context 会对近期见过的 delivery ID 做去重抑制,并在 hopsRemaining 归零时停止转发。插件可以检查只读的 inner 值并转换或丢弃其 Eventa,但不得替换路由身份或跳数状态

对于 iframe 到服务器的路由场景,推荐做法是:将 EventTarget 侧的 context 与插件的 BroadcastChannel context 连接,再将网关的 BroadcastChannel context 与其 WebSocket context 连接。适配器负责把 inner 值携带过运行时边界,无需方向性转发标记。

最后两条并发语义值得注意:context 不会对并发的 emit() 调用做串行化;请求/响应流的 pump 逐帧 await 只为保证单次调用内的流顺序;而取消是独立路由的,可能先于请求帧到达——handler 因此必须能容忍乱序的取消。

批量注册(Shorthands)

当一组 invoke 契约集中定义时,可用批量 API 一次性注册 handler 或创建 invoke 函数:

const events = {
  double: defineInvokeEventa<number, number>(),
  append: defineInvokeEventa<string, string>(),
}

defineInvokeHandlers(ctx, events, {
  double: input => input * 2,
  append: input => `${input}!`,
})

const { double, append } = defineInvokes(ctx, events)

适配器(Adapters):把任意传输层包装成 Eventa Context

每个适配器把特定传输层包装成一个 eventa context,使用模式完全一致:

import { createContext } from '@moeru/eventa/adapters/<adapter-name>'
const { context } = createContext(transportInstance)

完整适配器列表如下:

Adapter Import Path Transport
Electron Main @moeru/eventa/adapters/electron/main ipcMain + webContents
Electron Renderer @moeru/eventa/adapters/electron/renderer ipcRenderer
Web Worker (main) @moeru/eventa/adapters/webworkers Worker instance
Web Worker (worker) @moeru/eventa/adapters/webworkers/worker self (worker global)
Worker Threads (main) @moeru/eventa/adapters/worker-threads Node.js Worker
Worker Threads (worker) @moeru/eventa/adapters/worker-threads/worker parentPort
WebSocket Client @moeru/eventa/adapters/websocket/native WebSocket
WebSocket Server (H3) @moeru/eventa/adapters/websocket/h3 H3 WebSocket hooks
BroadcastChannel @moeru/eventa/adapters/broadcast-channel BroadcastChannel
EventTarget @moeru/eventa/adapters/event-target EventTarget
EventEmitter @moeru/eventa/adapters/event-emitter Node.js EventEmitter

airi 桌面端的 Electron 适配器实战

airi 桌面应用 stage-tamagotchi 的主进程按窗口拆分 RPC 服务。以主窗口为例(apps/stage-tamagotchi/src/main/windows/main/rpc/index.electron.ts):

import { defineInvokeHandler } from '@moeru/eventa'
import { createContext } from '@moeru/eventa/adapters/electron/main'
import { ipcMain } from 'electron'

const { context } = createContext(ipcMain, params.window)

// 各服务围绕同一个 context 注册 invoke handler
createWidgetsService({ context, widgetsManager: params.widgetsManager, window: params.window })
createAutoUpdaterService({ context, window: params.window, service: params.autoUpdater })
createMcpServersService({ context, manager: params.mcpStdioManager })

defineInvokeHandler(context, electronCenterMainWindow, () => centerWindowOnDisplay(params.window))

可以看到 main 侧 createContext(ipcMain, window) 返回的 context 被当作依赖注入进各个 service,每个 service 再基于共享契约(来自 apps/stage-tamagotchi/src/shared/eventa/ 下的 host.tsplugin/ 等模块)用 defineInvokeHandler 注册处理器。源码中还留有一条值得注意的 TODO:当前在 ipcMain.setMaxListeners(0) 上做了兜底,注释说明“eventa 重构支持按窗口命名空间 context 后即可移除”——从源码结构看,多窗口共享 ipcMain 的监听上限管理是现存的一个演进方向。

渲染进程侧则由 packages/electron-vueuse/src/composables/use-electron-eventa-context.ts 提供 Vue 组合式封装:getElectronEventaContext 用单例模式缓存 createContext(ipcRenderer).context(可从参数或 window.electron.ipcRenderer 自动解析),useElectronEventaInvoke 则在共享 context 之上 defineInvoke 出类型化的 invoke 函数,供组件直接 await 调用。

技能文档给出的最小 Electron 示例可概括为“三处代码、一份契约”:

// shared/events.ts — define events once
import { defineInvokeEventa } from '@moeru/eventa'
export const readdir = defineInvokeEventa<{ dirs: string[] }, { path: string }>('fs:readdir')

// main.ts — register handler
import { createContext } from '@moeru/eventa/adapters/electron/main'
const { context } = createContext(ipcMain, mainWindow.webContents)
defineInvokeHandler(context, readdir, async ({ path }) => ({ dirs: await fs.readdir(path) }))

// renderer.ts (preload) — call it
import { createContext } from '@moeru/eventa/adapters/electron/renderer'
const { context } = createContext(ipcRenderer)
const invokeReaddir = defineInvoke(context, readdir)
const result = await invokeReaddir({ path: '/usr' })

高级特性

技能文档列出了三个进阶能力:

  • Delivery 路由EventaInner<T> 在频道与适配器之间保持投递身份(delivery identity)与跳数预算(hop budget),是多跳转发与去重的底层支撑;
  • 匹配表达式matchBy(glob)matchBy(regex)and(...)or(...) 用于事件过滤,适合监听器按名称模式批量订阅;
  • WebSocket 生命周期:native 适配器额外导出 wsConnectedEventwsDisconnectedEvent,可把连接状态变化当作普通事件消费。

关键规则(Key Rules)

使用 Eventa 的五条纪律,直接决定代码质量:

  1. 事件定义必须放在共享模块 —— 双方导入同一事件定义以获得类型安全;
  2. defineInvokeEventa<Res, Req>() —— Response 类型在前,Request 类型在后
  3. Handler 可以安全地抛错 —— eventa 会把错误传播回调用方;
  4. 在边界处校验数据 —— eventa 转发你 emit 的任何负载,它不做校验;
  5. 只安装你需要的 peer 依赖 —— electron、h3、web-worker 等均为可选依赖,按需安装。

小结

@moeru/eventa 在 airi 中承担的角色是“一套类型化事件原语 + 可替换传输层 + 由此组合出的 RPC/流式调用”,它让 Electron 主进程与渲染进程、浏览器 Worker、服务端 WebSocket 之间共享同一份契约代码。结合 packages/electron-eventa/src/electron/window.ts 这类共享契约模块与 apps/stage-tamagotchi/src/main/windows/main/rpc/index.electron.ts 的 service 化注册方式,可以直接照搬到自己的桌面/实时应用中:先定义共享事件,再按运行环境选择适配器创建 context,最后用 invoke/stream 原语编写调用与处理器。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
898
5.82 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
921
1.84 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.8 K
1.02 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
531
596
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.02 K
519
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.36 K
1.46 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
391