Novu Cloud 实时通道:基于 Cloudflare Workers 与 Durable Objects 的 @novu/socket-worker 架构与本地开发指南
【免费下载链接】novuThe open-source communication infrastructure for agents and products项目地址: https://gitcode.com/GitHub_Trending/no/novu
本篇文章围绕 Novu 仓库中 enterprise/workers/socket/README.md 展开,系统讲解@novu/socket-worker——一个承载 Novu Cloud WebSockets(PartySocket)实时通道的 Cloudflare Worker + Durable Object 服务。你将掌握它的整体架构(Hono 路由、Durable Object 房间模型、EU 数据驻留)、本地联调方法(8787 端口、.dev.vars、与 API/Worker/Playground 的环境接线),以及 JWT 认证、内部 API 鉴权、WebSocket Hibernation 与 contextKeys 精确匹配等核心实现原理,可直接上手在本地跑通整套实时链路。
一、socket-worker 在 Novu 实时体系中的角色
Novu 的开源实时通道由apps/ws(Node.js 网关,基于 Socket.IO/PartySocket 生态)与本次要讲的@novu/socket-worker组成。后者定位为Novu Cloud 的 WebSocket 实时路径:当NOVU_ENTERPRISE=true时,Cloud 环境下的实时消息不再走传统 Node 网关,而是由 Cloudflare 边缘上的 Worker + Durable Object 承担 WebSocket 升级与消息投递,也就是 README 中标注的 "Cloudflare Worker + Durable Object for Novu Cloud WebSockets (PartySocket)"。
从 package.json 可以看到,该 Worker 的运行时依赖非常精简:
hono:HTTP 路由框架,处理 WebSocket 升级、内部消息接口与健康检查;ws/@types/jsonwebtoken:WebSocket 与 JWT 相关类型支撑;@tsndr/cloudflare-worker-jwt:在 Worker 运行时内完成 JWT 签名校验(Cloudflare Workers 无 Node 原生crypto完整 API,故使用专为 Worker 设计的 JWT 库);wrangler(^4.49.0):本地开发与多环境部署 CLI。
二、整体架构与请求路径
Worker 的入口是 src/index.ts,一个典型的 Hono 应用,对外暴露三条路由:
| 路由 | 方法 | 中间件 | 职责 |
|---|---|---|---|
/ | GET | authenticateJWT | 携带?token=发起 WebSocket 升级 |
/send | POST | authenticateInternalAPI | 内部服务向指定用户房间推送事件 |
/health | GET | 无 | 健康检查,返回OK |
未命中路由统一返回 404Not found;应用级错误经app.onError记录日志并返回 500。
2.1 WebSocket 升级链路
当客户端(如 playground/web-chat 中的 PartySocket 客户端)发起连接时,请求经过 middleware/auth.ts 校验?token=中的 JWT,然后进入 handlers/websocket.ts 的handleWebSocketUpgrade:
- 从 JWT payload 与请求中取出
userId、subscriberId、organizationId、environmentId、contextKeys; - 计算房间 ID:
roomId = ${environmentId}:${userId},即每个用户在其环境内拥有一个专属房间; - 根据
REGION变量决定是否使用 EU 数据驻留命名空间(WEBSOCKET_ROOM.jurisdiction('eu')); - 通过
idFromName(roomId)定位 Durable Object 实例并stub.fetch(...),把用户信息以X-User-Id、X-Subscriber-Id、X-Organization-Id、X-Environment-Id、X-JWT-Token、X-Context-Keys等自定义头透传给 DO。
2.2 消息推送链路
API/Worker 需要给在线用户推送实时事件时,向/send发起 POST,请求体结构(由handleSendMessage校验逻辑确认):
{ "userId": "用户 ID(字符串)", "environmentId": "环境 ID(字符串)", "event": "事件名,如 notification.inbox_received", "data": "任意业务负载", "contextKeys": ["可选", "上下文键数组"] }handleSendMessage会依次校验:userId与event必填、environmentId必填、三者必须为字符串;随后同样按environmentId:userId定位 Durable Object,并通过context.executionCtx.waitUntil(stub.sendToUser(...))异步投递,接口立即返回{ success: true, roomId, timestamp }。
三、本地开发环境搭建
该包属于 pnpm workspace(见根目录 pnpm-workspace.yaml),因此依赖安装统一在仓库根目录执行:
pnpm install启动开发服务器有两种等价方式:
# 方式一:仓库根目录运行(利用 workspace filter) pnpm dev:socket-worker # 方式二:进入本包目录直接运行 cd enterprise/workers/socket pnpm run dev根目录 package.json 中dev:socket-worker定义为pnpm --filter @novu/socket-worker dev,即pnpm run dev执行的是wrangler dev --env local。
端口约定(务必遵守):socket-worker 固定运行在http://127.0.0.1:8787,而本地thalamus-observer(另一 Cloudflare Worker,见 scripts/dev-environment-setup.sh 相关脚本)使用8788,两者互不冲突。若修改 socket-worker 端口,需同步修改下游所有指向它的环境变量。
3.1 首次运行前的密钥配置
.dev.vars被 gitignore,首次需要从模板复制并填写:
cd enterprise/workers/socket cp .dev.vars.example .dev.vars然后参考 .dev.vars.example 中的注释,从apps/api/src/.env取值:
JWT_SECRET=<apps/api/src/.env 中的 JWT_SECRET> INTERNAL_API_KEY=<apps/api/src/.env 中的 INTERNAL_SERVICES_API_KEY>关键约束:INTERNAL_API_KEY必须与 API 侧的INTERNAL_SERVICES_API_KEY完全一致,因为/send接口的调用方(API/Worker 内部服务)正是用它作为 Bearer 凭证;不一致将导致推送被 401 拒绝。JWT_SECRET则用于校验客户端 WebSocket 升级时携带的?token=签名,必须与签发 token 的 API 侧共享同一密钥。
localwrangler 环境会将API_URL默认设置为http://127.0.0.1:3000(见 wrangler.jsonc),用于 Durable Object 回拨 API 上报用户在线状态。
四、打通 API / Worker / Playground 的完整环境接线
README 给出了三条链路的联调配置。要让本地整套系统(API + Worker + Playground)都走 Cloudflare socket,需要:
1) API 与 Worker 侧(apps/api/src/.env和apps/worker/src/.env):
SOCKET_WORKER_URL=http://127.0.0.1:8787 NOVU_ENTERPRISE=true # INTERNAL_SERVICES_API_KEY 保持与 .dev.vars 的 INTERNAL_API_KEY 相同2) Playground 侧(playground/web-chat):
NEXT_PUBLIC_NOVU_SOCKET_URL=http://127.0.0.1:8787 NEXT_PUBLIC_NOVU_SOCKET_TYPE=cloud其中NEXT_PUBLIC_NOVU_SOCKET_TYPE=cloud指示 Playground 走 Cloudflare 云端 socket 路径而非本地 Node 网关——这正是 README 强调的 "Locally, Cloudflare sockets are the realtime path" 的含义。配置完成后,Playground 发起的 WebSocket 升级请求会携带 JWT 直连 8787 端口的 Worker。
五、wrangler 多环境配置与部署
wrangler.jsonc 定义了四个环境,差异点集中在名称、路由域名、Durable Object 绑定与变量:
| 环境 | Worker 名称 | 自定义域名 | API_URL | REGION |
|---|---|---|---|---|
local | socket-worker-local | 无(workers_dev) | http://127.0.0.1:3000 | global |
staging | socket-worker-staging | socket.novu-staging.co | https://api.novu-staging.co | global |
production-us | socket-worker-production-us | socket.novu.co | https://api.novu.co | global |
production-eu | socket-worker-production-eu | eu.socket.novu.co | https://eu.api.novu.co | eu |
所有环境都声明了同一个 Durable Object 绑定WEBSOCKET_ROOM → WebSocketRoom,并在migrations中以new_sqlite_classes: ["WebSocketRoom"](tagv1)注册(SQLite 后端 DO)。此外observability.enabled: true开启了 Cloudflare 观测。
对应 package.json 中的部署脚本:
pnpm run deploy # wrangler deploy(默认环境) pnpm run deploy:staging pnpm run deploy:production-us pnpm run deploy:production-eu pnpm run deploy:local pnpm run cf-typegen # 从 wrangler 配置生成 Worker 类型声明REGION=eu在生产 EU 环境中的作用非常关键:DO 命名空间会调用jurisdiction('eu')将连接与数据锁定在欧盟境内,以满足数据驻留要求(见下一节源码说明)。
六、核心实现原理:WebSocketRoom Durable Object
Durable Object 的实现集中在 src/durable-objects/websocket-room.ts,类WebSocketRoom是整条实时链路的"房间"载体。以下是几个值得深入理解的设计点。
6.1 基于 Hibernation API 的连接管理
构造函数中通过this.ctx.setWebSocketAutoResponse(new WebSocketRequestResponsePair('ping', 'pong'))配置自动心跳应答,使 Worker 可在空闲时休眠而连接保持存活。WebSocket 接受使用hibernation 兼容方式:
const tags = [`user:${userId}`, `env:${environmentId}`]; this.ctx.acceptWebSocket(server, tags); server.serializeAttachment({ jwtToken, connectedAt: Date.now(), contextKeys });- tags:给连接打上
user:与env:标签,后续可定向按标签检索连接; - serializeAttachment:把 JWT、连接时间、contextKeys 持久化到连接附件上(附件上限 2KB,JWT 通常 <1KB),这样 DO 休眠唤醒后依然能恢复连接元数据——这正是代码注释强调"No need to store JWT tokens in memory"的原因。
运行时通过三个钩子驱动生命周期:webSocketMessage(收到客户端消息,此处仅校验元数据存在性)、webSocketClose(关闭连接并触发下线上报)、webSocketError(记录错误日志)。
6.2 房间容量与并发保护
每个 DO 实例设定了MAX_CONNECTIONS = 100的硬上限。fetch在升级前检查this.ctx.getWebSockets().length,达到上限返回503 +Retry-After: 60,让客户端 60 秒后重试。此外还暴露了三个统计方法:getActiveConnectionsForUser、getTotalActiveConnections、getConnectionCapacity(返回{ current, max, available }),接口契约定义在 src/types/index.ts。
6.3 contextKeys 精确匹配:多上下文隔离投递
一个用户可能同时打开多个"上下文"(例如不同的工作区页面),/send携带的contextKeys与连接附件中的contextKeys做精确匹配后才投递。isExactMatch的规则是:
- 消息 contextKeys 为空数组 → 仅投递给 contextKeys 也为空的连接;
- 长度不一致 → 不投递;
- 否则逐个成员比较(
every(key => inboxContextKeys.includes(key))),全部命中才投递。
代码注释明确指出这套逻辑与ws.gateway.ts保持一致,保证了新旧实时通道在语义上的兼容。投递过程对消息体只做一次JSON.stringify预序列化({ event, data, timestamp }),再对命中连接并行发送,并用Promise.allSettled容错。
6.4 在线状态回拨 API
连接建立与断开时,DO 都会通过notifySubscriberOnlineState向API_URL的POST /v1/internal/subscriber-online-state上报(请求体含subscriberId、environmentId、isOnline、organizationId、timestamp,以Bearer ${jwtToken}鉴权)。该调用使用ctx.waitUntil包裹,确保 DO 可立即进入休眠而不阻塞连接建立;只有所有同用户连接都断开(剩余连接数 ≤ 0)时才上报离线。
七、安全模型:双层认证
7.1 客户端侧:JWT 校验
middleware/auth.ts 中authenticateJWT负责升级请求认证:
- 从
?token=查询参数取 JWT,缺失返回 401; - 用
JWT_SECRET通过@tsndr/cloudflare-worker-jwt验签并解码; - 从 payload 提取
_id(作为 userId)、subscriberId(缺省回退为 userId)、organizationId、environmentId、contextKeys,任一缺失返回 401; - 认证信息通过 Hono
context.set注入后续处理。
7.2 内部侧:常量时间比较
middleware/internal-auth.ts 保护/send:请求头需携带Authorization: Bearer <key>,与INTERNAL_API_KEY做常量时间比较(constantTimeEquals,逐字符异或累加,长度不等直接失败),从实现层面抵御时序侧信道攻击。若服务端未配置INTERNAL_API_KEY则返回 500。
八、类型契约与可观测性
src/types/index.ts 定义了完整的环境与元数据契约:
IEnv:WEBSOCKET_ROOM(DO 命名空间绑定)、JWT_SECRET、INTERNAL_API_KEY(必填),API_URL、REGION(可选);IConnectionMetadata:userId、environmentId、connectedAt、jwtToken、contextKeys——即 serializeAttachment 持久化的全部字段;IWebSocketRoom:sendToUser、连接数查询与容量查询接口。
运维层面,/health提供存活探针;wrangler 配置开启了observability;app.onError统一记录应用错误。日常排障可结合 Cloudflare Dashboard 的 Worker 日志查看[Internal API] Routing message to room: ...等关键日志行。
九、小结
@novu/socket-worker展示了如何用 Cloudflare Workers + Durable Objects 构建生产级实时通道:边缘就近升级 WebSocket、以environmentId:userId为粒度分房、Hibernation API 降本增效、JWT + 内部 API Key 双层鉴权、EU 数据驻留按区域隔离。对于希望深度定制 Novu Cloud 实时链路或自建同类基础设施的开发者,建议从 enterprise/workers/socket/src/index.ts(路由入口)→ src/handlers/websocket.ts(升级与推送)→ src/durable-objects/websocket-room.ts(房间核心)这条调用链入手阅读,再结合 wrangler.jsonc 完成本地与多环境部署验证。
【免费下载链接】novuThe open-source communication infrastructure for agents and products项目地址: https://gitcode.com/GitHub_Trending/no/novu
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考