fix(WebSocket): 断线不重连、评分结果会丢、登出后连接还活着
Some checks failed
Deploy / deploy (push) Has been cancelled

排查 WS 这一块时发现的一批问题,多数是「机制写了但从没生效过」。

## 重连

`disconnect()` 里 `enableAutoReconnect = false`,而 `connect()` 从不改回 true ——
登出再登录后,这条连接就永远失去了自动重连能力(configUpdate 那条 watch 上尤其
明显)。改成用 `closedByUser` 表达「用户主动断开」的意图,和 `enableAutoReconnect`
这个**配置**分开。

`scheduleDisconnect` 的回调里断完紧接着一句 `enableAutoReconnect = true`,而
close 是异步的 —— 等 onclose 跑到时标志已经翻回来了,1 秒后又自动连上。那个
「15 分钟空闲省资源」从来没真正断开过。现在只断开,不做事后翻转。

重连的 setTimeout 没存句柄,组件卸载后照样触发 `connect()`,在已销毁的组件上
又建一条连接。现在 `disconnect()` 里 clearTimeout。

退避从「线性 ×5 次」改成「指数 + 抖动、30 秒封顶、次数不封顶」。原来 1+2+3+4+5
只有 15 秒,后端 deploy 重启一次就超了,之后这条连接死到用户刷新为止。抖动是
为了避免一个班几十台机器在同一毫秒一起冲回刚起来的后端。另挂 online /
visibilitychange,网络恢复或切回标签页立刻重连,不必等退避走完。

所有 socket 回调改成闭包住局部 ws 并在入口 `if (ws !== this.ws) return`,旧连接
迟到的 onclose 不再污染新连接的状态 —— 也是让 `disconnect()` 能被 onclose 识别
出来的关键。

## 订阅重放

`pendingSubmissionId` 一发送成功就清空,它只解决了「还没连上就 subscribe」,
**没解决断线重连**。而真正会丢结果的恰恰是后者:服务端收到 subscribe 会回一份
当前状态,掉线期间错过的推送就是靠这次重放补回来的;不重新订阅,重连后只收得到
「将来」的事件,可结果已经是过去式了。改成订阅意图保留到显式 `unsubscribe()`,
每次 onConnected 都重发。

这套逻辑原来只有 SubmissionWebSocket 有,FlowchartWebSocket 是 send 失败打一行
日志了事 —— socket 一掉,那次评分结果就再也回不来,按钮一直转圈。提成公共基类
SubscribingWebSocket,两条通道共用。

流程图另加轮询兜底:提交后 5 秒 WS 还没出结果就每 3 秒拉一次,读 status 2/3
结算,3 分钟上限。判题那边一直有兜底,流程图这边没有,而 Redis pub/sub 是发完
不管的,worker 推的那一刻连接不在就永远丢了。

`useSubmissionMonitor` 里 `watch(wsStatus, ..., { immediate: true })` 的回调在
watch() **返回之前**就同步跑了,已经连着时 `unwatch` 还是 null,if 不成立,
watcher 永远停不掉:每提交一次泄漏一个,往后每次重连它们都会把各自那个早就判完
的旧 submissionId 重新订阅一遍。整块删掉,直接 subscribe —— 基类已经管了时序。

WS handler 原来不校验 submissionId。学生同时开着几道题的页面时,每条连接都订在
同一个用户 topic 上,别的页面的评分结果会被当成自己的。

## 会话

握手时校验过一次会话就再也不管了,这条连接却能挂几个小时:用户在别的标签页
登出、或者会话本身到期,旧 socket 照样收推送。加一条 60 秒一轮的巡检,用
Redis EXPIRE 一条命令同时完成「判断存在」和「续期」(续期是必要的:只开着页面
挂 WS 的人一次 HTTP 请求都不发,不该被算成不活跃踢下线)。Redis 抛错时整轮
放弃,绝不因为一次抖动把全班踢下线。

禁用只改数据库的 isDisabled 列、不动 Redis 里的会话,巡检永远发现不了。加
`session:revoked` 频道主动通知,**两种作用域不能混**:

    { token }   用户登出。只断这一张会话 —— 同一个人在别的设备上是另一张
                会话,按 userId 广播会把他手机上的登录一起踢掉
    { userId }  账号被禁用。所有设备都得断

先发一帧 force_logout 再隔 100ms 断开。只断不发的话前端只看到一次普通掉线,
会照常重连、页面上还显示着登录态。token 不进帧里 —— 那是 httpOnly cookie 的值,
推到 WS 上就等于交给了 JS,匹配全在服务端做。

前端在协议层拦截 force_logout(和 pong 一样,不下发给业务 handler),并主动
disconnect —— 否则会一路 401 重连到退避上限,正是这机制要消掉的浪费。表现刻意
和 utils/api.ts 里 account-disabled / login-required 两支保持一致:同一件事从
HTTP 和 WS 两条路进来,学生看到的不该有两个样子。

## 开销与安全

`void handleMessage(...)` 是裸的,里面有两次 DB 查询和一个会抛的 schema.parse,
库抖一下就是一个 unhandled rejection(隔壁 bridgeSubmissionEvents 两处都接住了,
只有这里漏了)。

ping 提到用户查询之前。原来的顺序是「先查 user 再看消息类型」,每个客户端每
30 秒都要为一次心跳打一趟数据库。禁用用户不会因此漏网:推送路径上 bridge 会查,
subscribe 这条真正读数据的路径下面照样查。

bridge 两条 per-user 通道都先看 `server.subscriberCount(topic)`,没人订阅就别
查库了 —— 判题高峰期绝大多数事件的目标用户此刻并不在线。

flowchart 评分失败原来把 error.message 原样推给学生、前端直接弹出来,AI provider
的地址和内部报错就这么进了浏览器。改成真实原因写服务端日志。

加每连接令牌桶(20 突发 + 每秒回填 2)。一条 subscribe 在服务端是一到两次数据库
查询,一个学生开着一条 socket 狂发就能压住库。正常流量离阈值几十倍远。

升级时校验 Origin。会话 cookie 是 SameSite=Lax、WS 握手不是导航,跨站页面本来就
带不上 cookie,所以这是防御纵深不是唯一防线。同源放行;本机开发(Vite 5173 →
API 3000)自动放行,且只在两边都是本机时成立 —— 生产环境 url.hostname 是正式
域名,这条永远不触发;跨域部署走 ALLOWED_WS_ORIGINS。不发 Origin 的一律放行:
真正的攻击面是带着受害者 cookie 的浏览器页面,而浏览器一定会带 Origin。

顺带:useConfigWebSocket 的 handler 从 onMounted 挪到同步注册(调用方在 setup
阶段就 connect() 了),删掉每条消息打完整内容的 console.log 和死字段
ws.data.username。

## 没动的

题目页上并没有两条 /ws/submissions —— Form.vue 里 SubmitFlowchart 和 SubmitCode
是 v-if/v-else,互斥。学生实际是 2 条连接:全站一条 /ws/config + 一条
/ws/submissions,正常,不必合并。

## 验证

55 个用例,分六组打桩跑(假 WebSocket + 假计时器;会话/限流/吊销三组对着真
Redis):

    重连语义        7   断开后不再自我复活、卸载后计时器已取消、退避封顶、
                        online 立即重连、旧 onclose 不污染新连接
    订阅重放        9   未就绪时补发、重连后重新订阅、unsubscribe 后不再重放
    会话巡检        8   EXPIRE 三态、只断失效的、同 token 只查一次、
                        Redis 抖动时一个都不踢
    Origin/限流    13   跨站与跨端口拒绝、生产不因 localhost 开后门、
                        突发额度、按时间回填
    强制登出       13   登出只断同 token(别的设备不受牵连)、禁用断所有设备、
                        巡检先通知再断
    前端登出        5   两支表现、收到后不再重连、两条通道只处理一次

前两组做了改动前/后对比,老代码该挂的都挂了 —— 「重连后自动重新订阅」正是这么
跑出来的,此前我以为 pendingSubmissionId 已经覆盖了这种情况。

apps/api 的 tsc 和 apps/web 的 vue-tsc 都干净,仓库既有测试照常通过。

**SubmitFlowchart.vue 的轮询兜底未经运行时验证** —— 在 SFC 内,没搭组件挂载
环境,只过了类型检查和人工核对。要验的话,停掉 worker 提交一次流程图,看 5 秒
后是否转入轮询、3 分钟后是否给出超时提示。

这批改动动了 WS 的行为面(限流会断连接、Origin 会拒绝、巡检会踢会话),上线前
建议手测:提交代码看判题、切流程图看评分、后台改配置看全站生效、开两个标签页
在一个里登出、禁用一个在线学生看另一端反应。

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
2026-08-27 05:07:26 -06:00
parent 47a43b9880
commit 64facc5701
12 changed files with 647 additions and 138 deletions

View File

@@ -1,6 +1,8 @@
import { flowchartUpdateSchema, submissionUpdateSchema } from "@oj2/contract"
import { and, eq } from "drizzle-orm"
import { touchSession } from "./auth/session"
import { config } from "./config"
import { db, schema } from "./db"
import {
parseSubmissionEvent,
@@ -9,13 +11,147 @@ import {
} from "./judge/events"
import { JudgeStatus } from "./judge/status"
import { createSubscriberRedis } from "./redis"
import { configTopic, configUpdateChannel, parseUserEvent, userEventChannel, userEventTopic } from "./events"
import {
configTopic,
configUpdateChannel,
parseSessionRevoked,
parseUserEvent,
sessionRevokedChannel,
userEventChannel,
userEventTopic,
} from "./events"
/** 本机的几种写法。开发时 Vite 代理会让 Origin5173和 Host3000对不上 */
const LOCAL_HOSTNAMES = new Set(["localhost", "127.0.0.1", "::1", "[::1]"])
/**
* WebSocket 升级的来源校验。
*
* 会话 cookie 是 SameSite=Lax而 WebSocket 握手不是导航,跨站页面本来就带不上
* 这个 cookie —— 所以这里是防御纵深,不是唯一防线。
*
* 不发 Origin 的一律放行:真正的攻击面是「带着受害者 cookie 的浏览器页面」,
* 而浏览器一定会带 Origin脚本客户端本来就能伪造任意请求头拦它没有意义。
*/
export function isAllowedWebSocketOrigin(origin: string | null, url: URL) {
if (!origin) return true
if (config.allowedWebSocketOrigins.includes(origin)) return true
let originUrl: URL
try {
originUrl = new URL(origin)
} catch {
return false
}
if (originUrl.host === url.host) return true
// 两边都是本机才放行。生产环境 url.hostname 是正式域名,这条永远不成立
return (
LOCAL_HOSTNAMES.has(originUrl.hostname) && LOCAL_HOSTNAMES.has(url.hostname)
)
}
export interface SubmissionSocketData {
userId: number
username: string
/** 同一个 Bun.serve 只能挂一个 websocket handler用它区分两条通道 */
kind: "submissions" | "config"
/** 握手时那张会话的 token留着定期确认它还没被登出 / 过期,见 sweepSessions */
token: string
/** 令牌桶open 时初始化,见 allowMessage */
rate?: { tokens: number; updatedAt: number }
}
/**
* 每条连接的消息限流。
*
* 一条 subscribe 在服务端是一到两次数据库查询,一个学生开着一条 socket 狂发就能
* 压住库。正常流量离这个阈值很远:心跳 30 秒一条,订阅一次提交也就一两条,
* 20 的突发额度 + 每秒 2 个的回填是几十倍的余量。
*/
const RATE_BURST = 20
const RATE_REFILL_PER_SECOND = 2
function allowMessage(ws: Bun.ServerWebSocket<SubmissionSocketData>) {
const now = Date.now()
const rate = (ws.data.rate ??= { tokens: RATE_BURST, updatedAt: now })
const refill = ((now - rate.updatedAt) / 1000) * RATE_REFILL_PER_SECOND
rate.tokens = Math.min(RATE_BURST, rate.tokens + refill)
rate.updatedAt = now
if (rate.tokens < 1) return false
rate.tokens -= 1
return true
}
/**
* 当前挂着的连接。Bun 不提供遍历连接的接口,要定期巡检就得自己登记。
* open 时加入、close 时移除,见 sweepSessions。
*/
const liveSockets = new Set<Bun.ServerWebSocket<SubmissionSocketData>>()
/** 会话巡检间隔。够快到登出后一分钟内断开,又不至于让 Redis 忙起来 */
const SESSION_SWEEP_INTERVAL = 60_000
/** 先把 force_logout 帧发出去,再断连接,留一拍给它出门 */
const FORCE_LOGOUT_CLOSE_DELAY = 100
/**
* 通知并断开一批连接。
*
* 之所以先发一帧再断:只断连接的话前端只看到一次普通掉线,会照常重连,页面上
* 还显示着登录态;收到 force_logout 才知道要清掉身份、弹登录框或者提示被禁用。
*/
function forceLogout(
targets: Bun.ServerWebSocket<SubmissionSocketData>[],
reason: string,
) {
if (targets.length === 0) return
const frame = JSON.stringify({ type: "force_logout", reason })
for (const ws of targets) ws.send(frame)
setTimeout(() => {
for (const ws of targets) ws.close(1008, "Session ended")
}, FORCE_LOGOUT_CLOSE_DELAY)
}
/**
* 定期把会话已经失效的连接断掉。
*
* 握手时校验过一次会话,但这条连接能挂几个小时 —— 期间用户可能在别的标签页登出,
* 或者会话本身到期。只靠消息触发的校验不够:一条连接完全可能除了心跳什么都不发,
* 而心跳是**故意**不查会话的(否则每客户端每 30 秒一趟 Redis 又回来了)。
*
* 注意这里不查 isDisabled管理员禁用只改数据库列、不删会话所以 token 校验
* 覆盖不到它。禁用由推送路径上的 bridgeSubmissionEvents 挡着 —— 被禁用的学生
* 收不到任何数据socket 还挂着只是根空管子。
*/
export async function sweepSessions() {
// 一个学生至少有配置和提交两条通道,多开几个标签页还会更多,而它们共用同一张
// 会话 —— 一轮里同一个 token 只查一次
const checked = new Map<string, boolean>()
const dead: Bun.ServerWebSocket<SubmissionSocketData>[] = []
for (const ws of liveSockets) {
const token = ws.data.token
let alive = checked.get(token)
if (alive === undefined) {
try {
alive = await touchSession(token)
} catch (error) {
// Redis 抖一下不该把全班踢下线:这一轮直接放弃,下一轮再说
console.error("Failed to verify websocket sessions", error)
return
}
checked.set(token, alive)
}
if (!alive) dead.push(ws)
}
// 会话没了有两种可能:在别的标签页登出了,或者会话自己到期。对用户都是
// 「要重新登录」,走 session-ended 这一支
forceLogout(dead, "session-ended")
}
export function startSessionSweep() {
const timer = setInterval(() => {
void sweepSessions()
}, SESSION_SWEEP_INTERVAL)
timer.unref()
return timer
}
function objectValue(value: unknown): Record<string, unknown> {
@@ -27,6 +163,8 @@ function objectValue(value: unknown): Record<string, unknown> {
export function submissionWebSocketHandler(): Bun.WebSocketHandler<SubmissionSocketData> {
return {
open(ws) {
liveSockets.add(ws)
ws.data.rate = { tokens: RATE_BURST, updatedAt: Date.now() }
if (ws.data.kind === "config") {
ws.subscribe(configTopic)
return
@@ -35,9 +173,20 @@ export function submissionWebSocketHandler(): Bun.WebSocketHandler<SubmissionSoc
ws.subscribe(userEventTopic(ws.data.userId))
},
message(ws, message) {
void handleMessage(ws, String(message))
if (!allowMessage(ws)) {
ws.close(1008, "Too many messages")
return
}
// handleMessage 里有 DB 查询和会抛的 schema.parse。以前是裸的 `void`
// 库抖一下就是一个 unhandled rejection隔壁 bridgeSubmissionEvents 两处
// 都接住了,只有这里漏了)
handleMessage(ws, String(message)).catch((error) => {
console.error("Failed to handle websocket message", error)
ws.send(JSON.stringify({ type: "error", message: "Internal error" }))
})
},
close(ws) {
liveSockets.delete(ws)
if (ws.data.kind === "config") {
ws.unsubscribe(configTopic)
return
@@ -52,6 +201,34 @@ async function handleMessage(
ws: Bun.ServerWebSocket<SubmissionSocketData>,
raw: string,
) {
let message: { type?: unknown; timestamp?: unknown; submissionId?: unknown }
try {
message = JSON.parse(raw) as typeof message
} catch {
ws.send(JSON.stringify({ type: "error", message: "Invalid JSON" }))
return
}
// 心跳不查库。原来的顺序是「先查 user 再看消息类型」,于是每个客户端每 30 秒
// 都要为一次 ping 打一趟数据库;一个题目页还开着两条连接,全班在线时纯空转。
// 禁用用户不会因此漏网:往用户 topic 推之前 bridgeSubmissionEvents 会查一次,
// 而 subscribe 这条真正读数据的路径下面照样查。
if (message.type === "ping") {
ws.send(JSON.stringify({ type: "pong", timestamp: message.timestamp }))
return
}
if (message.type !== "subscribe" || typeof message.submissionId !== "string") {
ws.send(JSON.stringify({ type: "error", message: "Invalid message" }))
return
}
// 会话可能在连接期间就失效了:用户在别的标签页登出,或者会话自己到期。
// 握手时校验过一次不算数 —— 这条连接能挂几个小时。
if (!(await touchSession(ws.data.token))) {
ws.close(1008, "Session expired")
return
}
const [activeUser] = await db
.select({ id: schema.user.id })
.from(schema.user)
@@ -67,23 +244,6 @@ async function handleMessage(
return
}
let message: { type?: unknown; timestamp?: unknown; submissionId?: unknown }
try {
message = JSON.parse(raw) as typeof message
} catch {
ws.send(JSON.stringify({ type: "error", message: "Invalid JSON" }))
return
}
if (message.type === "ping") {
ws.send(JSON.stringify({ type: "pong", timestamp: message.timestamp }))
return
}
if (message.type !== "subscribe" || typeof message.submissionId !== "string") {
ws.send(JSON.stringify({ type: "error", message: "Invalid message" }))
return
}
const [submission] = await db
.select({
id: schema.submission.id,
@@ -112,7 +272,7 @@ async function handleMessage(
const replay = flowchart.status === 2
? { type: "flowchart_evaluation_completed", submissionId: flowchart.id, score: flowchart.score ?? undefined, grade: flowchart.grade ?? undefined }
: flowchart.status === 3
? { type: "flowchart_evaluation_failed", submissionId: flowchart.id, error: "Evaluation failed" }
? { type: "flowchart_evaluation_failed", submissionId: flowchart.id }
: { type: "flowchart_evaluation_update", submissionId: flowchart.id }
ws.send(JSON.stringify(flowchartUpdateSchema.parse(replay)))
return
@@ -147,9 +307,27 @@ export async function bridgeSubmissionEvents(
server.publish(configTopic, raw)
return
}
if (channel === sessionRevokedChannel) {
const revoked = parseSessionRevoked(raw)
if (!revoked) return
// 按 token 还是按 userId取决于是「这张会话登出了」还是「这个账号被禁用了」
forceLogout(
[...liveSockets].filter((ws) =>
revoked.token !== undefined
? ws.data.token === revoked.token
: ws.data.userId === revoked.userId,
),
revoked.reason,
)
return
}
if (channel === userEventChannel) {
const event = parseUserEvent(raw)
if (!event) return
const topic = userEventTopic(event.userId)
// 这台实例上没人订阅就到此为止:判题高峰期绝大多数事件的目标用户此刻并不
// 在线,查一次库只为了 publish 给零个订阅者
if (server.subscriberCount(topic) === 0) return
void (async () => {
const [activeUser] = await db
.select({ id: schema.user.id })
@@ -157,7 +335,7 @@ export async function bridgeSubmissionEvents(
.where(and(eq(schema.user.id, event.userId), eq(schema.user.isDisabled, false)))
.limit(1)
if (!activeUser) return
server.publish(userEventTopic(event.userId), JSON.stringify(event.data))
server.publish(topic, JSON.stringify(event.data))
})().catch((error) => {
console.error("Failed to bridge user event", error)
})
@@ -166,6 +344,8 @@ export async function bridgeSubmissionEvents(
if (channel !== submissionUpdateChannel) return
const event = parseSubmissionEvent(raw)
if (!event) return
const topic = userSubmissionTopic(event.userId)
if (server.subscriberCount(topic) === 0) return
void (async () => {
const [activeUser] = await db
.select({ id: schema.user.id })
@@ -178,10 +358,7 @@ export async function bridgeSubmissionEvents(
)
.limit(1)
if (!activeUser) return
server.publish(
userSubmissionTopic(event.userId),
JSON.stringify(event.data),
)
server.publish(topic, JSON.stringify(event.data))
})().catch((error) => {
console.error("Failed to bridge submission event", error)
})
@@ -189,6 +366,11 @@ export async function bridgeSubmissionEvents(
subscriber.on("error", (error) => {
console.error("Submission event subscriber error", error)
})
await subscriber.subscribe(submissionUpdateChannel, userEventChannel, configUpdateChannel)
await subscriber.subscribe(
submissionUpdateChannel,
userEventChannel,
configUpdateChannel,
sessionRevokedChannel,
)
return subscriber
}