Happy 多副本架构实战:基于 Socket.IO Redis Streams Adapter 的跨 Pod RPC 与广播路由
【免费下载链接】happyMobile and Web client for Codex and Claude Code, with realtime voice, encryption and fully featured项目地址: https://gitcode.com/gh_mirrors/happy20/happy
导读
本文深入剖析 happy 仓库中 happy-server 如何从单副本演进为多 Kubernetes 副本运行:Socket.IO Redis streams adapter 如何把io.to(...).emit(...)在副本之间转发,RPC(Web 客户端 → 本地 daemon)如何通过rpc:<userId>:<method>房间完成跨副本寻址,以及 Pod 被杀、瞬时断连、网络分区等"混乱场景"下的真实行为。读完本文,你将掌握一套不依赖 Redis Key/TTL、无 keep-alive 刷新路径的 Socket.IO 集群路由设计,理解四个历史 Bug 的根因与修复方案,并能用仓库自带的 minikube 集成测试复现和验证全部行为。
这是 docs/multi-process.md 的展开版,高级控制流概览见 docs/realtime-sync-and-rpc.md,推荐先读后者建立整体认知,再回到本文看故障模式与实现细节。
状态说明:本文描述的代码均已合入
main,但生产上是否切换到多副本是独立决策——handy.yaml当前已经配置了replicas: 3(见下文部署章节),是否在生产环境开启由运维决定。
一、为什么需要多副本:单点问题的现实压力
happy-server 承担两类实时流量:realtime sync(服务器把加密消息、会话状态等 update/ephemeral 事件推送给 Web、桌面、移动客户端)和point-to-point RPC(Web 客户端把bash、apply_patch等工具调用路由到用户机器上的 daemon 进程)。当服务以单副本运行在 Kubernetes 时:
- 滚动发布、Pod 重启、节点迁移都会造成连接中断窗口;
- 单 Pod 成为资源瓶颈,websocket 长连接与 Redis stream 消费无法水平扩展;
- daemon 重连、Pod 被杀等瞬态抖动会直接放大为用户可见的 RPC 失败。
多副本化的核心难题在于:一个用户的 daemon 可能连接在 Pod A,而他的 Web 客户端请求落在 Pod B,两者之间必须完成跨 Pod 的 RPC 转发与广播扇出。本文所有内容都围绕这条主线展开。
二、架构总览:Streams Adapter + 房间路由,零 Redis Key、零 TTL
docs/multi-process.md给出的核心设计可以浓缩为一句话:
happy-server 使用Socket.IO Redis streams adapter通过单一 Redis stream在副本之间转发
io.to(...).emit(...);RPC 路由(web → daemon)走名为rpc:<userId>:<method>的Socket.IO rooms;服务器通过io.in(room).fetchSockets()(cluster-adapter 提供的、可跨副本工作的原语)解析 daemon socket,并向单个 RemoteSocket 发送请求。没有 Redis key、没有 TTL、没有 Lua-CAS 清理、没有 keep-alive 刷新路径——成员关系就是标准 Socket.IO room 状态,断连时自动清理。
与上一版"Redis key + 60 秒 TTL + 心跳刷新"的方案相比,这套设计最大的变化是把"谁拥有某个 RPC 方法"这个事实交给 Socket.IO 自身的房间成员机制来管理,让框架替你处理成员生命周期,而不是自己维护一份极易腐化的外部状态。
2.1 三份关键代码的位置
. ├── packages/happy-server/sources/app/ │ ├── api/socket.ts io.Server 初始化;设置 REDIS_URL 时挂载 │ │ streams adapter;connectionStateRecovery │ │ 处于注释状态 │ ├── api/socket/rpcHandler.ts 整个 RPC 路由层(约 260 行,单一路径) │ ├── api/socket/machineUpdateHandler.ts 不再接触 RPC 状态 │ ├── api/socket/sessionUpdateHandler.ts 不再接触 RPC 状态 │ └── events/eventRouter.ts 通过房间做广播扇出 │ └── packages/happy-server/deploy/handy.yaml k8s Deployment + Service + Redis StatefulSet下面依次展开每条链路。
三、Socket 服务器初始化:Adapter 的挂载条件与连接鉴权
入口在 api/socket.ts,startSocket(app)创建io.Server时的关键配置:
| 配置项 | 值 | 说明 |
|---|---|---|
transports | ['websocket', 'polling'] | websocket 优先,polling 兜底(与生产客户端一致) |
pingTimeout/pingInterval | 45000/15000 | 心跳参数,RPC_RECONNECT_GRACE_MS的"2×心跳周期"以此推算 |
path | '/v1/updates' | 唯一的 Socket.IO 端点 |
serveClient | false | 不向客户端下发 Socket.IO 客户端文件 |
connectionStateRecovery | 注释掉 | 见"有意推迟"章节 |
Adapter 挂载的开关就是环境变量REDIS_URL:
// packages/happy-server/sources/app/api/socket.ts if (process.env.REDIS_URL) { const streamClient = new Redis(process.env.REDIS_URL); io.adapter(createAdapter(streamClient, { maxLen: 200000, readCount: 2000 })); log({ module: 'websocket' }, 'Redis streams adapter enabled for multi-process support'); ... }注意当前源码中的maxLen: 200000、readCount: 2000,与文档早期描述的~50000相比已上调——XADD 时由 Redis 自动修剪,无需人工清理(详见"Adapter 细节")。
3.1 鉴权中间件:为什么必须放在io.use
socket.ts还演示了一个容易踩坑的关键细节:鉴权必须在中间件(io.use)中完成,而不是在connect事件里异步 verifyToken。否则会存在一个窗口:客户端的rpc-register、rpc-call事件在 handler 挂载之前到达,被静默丢弃。中间件校验handshake.auth.token,并据此把userId、clientType(session-scoped/user-scoped/machine-scoped)、sessionId、machineId写入socket.data;session-scoped缺sessionId、machine-scoped缺machineId直接拒绝。
鉴权通过后,连接按作用域分三类(与 docs/realtime-sync-and-rpc.md 中的房间模型一一对应):
user-scoped:App/Web 客户端与账户级监听器;session-scoped:一个会话一个实时进程;machine-scoped:一台机器一个 daemon。
四、rpc-call 完整控制流:跨副本寻址 + 重连宽限 + 在飞探测
RPC 路由层全部集中在 rpcHandler.ts,docs/multi-process.md给出了这条单一路径的控制流:
rpc-call from web client . ├── input validation │ └── method name → invalid → callback({ok:false, error:'Invalid parameters'}) │ ├── 1. resolve target via cluster adapter │ └── fetchRoomSockets(io, 'rpc:<userId>:<method>') │ ├── io.in(room).timeout(...).fetchSockets() │ ├── on success → returns [...] │ └── on failure (peer replica unresponsive, fast adapter timeout) │ └── log + return [] (treat as "nobody here") │ │ │ ├── returns [target] → go to step 2 │ └── returns [] → go to wait-for-reconnect │ ├── wait-for-reconnect grace (only when no target found) │ └── waitForRoomMember(io, room, RPC_RECONNECT_GRACE_MS) │ └── poll every RPC_RECONNECT_POLL_MS via fetchRoomSockets: │ ├── room gained a member → return [target] │ └── deadline reached → return [] │ │ │ ├── grace produced [target] → go to step 2 │ └── grace produced [] │ └── callback({ok:false, error:'RPC method not available'}) │ ├── 2. sanity checks on resolved target │ ├── multiple sockets in room → log warn, use first │ └── target.id === socket.id → callback({ok:false, error:'same socket'}) │ ├── 3. fire emit + race a presence poll │ ├── ackPromise = target.timeout(RPC_CALL_TIMEOUT_MS).emitWithAck('rpc-request', ...) │ │ (cluster adapter routes cross-replica via Redis stream) │ │ │ └── presencePoll = while (alive) │ └── sleep RPC_PRESENCE_POLL_MS, fetchRoomSockets again │ ├── target still in room → keep watching │ └── target absent → throw 'RPC target disconnected' │ ├── Promise.race(ackPromise, presencePoll) │ ├── ackPromise resolves → callback({ok:true, result}) │ ├── ackPromise throws (timeout / err) → callback({ok:false, error: msg}) │ └── presencePoll throws → callback({ok:false, error:'RPC target disconnected'}) │ └── finally └── presenceAlive = false (stops the poll cleanly on success or failure)4.1 关键原语:带超时的fetchRoomSockets
源码把io.in(room).fetchSockets()封装成带调用方超时与失败兜底的函数:
// packages/happy-server/sources/app/api/socket/rpcHandler.ts async function fetchRoomSockets(io: Server, room: string, timeoutMs: number, context: 'lookup' | 'presence' = 'lookup'): Promise<RoomSockets> { try { return await io.in(room) .timeout(timeoutMs) .fetchSockets(); } catch (error) { rpcFetchSocketsTimeouts.inc({ context }); log({ module: 'websocket' }, `fetchSockets failed for ${room} (timeout=${timeoutMs}ms): ${error}`); return []; // 视为"没人在这",走重连宽限 } }两点设计意图值得注意:
- 失败降级为
[]而不是抛错:对端副本无响应时,调用方会走 wait-for-reconnect 宽限窗口,最终返回"RPC method not available",而不是让整个调用链崩掉; - lookup 与 presence 使用不同超时:in-flight 探测必须远小于 30 秒调用上限,否则死掉的副本会把每次轮询都拖满。
4.2 重连宽限与指数退避(当前源码比文档更进一步)
文档描述宽限窗口为 10 秒、轮询 200ms。当前main上的源码已经演进为指数退避 + 更大的窗口:
// packages/happy-server/sources/app/api/socket/rpcHandler.ts const RPC_CALL_TIMEOUT_MS = 30_000; const RPC_PRESENCE_POLL_MS = 2_000; // 重连宽限期内 fetchSockets 的超时,指数退避:2s → 4s → 8s const RPC_LOOKUP_FETCH_TIMEOUTS_MS = [2_000, 4_000, 8_000]; // in-flight 探测超时:必须 << RPC_CALL_TIMEOUT_MS,500ms 让死亡检测保持 ~2 次轮询内响应 const RPC_PRESENCE_FETCH_TIMEOUT_MS = 500; const RPC_RECONNECT_GRACE_MS = 15_000; const RPC_RECONNECT_POLL_MS = 200;waitForRoomMember的循环逻辑是:每次轮询的 fetch 超时按2s → 4s → 8s递增(RPC_LOOKUP_FETCH_TIMEOUTS_MS[Math.min(polls, 2)]),配合 200ms 睡眠,单次迭代耗时约 2.2s / 4.2s / 8.2s,15 秒窗口内可完成约 3 次尝试。退避的动机是减少慢速 Redis 下的流压力:超时 → 重试 → 超时的螺旋会被大量并发请求放大,退避让早期快速失败、后期给足时间,同时rpcLookupRetries直方图(桶[0..7])记录"第几次轮询才找到 daemon",便于观测宽限窗口的真实占用。
4.3 在飞探测:两次连续空轮询才判定断连
这是修复 Bug #1(in-flight RPC 吃满 30 秒超时)的核心机制,当前实现比文档描述更稳健——要求连续 2 次空轮询才宣告目标失联,避免瞬时 Redis/adapter 超时造成误杀:
const presencePoll = (async () => { let consecutiveMisses = 0; while (presenceAlive) { await sleep(RPC_PRESENCE_POLL_MS); if (!presenceAlive) return; const stillThere = await fetchRoomSockets(io, room, RPC_PRESENCE_FETCH_TIMEOUT_MS, 'presence'); if (!stillThere.some(s => s.id === target.id)) { consecutiveMisses++; if (consecutiveMisses >= 2) throw new Error('RPC target disconnected'); } else { consecutiveMisses = 0; } } })();为什么必须有这个探测?注释给出了根因:emitWithAck无法感知目标 socket 已死——daemon 的 Pod 在调用中途被杀时,cluster adapter 发出的 BROADCAST 请求在队列里等待永远不会到来的 BROADCAST_ACK,只能等到用户设置的 30 秒超时。adapter 的心跳检测约 10 秒才能发现 pod 失联,且不会主动取消挂起的广播。轮询fetchSockets是唯一能在 ~2-4 秒内判定"目标 socket 已消失"并快速中止的方法。
Promise.race([ackPromise, presencePoll])在finally中把presenceAlive置 false,保证无论成功还是失败轮询都会干净停止。三种终态:
- ack 先返回 →
{ok: true, result}; - ack 超时/报错 →
{ok: false, error: msg}; - presence 探测抛错 →
{ok: false, error: 'RPC target disconnected'}。
所有结果同时写入 Prometheus 指标rpc_calls_total(按 method/result 计数)与rpc_call_duration_seconds(直方图,桶[0.05..30]),method 会先经baseMethodName剥掉machineId/sessionId前缀(wire 格式如"cm9xyz123:bash"→"bash"),保证指标不因机器 ID 爆炸。
4.4 同一个房间出现多个 socket 怎么办
一个机器只有一个 daemon、每个方法只注册一次,正常不会出现。万一出现(例如异常状态下的重复注册),源码的处理是打 warn 日志并取targets[0],与上一版 Redis last-write-wins 的爆炸半径相同。另外,如果解析出的目标恰好是调用者自己(target.id === socket.id),直接返回'Cannot call RPC on the same socket',防止自调用死循环。
五、Daemon 生命周期:注册、应答、断连、重连
docs/multi-process.md用一张图总结了 daemon 侧的全部职责:
daemon (machine-scoped or session-scoped) . ├── connect to handy-server │ └── server: socket.handshake.auth.token → auth.verifyToken │ └── attaches rpcHandler / *UpdateHandler / etc │ ├── emit('rpc-register', { method }) │ └── server: socket.join('rpc:<userId>:<method>') │ └── ack: emit('rpc-registered', { method }) │ (Socket.IO room state, NO Redis key, NO TTL) │ ├── on('rpc-request', (data, cb) => …) │ └── handler runs, cb(result) returns the value via the cluster adapter │ ├── disconnect (any reason) │ └── Socket.IO automatically removes the socket from all rooms │ (cluster adapter syncs via heartbeat; no manual cleanup needed) │ └── auto-reconnect └── on 'connect': re-emit rpc-register (the only client-side responsibility)对应到源码,rpcHandler注册了三个事件:
rpc-register:校验method是字符串后socket.join(rpcRoom(userId, method)),回rpc-registered。房间名 =rpc:<userId>:<method>(RPC_ROOM_PREFIX = 'rpc:')。rpc-unregister:socket.leave(...),回rpc-unregistered。rpc-call:上文第四节描述的完整路由。- 没有 disconnect handler——注释明确写了:"Socket.IO removes the socket from all rooms automatically, and the cluster adapter syncs the removal to other replicas." 这正是"无 TTL、无手动清理"设计的体现:房间成员就是权威事实。
daemon 侧的客户端职责只有一条:每次重连成功后在connect事件里重新 emitrpc-register。旧方案中"daemon 重连后忘记重新注册"正是 Bug #2 与 #3 的一部分成因。
六、广播扇出:eventRouter 与四类房间
广播路径由 eventRouter.ts 承担:
eventRouter.emitUpdate / emitEphemeral . └── io.to(rooms).emit('update' | 'ephemeral', payload) ├── streams adapter: XADD on the 'socket.io' Redis stream │ (MAXLEN 自动修剪) └── every replica's XREAD loop picks up the entry └── delivers to its local sockets that match the room set (sockets that disconnected before the emit miss it; client falls through to apiSocket onReconnected → REST refetch)EventRouter在socket.ts中通过eventRouter.init(io)初始化,连接建立/断开时调用addConnection/removeConnection。addConnection把 socket 加入房间:
// packages/happy-server/sources/app/events/eventRouter.ts socket.join(`user:${userId}`); // 所有该用户的 socket switch (connection.connectionType) { case 'user-scoped': socket.join(`user:${userId}:user-scoped`); break; case 'session-scoped':socket.join(`user:${userId}:session:${connection.sessionId}`); break; case 'machine-scoped':socket.join(`user:${userId}:machine:${connection.machineId}`); break; }removeConnection是空实现——Socket.IO 断连自动清房。eventRouter使用的房间全集:
. ├── user:<userId> all of a user's sockets ├── user:<userId>:user-scoped only the web/desktop clients ├── user:<userId>:session:<sessionId> session-scoped subscribers └── user:<userId>:machine:<machineId> one specific machineRecipientFilter决定发往哪些房间(getRoomsForFilter):
| Filter | 目标房间 | 用途 |
|---|---|---|
all-user-authenticated-connections | user:<userId> | 默认:该用户全部连接 |
user-scoped-only | user:<userId>:user-scoped | 如 daemon 上下线状态推送 |
all-interested-in-session | user:<userId>:session:<sid>+user:<userId>:user-scoped | 会话订阅者(Socket.IO 自动去重) |
machine-scoped-only | user:<userId>:machine:<mid>+user:<userId>:user-scoped | 特定机器 + Web 端 |
emit还支持skipSenderConnection,此时改用socket.broadcast.to(rooms)排除发送者自身。
6.1 一个值得注意的跨副本用例:hasActiveUiClient
eventRouter.hasActiveUiClient(userId)用fetchSockets()判断用户当前是否在看某个 Happy UI 客户端(用于抑制冗余推送)。它对user:<userId>房间做带 2 秒超时的fetchSockets(),然后检查clientType === 'user-scoped' && appState === 'active'。注释里强调了两条刻意为之的语义:
- 只有
user-scoped是通知界面;session-scoped是编码 agent 本身、machine-scoped是 daemon,都不展示内容——若把 session-scoped 算进去,运行中的会话自己的 socket 会抑制它正请求的那条推送; - 只有显式上报过
app-state: active才算数,未上报视为未知而非在场("presence must be proven, not assumed")。
这也印证了fetchSockets()作为 cluster-adapter 原语在 RPC 之外的第二个跨副本用途。
七、前车之鉴:四个历史 Bug 与根因复盘
docs/multi-process.md明确指出:上一版方案把 RPC 路由状态存成rpc:user:<u>:method:<m>→ socketId 的Redis key,60 秒 TTL,由machine-alive/session-alive心跳续期。这套设计有四个致命问题,完整复盘与复现命令见 deploy/integration-tests/POSTMORTEM.md:
├── #1 In-flight RPC eats the full 30s timeout when the target pod dies │ io.to(deadSocketId).emitWithAck() has no fast-fail. │ FIX: presence poll aborts within ~1s(当前源码为 2 次空轮询、~2-4s) │ ├── #2 Reconnect race │ Between the daemon's disconnect cleanup and re-register, ~5–7% of │ cross-pod RPCs fail with either "method not available" (key │ deleted) or "target not reachable" (key still pointed at dead │ socketId). │ FIX: atomic socket.join / auto-leave on disconnect, no race window │ ├── #3 Silent TTL expiry(smoking gun) │ Daemon stays connected, registration vanishes after 60s if the │ keep-alive event was missed for any reason. Daemon never knows; │ stays broken until reconnect. │ FIX: no TTL exists anymore │ └── #4 Streams adapter "unbounded growth" FALSE ALARM. The adapter trims with MAXLEN on every XADD. Crossing this off the list.7.1 Bug #1:in-flight RPC 吃满 30 秒(复现时间线)
POSTMORTEM 中hammer.mjs pod-kill-mid-rpc的实测:
[+ 1.85s] firing rpc-call (will block 5s in handler) [+ 1.89s] daemon got rpc-request, sleeping 5s [+ 2.85s] killing daemon pod handy-server-67b86c7b7c-2bc6f [+ 2.94s] socket disconnect: transport close [+ 3.27s] socket reconnected [+31.85s] rpc-call result: ok=false latency=30002ms err=operation has timed outdaemon 的 Pod 被杀,socket 90ms 内断开、0.4s 后在另一 Pod 重连,但调用方整整挂了 30 秒。根因是emitWithAck对"目标 socketId 在整个集群已不存在"没有快速失败路径,且被 SIGKILL 的 Pod 连 Redis key 的清理 handler 都不会执行——死 key 无人回收。生产表现:每次 Pod 回收 + daemon 重连 → 每个并发 Web RPC 白等 30 秒,多次客户端重试叠加成数分钟级卡顿(即用户报告"运行 ls 三次花了三分钟")。
7.2 Bug #2:重连风暴竞态,~6% RPC 失败
results: success=178 fail=12 err: RPC method not available ×7 ← Redis key 已删、尚未重建 err: RPC target not reachable ×5 ← key 还指着已死的 socketId根因是 disconnect handler 的 Lua CAS 删 key 与 reconnect 的 SET 之间没有原子语义,设计上就无法在 daemon socket 过渡期间维持 RPC 可用。修复后 join/leave 由 Socket.IO 原子完成,竞态窗口消失。
7.3 Bug #3(决定性证据):静默 TTL 过期
[+ 55.35s] t=+55s rpc: ok=true [+ 65.35s] t=+65s rpc: ok=false err=RPC method not available [+ 75.36s] t=+75s rpc: ok=false err=RPC method not availabledaemon全程保持连接,但 60 秒一到注册凭空消失(TTL 刷新只发生在machine-alive/session-alive里),daemon 对此一无所知,也没有任何路径会重新注册——一旦发生就持续坏到下次重连。这解释了"UI 显示 daemon 在线但 RPC 报 method not available"的诡异现象。修复:TTL 根本不存在了,房间成员由断连自动清理,这条故障类别整体被消除。
7.4 Bug #4:虚惊一场
RedisXINFO STREAM socket.io显示流持续增长,但 adapter 每次 XADD 都按MAXLEN自动修剪(当前配置 200000),有界增长,划掉此项。POSTMORTEM 顺带指出两个观察:流上groups: 0(adapter 不用 consumer group,各副本用内存游标,Pod 重启后从$最新位置续读,重启窗口内写入的条目会跨副本丢失);这也是connectionStateRecovery与 REST 重取之所以重要的背景。
7.5 POSTMORTEM 的教训
复盘特别强调了一个方法论教训:为什么原始提交"测过 10/10 通过"却依然全是坑——因为测试只覆盖了稳态(连接就绪、无 churn),而四个 bug 全部活在转换态(Pod 被杀、重连、TTL 轮转)。这正是 deploy/integration-tests/ 里大量破坏性测试(rolling-deploy、dead-daemon)存在的意义。
八、部署形态:handy.yaml 与 Redis 基础设施
handy.yaml 是完整的 k8s 部署清单,当前内容比文档早期描述的replicas: 1已进一步演进(当前为replicas: 3),并且配置了发布策略与高可用约束:
apiVersion: apps/v1 kind: Deployment metadata: name: handy-server spec: strategy: type: RollingUpdate rollingUpdate: maxUnavailable: 0 # 滚动期间不允许低于期望副本数 maxSurge: 2 replicas: 3 selector: matchLabels: { app: handy-server } template: spec: containers: - name: handy image: docker.korshakov.com/handy-server:{version} ports: [{ containerPort: 3005 }] env: - name: NODE_ENV value: production - name: PORT value: "3005" - name: REDIS_URL # ← 多副本开关:设置后 socket.ts 挂载 streams adapter value: redis://happy-redis:6379 envFrom: - secretRef: { name: handy-secrets } livenessProbe: { httpGet: { path: /health, port: 3005 }, initialDelaySeconds: 30, periodSeconds: 60, timeoutSeconds: 30, failureThreshold: 10 } readinessProbe: { httpGet: { path: /health, port: 3005 }, initialDelaySeconds: 30, periodSeconds: 60, timeoutSeconds: 30, failureThreshold: 10 }同文件还包含:
- PodDisruptionBudget(
minAvailable: 1):保证主动驱逐(节点维护、滚动发布)时始终至少 1 个副本可用; - Service:
handy-server暴露3000 → 3005,默认ClusterIP;测试时可 patch 为LoadBalancer配合minikube tunnel获得真实 LB 行为(见下节); - Redis StatefulSet(
happy-redis,单副本、redis:7-alpine、appendonly + 持久卷):streams adapter 的依赖,多副本横向扩展后它仍是单一事实源; - ExternalSecret:从 Vault 拉取
/handy-db、/handy-master等密钥。
需要强调:REDIS_URL是多副本与单副本的行为分界——不设置时io.adapter不被调用,Socket.IO 退回单进程内存广播,所有房间成员与 RPC 只在本地可见。因此生产开启多副本 = 部署REDIS_URL+replicas > 1两个条件同时满足,缺一不可。
九、测试与验证:minikube 上的完整复现方案
docs/multi-process.md记录了本地 minikube(2 副本 handy-server + Redis + Postgres,minikube tunnel暴露真实LoadBalancer)上的全部测试矩阵,harness 均位于 deploy/integration-tests/。
9.1 从零搭建测试环境
packages/happy-server/deploy/integration-tests/local.sh # 启动 minikube、构建镜像、跑 prisma migrate、部署 kubectl get pods -l app=handy-server # 确认副本数 kubectl patch svc handy-server -p '{"spec":{"type":"LoadBalancer"}}' minikube tunnel & # 暴露 :3000 node packages/happy-server/deploy/integration-tests/test-rpc-cross-replica.mjslocal.sh的完整流程是:检查/启动 minikube →eval $(minikube docker-env)→ 用仓库根目录的Dockerfile.server构建happy-server:local→kubectl kustomize overlays/local部署 → 用一次性 Job 跑 Prisma migrate →rollout restart deployment/handy-server。注意local.sh内还顺带演示了两种访问方式:kubectl port-forward svc/handy-server 3005:3000或kubectl logs -l app=handy-server --all-containers -f双副本日志。
9.2 一键集成测试:run-all.sh
仓库已经提供了 10 项测试的一键运行器 run-all.sh,用法:
./run-all.sh # 假设集群已部署,跑全部测试 ./run-all.sh --deploy # 先构建部署再测试 ./run-all.sh --safe-only # 跳过杀 Pod 的破坏性测试10 项测试分为两组:
- 安全测试(8 项):
stress-prod-realistic(5000 条/秒事件)、stress-rpc-registration的 7 个场景——fire-and-forget、register-race-timing(注册竞态时序)、reconnect-no-ack(重连无 ack)、rapid-sessions(会话快速启停)、high-concurrency(50 个 daemon)、ios-session-flow、cross-replica-3pod(3 副本跨 Pod); - 破坏性测试(2 项):
rolling-deploy(滚动发布杀一个 Pod)、test-rpc-dead-daemon(杀 daemon Pod),跑完后wait_for_pods等副本恢复再继续。
脚本中还有一个值得学习的环境细节:访问服务器优先用minikube service handy-server --url(走 kube-proxy iptables 规则、可存活于 Pod 被杀),而port-forward是"单 Pod 隧道",杀 Pod 类测试会因此误报失败——所以setup_server_url会在 Service 是ClusterIP时打黄字警告。
9.3 最终验证矩阵(来自 multi-process.md)
修复后的全量测试结果:
├── steady-state cross-pod RPC 50/50 + 20/20 ✅ (after ~5s warmup) ├── pod-kill-mid-rpc 1612ms fast-fail ✅ (was 30000ms) ├── brief-disconnect SUCCESS in 2011ms ✅ ├── long-disconnect bounded 10542ms ✅ (10s grace + ~0.5s) ├── ttl-expiry (smoking gun) ALL 5 calls pass through +75s ✅ ├── reconnect-storm (5 cycles) 96–97% success ✅ (only inherent │ in-flight failures, ~3%) ├── broadcast multi-process 20/20 fan-out, 5/5 unaffected ✅ ├── network-loss 60s loop 85/85 zero failures ✅ └── missed-events parity event lost via socket, in DB, recovered=undefined ✅ (matches main)注意:文档记录的是当时 10s 宽限下的数字(pod-kill-mid-rpc1612ms、long-disconnect10542ms);当前源码宽限已上调至 15s 并引入指数退避,量级语义不变——快速失败从 30s 降到数秒、短暂断连由宽限吸收、TTL 类故障整体消失。ttl-expiry场景在修复后已无 TTL 可过期,"ALL 5 calls pass through +75s"正是该故障类别被根除的直接证据。
十、可调常量一览(当前 main 实际值)
docs/multi-process.md给出常量表,当前源码的实际取值如下(以源码为准,文档为早期快照):
| 常量 | 文档值 | 当前 main 值 | 作用 |
|---|---|---|---|
RPC_RECONNECT_GRACE_MS | 10_000 | 15_000 | 空房间时等待 daemon 重连的窗口;文档解释为"2×心跳周期"(心跳 15s×2),当前 15s + 指数退避覆盖约 3 次探测 |
RPC_RECONNECT_POLL_MS | 200 | 200 | 宽限窗口内的轮询节奏 |
RPC_PRESENCE_POLL_MS | 1_000 | 2_000 | in-flight 期间探测间隔(配合连续 2 次空轮询,实测快速失败 ~2-4s) |
RPC_PRESENCE_FETCH_TIMEOUT_MS | 500 | 500 | 单次跨副本 fetchSockets 上限,防止一个无响应副本拖慢每次轮询 |
RPC_CALL_TIMEOUT_MS | 30_000 | 30_000 | emitWithAck 上界,与 main 一致;两条路径都不支持 >30s 的 RPC |
| — | — | RPC_LOOKUP_FETCH_TIMEOUTS_MS | 新增:宽限探测的指数退避序列[2000, 4000, 8000] |
adaptermaxLen | ~50_000 | 200_000 | Redis stream 上限,每次 XADD 自动修剪 |
十一、Adapter 细节与已知边界
├── streams adapter discovery │ Pod 启动约 5s 后(adapter 默认 heartbeatInterval),跨副本 fetchSockets() │ 才能看到全部房间。新滚动发布后头几个 RPC 可能命中宽限窗口; │ 宽限被定为 2 个心跳周期即为此因。 │ ├── MAXLEN ~ 200000(当前源码) │ 在 socket.ts 中配置,每次 XADD 自动修剪,无需人工清理。 │ ├── fetchSockets() 跨副本 │ 默认每请求 5 秒超时;presence poll 显式传 timeout(500), │ 避免单个无响应副本把每次轮询拖 5 秒。 │ ├── emitWithAck from a RemoteSocket │ 跨副本可用——streams adapter 继承 ClusterAdapterWithHeartbeat, │ 实现了 BROADCAST_ACK 与 FETCH_SOCKETS_RESPONSE。 │ └── 同一 RPC 房间出现多个 socket │ 理论不会发生(一机一 daemon、一方法一注册)。发生了就 warn + 取 │ targets[0],爆炸半径与旧 Redis last-write-wins 一致。socket.ts还内置了一个流滞后观测:包装adapter.onRawMessage记录最后读取的 stream offset,每 5 秒用XINFO STREAM socket.io对比 stream HEAD,写入redisStreamLagMsGauge指标——多副本下 Redis stream 消费滞后是首要可观测性信号(docs/realtime-sync-and-rpc.md 的 Debugging 一节也把 stream lag 列为必查项)。
十二、有意推迟(不做的清单,含理由)
docs/multi-process.md明确列出的 deferred 项,每一条都有"为什么现在不做":
connectionStateRecovery(注释状态):socket.ts中注释掉的配置为maxDisconnectionDuration: 2 * 60 * 1000。streams adapter 支持它(已验证可用,missed-events.mjs证明强制engine.close()后重连可recovered=true重放事件),但当前选择先与多进程化之前的main行为对齐——客户端每次重连仍走完整 REST 重取(apiSocket.onReconnected)。开启它才能让短暂断连跳过重型 REST refetch。- in-flight RPC 跨 daemon 重连的连续性:与上一条耦合。若开启 recovery 且 presence poll 变成"等同一 socketId 回来 N 秒再失败",则 daemon 短暂网络抖动时:daemon 的 handler 继续跑、ack 包躺在客户端 sendBuffer、重连后冲刷出去、调用方拿到结果。今天 presence poll 只要房间空了就快速失败,正好杀掉这个场景。本 PR 不做。
- LB 层的用户亲和路由:经 streams adapter 的跨 Pod RPC 开销约 3–6ms,JWT 感知路由(Envoy / Istio / nginx-lua)属于比修复本身更大的基建改造,列入未来工作。
- UI "reconnecting…" 指示器:服务器现在会等 daemon 10–15 秒,但客户端 UI 还不展示这个等待,属于
apiSocket侧改动,与本文 PR 分离。 - 调小 adapter discovery 窗口:5s 是 streams adapter 默认 heartbeatInterval;调小能缩小新 Pod 启动竞态但增加 Redis 流量。
- 长运行 RPC(>30s):main 与本 PR 都不支持。CLI 的 bash 命令自带 30s 上限,与服务器 30s emit 超时"死磕平局";要放宽需同时调服务器与(可能新增的)客户端超时。
十三、总结与阅读路径
多副本 happy-server 的设计核心可以概括为三个词:房间即事实、无 TTL、快速失败。把 RPC 路由身份从"易腐化的外部 Redis key"迁移到"Socket.IO 房间成员"后,四个历史 Bug 中的三个被结构性消除(TTL 静默过期、重连竞态、死 socketId 悬挂),第四个(in-flight 吃满超时)由 presence poll 兜底快速失败;广播扇出则完全交给 streams adapter 的 XADD/XREAD 机制,辅以 MAXLEN 自动修剪与流滞后指标。
建议的深入阅读顺序:
- 控制流总览:docs/realtime-sync-and-rpc.md
- 本文主体:docs/multi-process.md
- 服务器子系统定位:docs/backend-architecture.md、docs/cli-architecture.md
- 线协议与事件名:docs/protocol.md
- 核心实现:socket.ts、rpcHandler.ts、eventRouter.ts
- 部署清单:handy.yaml
- 完整 Bug 复盘:POSTMORTEM.md
- 一键测试:run-all.sh、local.sh
如果你要复现或验证这套行为,最直接的路径是run-all.sh --deploy:它会自动完成 minikube 集群拉起、镜像构建、数据库迁移、双/三副本部署和全部 10 项测试,其中的破坏性测试会真实地杀掉 Pod,让你亲眼看到"30 秒超时 → 数秒快速失败"的差异。
【免费下载链接】happyMobile and Web client for Codex and Claude Code, with realtime voice, encryption and fully featured项目地址: https://gitcode.com/gh_mirrors/happy20/happy
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考