Rivet Actors 实时协作光标实战:用 Raw WebSocket 处理器构建自定义协议的多人画布
【免费下载链接】actorsRivet Actors are the primitive for stateful workloads. Built for AI agents, collaborative apps, and durable execution.项目地址: https://gitcode.com/GitHub_Trending/riv/actors
本篇技术指南以仓库中的 cursors-raw-websocket 示例 为主体,讲解如何绕开 RivetKit 高层 WebSocket 抽象,直接使用onWebSocket底层处理器实现实时光标追踪与多人协作画布。你将掌握 actor 内自定义 WebSocket 协议设计、多客户端广播、持久化画布状态与多房间隔离的完整落地方案,并能直接复制到自己的 Rivet Actors 项目中运行。
示例概览与核心特性
该示例演示了基于 Rivet Actors 的实时协作画布:多个用户同时进入同一房间,彼此可以看到对方的实时光标位置,并可在画布任意位置添加文本标签,标签会持久化保存。它区别于使用 RivetKit 高层onConnect/onMessage抽象的示例(如 cursors 示例),全部交互都建立在原生的 WebSocket 消息之上。
示例(package.json)定位为example-cursors-raw-websocket,核心依赖为rivetkit与@rivetkit/react,服务端与前端同仓开发。其四大特性如下:
- Raw WebSocket handlers:在 actor 上使用
onWebSocket获得底层 WebSocket 控制权,可实现自定义协议、二进制流或任意消息格式; - Real-time cursor tracking:光标位置实时广播给房间内所有已连接用户;
- Persistent canvas state:文本标签自动写入 actor state,跨会话持久保存;
- Multiple rooms:每个房间都是一个独立的 actor 实例,状态互相隔离。
快速开始
克隆仓库并进入示例目录后,按以下步骤启动:
git clone https://github.com/rivet-dev/rivet.git cd rivet/examples/cursors-raw-websocket npm install npm run devnpm run dev会通过 concurrently 同时启动两个进程:
- 服务端:
tsx --watch src/index.ts,监听localhost:6420,运行 actor 注册表; - 前端:
vite,提供 React 界面,并将/actors、/metadata、/health等路径代理到服务端(见 vite.config.ts)。
浏览器打开 Vite 输出的地址即可进入默认房间general。可以开两个浏览器窗口对比观察光标同步效果。
Actor 端实现:Raw WebSocket 处理器
示例的 actor 定义在 src/index.ts(README 中提到的src/backend/registry.ts在本仓库中对应此文件)。整个实现只有一个cursorRoomactor,结构如下:
export const cursorRoom = actor({ state: { textLabels: [] as TextLabel[], }, createVars: (): Vars => { return { websockets: new Map(), }; }, actions: { getOrCreate: () => ({ status: "ok" }), getRoomState: (c) => { /* 返回当前所有光标与文本标签 */ }, }, onWebSocket: async (c, websocket: UniversalWebSocket) => { /* 核心处理逻辑 */ }, }); export const registry = setup({ use: { cursorRoom }, }); registry.start();onWebSocket 在 RivetKit 中的定义
onWebSocket是 RivetKit actor 配置的顶层生命周期钩子,其类型定义位于 rivetkit-typescript/packages/rivetkit/src/actor/config.ts:
/** * Called when a raw WebSocket connection is established to the actor. * * This handler receives WebSocket connections made to `/actors/{actorName}/websocket/*` endpoints. * Use this hook to handle custom WebSocket protocols, binary streams, or other WebSocket-based communication. * * @param c The WebSocket context with access to the connection * @param websocket The actor-facing raw WebSocket connection * @param opts Additional options including the original HTTP upgrade request */ onWebSocket?: ( c: WebSocketContext<...>, websocket: UniversalWebSocket, ) => void | Promise<void>;几点关键信息:
- 该钩子接收的是打到
/actors/{actorName}/websocket/*端点的原生 WebSocket 连接; - 第二个参数
websocket的类型是UniversalWebSocket,提供send()与addEventListener("message" | "close", ...)等与标准 WebSocket 近似的接口,其抽象定义见 rivetkit-typescript/packages/rivetkit/src/common/websocket-interface.ts; - 处理器可以返回
void或Promise<void>,因此既支持同步逻辑,也支持在async中完成握手后再发消息(仓库的驱动测试套件中rawWebSocketAsyncOpenActor即演示了这种异步打开模式,见 raw-websocket.ts)。
在 actor 配置的 zod schema 中,onWebSocket也被描述为“Called for raw WebSocket connections to/actors/{name}/websocket/*endpoints”(见 config.ts),与上方的类型注释互相印证。
连接握手:校验 sessionId 并发送初始状态
onWebSocket处理器第一步是校验请求并识别客户端身份:
onWebSocket: async (c, websocket: UniversalWebSocket) => { if (!c.request) { websocket.close(1008, "no request"); return; } const url = new URL(c.request.url); const sessionId = url.searchParams.get("sessionId"); if (!sessionId) { websocket.close(1008, "Missing sessionId"); return; } // ...c.request是 HTTP 升级时的原始请求,可通过new URL(c.request.url)解析查询参数;- 连接必须携带
sessionId,缺失时以 WebSocket 关闭码1008(Policy Violation)拒绝,避免无效连接占用资源; - 通过后,将
{ socket: websocket, cursor: null }存入c.vars.websockets的 Map,Map以sessionId为键。
新连接接入后立即收到一条init消息,包含当前房间全部光标与文本标签,用于前端一次性渲染初始画面:
const cursors: Record<string, CursorPosition> = {}; for (const [id, { cursor }] of c.vars.websockets.entries()) { if (cursor) { cursors[id] = cursor; } } websocket.send( JSON.stringify({ type: "init", data: { cursors, textLabels: c.state.textLabels }, }), );消息分发:三种入站消息
所有入站消息都挂在websocket.addEventListener("message", ...)上,统一按message.type分发,消息体为 JSON 字符串:
| 消息类型 | 数据载荷 | 服务端行为 |
|---|---|---|
updateCursor | { userId, x, y } | 更新该会话光标,向所有连接广播cursorMoved |
updateText | { id, userId, text, x, y } | 按id创建或更新文本标签,写入 actor state,广播textUpdated |
removeText | { id } | 从 actor state 中删除该标签,广播textRemoved |
其中updateCursor的完整逻辑为:从消息中取出userId/x/y,加上Date.now()时间戳构造成CursorPosition,写回当前会话的cursor字段,再遍历c.vars.websockets向包括发送者在内的所有连接广播:
case "updateCursor": { const { userId, x, y } = message.data; const cursor: CursorPosition = { userId, x, y, timestamp: Date.now(), }; const session = c.vars.websockets.get(sessionId); if (session) { session.cursor = cursor; } for (const { socket } of c.vars.websockets.values()) { socket.send(JSON.stringify({ type: "cursorMoved", data: cursor, })); } break; }注意这里cursorMoved广播给发送者自身,前端在收到自身消息时同样更新本地状态,因此“移动即所见”不需要额外的回显确认,实现简洁。
updateText的分支体现了 actor state 的读写方式:先按id在c.state.textLabels中查找,存在则替换、不存在则追加,随后广播;removeText则直接用filter从 state 中移除。由于 state 由 RivetKit 持久化,文本标签天然具备跨会话持久能力。
连接断开:清理会话并广播 cursorRemoved
close事件处理器负责两件事:若断开前该会话有光标位置,则向其他连接广播cursorRemoved(注意此处遍历时通过id !== sessionId排除了发送者自身);随后从c.vars.websockets中删除该会话:
websocket.addEventListener("close", () => { const session = c.vars.websockets.get(sessionId); if (session?.cursor) { for (const [id, { socket }] of c.vars.websockets.entries()) { if (id !== sessionId) { socket.send(JSON.stringify({ type: "cursorRemoved", data: session.cursor, })); } } } c.vars.websockets.delete(sessionId); });自定义协议消息总览
综合服务端与前端(frontend/App.tsx),整个示例的 WebSocket 协议可归纳为下表,这正是 README 所强调的“custom WebSocket protocol implementation”:
| 方向 | 类型 | 数据载荷 | 说明 |
|---|---|---|---|
| 服务端 → 客户端 | init | { cursors, textLabels } | 新连接初始化全量状态 |
| 客户端 → 服务端 | updateCursor | { userId, x, y } | 上报光标位置 |
| 服务端 → 客户端 | cursorMoved | CursorPosition | 广播光标移动(含发送者) |
| 服务端 → 客户端 | cursorRemoved | CursorPosition | 广播某会话光标下线 |
| 客户端 → 服务端 | updateText | { id, userId, text, x, y } | 创建/更新文本标签 |
| 服务端 → 客户端 | textUpdated | TextLabel | 广播文本标签变更 |
| 客户端 → 服务端 | removeText | { id } | 删除文本标签 |
| 服务端 → 客户端 | textRemoved | id | 广播文本标签被删除 |
两侧的数据结构完全一致:CursorPosition为{ userId, x, y, timestamp },TextLabel为{ id, userId, text, x, y, timestamp },类型定义见 src/index.ts,前端直接通过import type { CursorPosition, registry, TextLabel } from "../src/index.ts"复用同一份类型,保证协议前后端零漂移。
前端实现:连接、坐标映射与渲染
建立连接
前端通过@rivetkit/react的createClient获取类型安全的客户端(App.tsx):
const rivetUrl = "http://localhost:6420"; const rivetWsUrl = "ws://localhost:6420"; const client = createClient<typeof registry>(rivetUrl);连接流程分为两步:
- 先通过 action 解析/创建房间对应的 actor:
const actorId = await client.cursorRoom.getOrCreate(roomId).resolve();。getOrCreate正是 src/index.ts 中定义的 action,其作用是由房间名确定稳定的 actor ID——房间名不同则 actor 不同,从而天然实现房间隔离; - 再直接以原生
WebSocket连接到 Raw WebSocket 端点,并通过 URL 查询参数携带sessionId:
const wsUrl = `${rivetWsUrl}/gateway/${actorId}/websocket?sessionId=${encodeURIComponent(sessionId)}`; ws = new WebSocket(wsUrl);这段代码体现了“Raw WebSocket”与高层抽象的本质区别:这里不依赖@rivetkit/react的useActor/useConnection钩子,而是手动拼接连接地址、手动发送 JSON、手动维护连接状态。sessionId与userId在组件初始化时随机生成(App.tsx),其中userId用于消息负载与光标着色,sessionId用于服务端 Map 索引。
消息处理与状态更新
前端为每个消息类型维护独立的 React state 分支,与服务端协议一一对应:
init:一次性写入cursors与textLabels;cursorMoved:以message.data.userId为键更新光标集合;cursorRemoved:删除对应userId的光标;textUpdated:按id更新或追加文本标签;textRemoved:过滤掉对应id。
由于光标集合的键是userId而服务端 Map 的键是sessionId,同一用户在重连后会生成新的sessionId,但userId保持不变,因此前端的cursors状态不会因重连产生重复条目——这是该示例在状态建模上一个值得借鉴的设计。
画布坐标与交互
画布是一个固定逻辑尺寸(1920x1080)的虚拟空间,通过scale因子等比缩放适配视口(App.tsx)。所有交互都经过screenToCanvas将屏幕坐标换算为画布坐标后发送:
onMouseMove:发送updateCursor;onClick:在点击处进入输入态;- 输入过程中实时发送
updateText(边打字边同步),按Enter提交并结束编辑,空文本或按Escape/ 失焦则发送removeText撤销该标签。
每个用户的颜色由getColorForUser对userId做字符串哈希后从固定调色板中选取,保证同一用户在不同窗口颜色一致(App.tsx)。光标以 SVG 箭头呈现,并显示“you”或对应用户 ID 的标签,方便多人协作时辨认身份。
持久画布状态与多房间隔离
示例的两个进阶特性都源于 actor 的运行时模型:
- 持久画布状态:
textLabels声明在 actor 的state字段中(src/index.ts),由 RivetKit 负责持久化。任何客户端刷新或重连后,通过init消息拿到的是服务端最新 state,因此已添加的文本标签不会因会话中断而丢失; - 多房间隔离:前端以
roomId调用getOrCreate(roomId)解析 actor,不同房间名对应不同的 actor 实例,各自的state与vars.websockets完全隔离。这也印证了 README 中“Each room is a separate actor instance with isolated state”的描述。
与光标不同,cursors只存放在c.vars.websockets(进程内内存)中,不进入持久化 state:光标是瞬态信息,随连接生命周期存在,断开即广播移除,这是对 actorstate(持久)与vars(易失)两种存储的合理分工。
测试验证
示例内置了基于 vitest 的测试(tests/cursors.test.ts),通过rivetkit/test的setupTest在测试环境中加载registry并验证 actor 行为:
test("Cursor room can be created and initialized", async (ctx: any) => { const { client } = await setupTest(ctx, registry); const room = client.cursorRoom.getOrCreate(["test-room"]); const result = await room.getOrCreate(); expect(result).toEqual({ status: "ok" }); }); test("Cursor room can get initial room state", async (ctx: any) => { const { client } = await setupTest(ctx, registry); const room = client.cursorRoom.getOrCreate(["test-state"]); const state = await room.getRoomState(); expect(state).toEqual({ cursors: {}, textLabels: [] }); });两个用例分别覆盖:actor 可通过getOrCreate正常创建并响应 action;新房间的初始getRoomState返回空光标集合与空文本标签。运行npm test即可执行。此外,仓库驱动测试套件中的 raw-websocket.test.ts 对onWebSocket的消息回显、二进制处理、rivetMessageIndex消息序号等底层行为做了更全面的验证,可作为深入理解 Raw WebSocket 语义的参考。
总结与延伸阅读
本示例用约 200 行服务端代码加一个 React 组件,完整实现了一个可运行的多人在线协作画布。其核心价值在于展示 Rivet Actors 的onWebSocket钩子如何赋予开发者底层 WebSocket 的完全控制权:自定义消息协议、手动广播、请求级鉴权(如sessionId校验)、以及 actor state 与 vars 的合理分工。当协作功能需要非标准协议、二进制流或精细控制连接生命周期时,Raw WebSocket 处理器是比高层抽象更合适的起点。
进一步阅读仓库内相关资料:
- RivetKit onWebSocket 类型与文档
- UniversalWebSocket 抽象接口
- Raw WebSocket 驱动测试用例
- 原生 WebSocket 高层抽象对照示例
- WebSockets 概念文档、state 概念文档、events 概念文档
该示例以 MIT 协议开源(见 README.md),可直接在自有项目中复用其中的协议设计与连接管理思路。
【免费下载链接】actorsRivet Actors are the primitive for stateful workloads. Built for AI agents, collaborative apps, and durable execution.项目地址: https://gitcode.com/GitHub_Trending/riv/actors
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考