Novu Cloud 实时通道:基于 Cloudflare Workers 与 Durable Objects 的 @novu/socket-worker 架构与本地开发指南
2026/9/10 6:11:31 网站建设 项目流程

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 应用,对外暴露三条路由:

路由方法中间件职责
/GETauthenticateJWT携带?token=发起 WebSocket 升级
/sendPOSTauthenticateInternalAPI内部服务向指定用户房间推送事件
/healthGET健康检查,返回OK

未命中路由统一返回 404Not found;应用级错误经app.onError记录日志并返回 500。

2.1 WebSocket 升级链路

当客户端(如 playground/web-chat 中的 PartySocket 客户端)发起连接时,请求经过 middleware/auth.ts 校验?token=中的 JWT,然后进入 handlers/websocket.ts 的handleWebSocketUpgrade

  1. 从 JWT payload 与请求中取出userIdsubscriberIdorganizationIdenvironmentIdcontextKeys
  2. 计算房间 IDroomId = ${environmentId}:${userId},即每个用户在其环境内拥有一个专属房间;
  3. 根据REGION变量决定是否使用 EU 数据驻留命名空间(WEBSOCKET_ROOM.jurisdiction('eu'));
  4. 通过idFromName(roomId)定位 Durable Object 实例并stub.fetch(...),把用户信息以X-User-IdX-Subscriber-IdX-Organization-IdX-Environment-IdX-JWT-TokenX-Context-Keys等自定义头透传给 DO。

2.2 消息推送链路

API/Worker 需要给在线用户推送实时事件时,向/send发起 POST,请求体结构(由handleSendMessage校验逻辑确认):

{ "userId": "用户 ID(字符串)", "environmentId": "环境 ID(字符串)", "event": "事件名,如 notification.inbox_received", "data": "任意业务负载", "contextKeys": ["可选", "上下文键数组"] }

handleSendMessage会依次校验:userIdevent必填、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/.envapps/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_URLREGION
localsocket-worker-local无(workers_dev)http://127.0.0.1:3000global
stagingsocket-worker-stagingsocket.novu-staging.cohttps://api.novu-staging.coglobal
production-ussocket-worker-production-ussocket.novu.cohttps://api.novu.coglobal
production-eusocket-worker-production-eueu.socket.novu.cohttps://eu.api.novu.coeu

所有环境都声明了同一个 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 秒后重试。此外还暴露了三个统计方法:getActiveConnectionsForUsergetTotalActiveConnectionsgetConnectionCapacity(返回{ 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 都会通过notifySubscriberOnlineStateAPI_URLPOST /v1/internal/subscriber-online-state上报(请求体含subscriberIdenvironmentIdisOnlineorganizationIdtimestamp,以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)、organizationIdenvironmentIdcontextKeys,任一缺失返回 401;
  • 认证信息通过 Honocontext.set注入后续处理。

7.2 内部侧:常量时间比较

middleware/internal-auth.ts 保护/send:请求头需携带Authorization: Bearer <key>,与INTERNAL_API_KEY常量时间比较constantTimeEquals,逐字符异或累加,长度不等直接失败),从实现层面抵御时序侧信道攻击。若服务端未配置INTERNAL_API_KEY则返回 500。

八、类型契约与可观测性

src/types/index.ts 定义了完整的环境与元数据契约:

  • IEnvWEBSOCKET_ROOM(DO 命名空间绑定)、JWT_SECRETINTERNAL_API_KEY(必填),API_URLREGION(可选);
  • IConnectionMetadatauserIdenvironmentIdconnectedAtjwtTokencontextKeys——即 serializeAttachment 持久化的全部字段;
  • IWebSocketRoomsendToUser、连接数查询与容量查询接口。

运维层面,/health提供存活探针;wrangler 配置开启了observabilityapp.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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询