airi 中的 @moeru/eventa:跨进程事件通信与类型安全 RPC 实战指南
本篇技术指南围绕 airi 仓库中的 @moeru/eventa 技能文档(.agents/skills/eventa/SKILL.md)展开,系统讲解这个“传输感知的事件库”如何支撑 Electron IPC、WebSocket、Web Worker、BroadcastChannel 等多种进程间通信,并演示一元 RPC、三种流式调用模式、取消机制与多跳频道路由。读完后,你将能够复用 airi 中的真实代码模式,为任意 Electron/Node/浏览器应用定义类型安全的事件与 RPC 接口。
核心设计思想
Eventa 的整个 API 围绕三个基本理念构建,理解它们是掌握这个库的前提:
- 事件是一等公民 —— 定义一次带类型的事件,到处复用;
- 传输层可替换 —— 同一份事件定义可以在 Electron IPC、WebSocket、Web Workers、BroadcastChannel、EventEmitter、EventTarget 和 Worker Threads 之间无缝切换;
- RPC 只是事件的组合 —— invoke(一元调用)和 stream(流式调用)模式都从同一套事件原语组合而来,而不是独立的一套机制。
这种设计在 airi 中得到了大量验证:当前仓库通过 pnpm-workspace.yaml 将 @moeru/eventa 锁定在 catalog 版本 1.0.0,并被 packages/electron-eventa、packages/better-ws、server/packages/server-sdk-shared、packages/plugin-sdk、packages/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 为音频流水线集中定义了一组事件,如 speechSegmentEvent、speechTtsRequestEvent、speechPlaybackStartEvent,并将它们聚合成 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
使用多跳频道时有几条必须掌握的规则(直接来自技能文档):
- 不存在
a到c的直连边,需要扇出时应使用多条显式 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.ts、plugin/ 等模块)用 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 适配器额外导出
wsConnectedEvent与wsDisconnectedEvent,可把连接状态变化当作普通事件消费。
关键规则(Key Rules)
使用 Eventa 的五条纪律,直接决定代码质量:
- 事件定义必须放在共享模块 —— 双方导入同一事件定义以获得类型安全;
defineInvokeEventa<Res, Req>()—— Response 类型在前,Request 类型在后;- Handler 可以安全地抛错 —— eventa 会把错误传播回调用方;
- 在边界处校验数据 —— eventa 转发你 emit 的任何负载,它不做校验;
- 只安装你需要的 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 原语编写调用与处理器。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0629
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python07
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00