覆盖 Phase 2e-2 代码审查报告(docs/reviews/2026-04-23-phase2e-2-code-review.md) P2/Nit 收尾批次,均在本仓库完成闭环;剩余 5 项登记推迟至 Phase 2f/3。 ===== P2 七项 ===== - P2-1 cleanupUserResources 补 transport 清理 · media-server 新增 DELETE /internal/v1/transports/:id + transport.service.closeTransport · Go MediaOrchestrator 接口新增 CloseTransport;HTTP 实现按 doCloseRequest 走 4xx 幂等 + 指数退避 · meeting_signal_service.cleanupUserResources 新增 "transport:<id>" 分支 - P2-2 preview.vue 快速切设备竞态 · previewSeq 序号 + nextTick 后 stale 判断,丢弃过期结果 · onVideoChange/onAudioChange 走 scheduleRestartPreview 200ms 防抖 · onBeforeUnmount 清理 changeDebounceTimer - P2-3 room.vue onLoad redirectTo 后补 return · 引入 redirectingToJoin 守卫,跳转页不再执行 onMounted 初始化 - P2-6 generateUniqueRoomCode 重试上限监控 · 重试后成功:Warn 日志(码空间健康度告警) · 重试耗尽:Error 日志 + ErrRoomCodeConflict · 修正 logs.Error 调用签名(去掉多余的 nil) - P2-7 SendChatMessage 服务端长度 + 频率限制 · 新增 ErrChatContentEmpty / ErrChatContentTooLong / ErrChatRateLimited · utf8.RuneCountInString 校验 500 字符上限 · Redis INCR + EXPIRE 滑动窗口(30 条/60s,首次写入 EXPIRE 兜底) · controller.handleError 映射为 HTTP 400 - P2-8 MEETING_ENDED_REASON_LABEL 覆盖复核 · 新增前端专属常量 MEETING_ENDED_REASON_KICKED + 文案 · store/meeting.js _onMemberKicked 使用常量 · 同步修复后端 OnWSDisconnect 硬编码 "ws_disconnect" → constants.MeetingLeftReasonDisconnect - P2 已修 P2-1/2/3/6/7/8;P2-4(WS token 迁出 URL query)与 P2-5(Chat 服务拆分)登记推迟 ===== Nit 七项 ===== - Nit: kind:id 解析改用 strings.SplitN · cleanupUserResources / pushExistingRoomState 两处同步 - Nit: resourceTTL 中央化 · 新增 constants.MeetingResourceTrackTTLSeconds(3600) · meeting_signal_service.go resourceTTL 由 const 改 var 并引用常量 - Nit: ws/handler.go CheckOrigin 白名单 · NewHandler 新增 serverCfg 依赖;checkOrigin 支持同源放行 / dev 模式放行 / release 模式白名单严格匹配 · config.ServerConfig 新增 WSAllowedOrigins(逗号分隔)+ AllowedOrigins() / IsRelease() 辅助方法 · provider.go 新增 provideServerConfig,wire_gen.go 同步 - Nit: http_media_orchestrator.go 超时 + CreateRouter 重试 · 默认 TimeoutMS 5000→10000ms 兼容 Worker 冷启动 · 新增 CreateRouterRetry(默认 1 次,300ms 退避),仅对非 404 错误重试 · config.dev.yaml / config.docker.yaml 同步写入显式配置 - Nit: deploy-public.sh REDIS_PASSWORD × redis.conf 联动校验 · 检测 REDIS_PASSWORD 与 redis.conf 的 requirepass 配对一致性 · redis.conf 增加公网部署 requirepass 使用说明 - Nit: media-server internal-auth isPrivatePath 按 path 匹配 · 剔除 query/hash 后再与白名单 startsWith,避免 "?" 语义混淆 - Nit: mediasoup-client.js in-flight 锁走读确认 · finally 分支已覆盖 resolve/reject 两路,追加注释强化语义 - Nit 走读复核:_onMemberLeft 整槽关闭 vs _onProducerNew(closed=true) 精确匹配 producerId · 粒度正确,无需改动(登记结论) ===== 构建验证 ===== - go vet ./... / go build ./... 通过 - frontend npm run build:h5 通过(仅 uni-app legacy warning,无 error) - media-server npx tsc --noEmit 通过 ===== 审查追踪小节 ===== docs/reviews/2026-04-23-phase2e-2-code-review.md 追加 "Task 16 修复追踪(2026-04-24 更新)": - 已修复一览(本批次 14 处 + 历次 commitcdaa39d/ea2bf96/f5ae095/ 5ed14c2) - 推迟登记表(P2-4 / P2-5 / 端口收敛 / appData 校验 / RFC3339 时间格式,共 5 项) Made-with: Cursor
262 lines
8.9 KiB
Go
262 lines
8.9 KiB
Go
// Package ws 提供 WebSocket 连接处理
|
||
// 负责 HTTP → WebSocket 升级、JWT 认证、消息路由分发
|
||
package ws
|
||
|
||
import (
|
||
"context"
|
||
"net/http"
|
||
"net/url"
|
||
"strings"
|
||
|
||
"github.com/echochat/backend/config"
|
||
"github.com/echochat/backend/pkg/logs"
|
||
"github.com/echochat/backend/pkg/utils"
|
||
"github.com/echochat/backend/pkg/ws"
|
||
"github.com/gin-gonic/gin"
|
||
"github.com/gorilla/websocket"
|
||
"go.uber.org/zap"
|
||
)
|
||
|
||
// TokenValidator 有状态 JWT 验证接口(检查 Token 是否在 Redis 中有效)
|
||
// 由 auth.AuthService 实现,用于防止已登出用户建立 WebSocket 连接
|
||
type TokenValidator interface {
|
||
ValidateAccessToken(ctx context.Context, userID int64, clientType, token string) bool
|
||
}
|
||
|
||
// OfflineMessagePusher 离线消息推送接口
|
||
// 由 im.handler.OfflinePusher 实现,WebSocket 连接建立后触发推送
|
||
type OfflineMessagePusher interface {
|
||
PushOfflineMessages(ctx context.Context, userID int64)
|
||
}
|
||
|
||
// NotifyConnectHook 通知未读补偿钩子接口
|
||
// 由 notify.service.NotifyService 隐式实现
|
||
// WebSocket 连接建立后触发,向客户端推送 notify.unread.total 事件
|
||
// 用于断线重连场景下的徽标状态同步
|
||
type NotifyConnectHook interface {
|
||
PushUnreadTotalOnConnect(ctx context.Context, userID int64)
|
||
}
|
||
|
||
// MeetingDisconnectHook 会议 WS 断线钩子接口(Phase 2e-2 Task 8)
|
||
// 由 meeting.service.MeetingSignalService 隐式实现
|
||
// 触发时机:用户最后一条 WS 连接被移除(最终下线)
|
||
// 职责:清理该用户在会议中的媒体资源;若为 host 则启动 host 宽限期
|
||
// 抽象为接口避免 ws 包反向依赖 meeting 包引发循环引用
|
||
type MeetingDisconnectHook interface {
|
||
OnWSDisconnect(ctx context.Context, userID int64)
|
||
}
|
||
|
||
// Handler WebSocket 连接处理器
|
||
type Handler struct {
|
||
hub *ws.Hub
|
||
pubsub *ws.PubSub
|
||
jwtCfg *config.JWTConfig
|
||
serverCfg *config.ServerConfig // Task 16 Nit:CheckOrigin 白名单需要
|
||
onlineService *OnlineService
|
||
tokenValidator TokenValidator
|
||
offlinePusher OfflineMessagePusher
|
||
notifyConnectHook NotifyConnectHook
|
||
meetingDisconnectHook MeetingDisconnectHook // Task 8 注入
|
||
upgrader websocket.Upgrader
|
||
}
|
||
|
||
// NewHandler 创建 WebSocket Handler 实例
|
||
// Task 16 Nit:新增 serverCfg 参数,按 server.ws_allowed_origins + server.mode 收敛 CheckOrigin
|
||
func NewHandler(hub *ws.Hub, pubsub *ws.PubSub, jwtCfg *config.JWTConfig, serverCfg *config.ServerConfig, onlineService *OnlineService, tokenValidator TokenValidator) *Handler {
|
||
h := &Handler{
|
||
hub: hub,
|
||
pubsub: pubsub,
|
||
jwtCfg: jwtCfg,
|
||
serverCfg: serverCfg,
|
||
onlineService: onlineService,
|
||
tokenValidator: tokenValidator,
|
||
}
|
||
h.upgrader = websocket.Upgrader{
|
||
ReadBufferSize: 1024,
|
||
WriteBufferSize: 1024,
|
||
CheckOrigin: h.checkOrigin,
|
||
}
|
||
return h
|
||
}
|
||
|
||
// checkOrigin 按配置收敛 WebSocket 握手 Origin:
|
||
// - 同源(Origin 为空 或 Origin.Host == Request.Host)→ 放行
|
||
// - dev 模式(server.mode != release)+ 未配置白名单 → 放行全部(便于本地开发调试)
|
||
// - release 模式 / 配置了白名单 → 仅放行白名单匹配的 Origin
|
||
//
|
||
// 精确匹配 scheme+host+port,不做前缀/通配,避免误匹配
|
||
func (h *Handler) checkOrigin(r *http.Request) bool {
|
||
origin := strings.TrimSpace(r.Header.Get("Origin"))
|
||
// 空 Origin:非浏览器场景(如 curl / server-to-server),默认放行
|
||
if origin == "" {
|
||
return true
|
||
}
|
||
u, err := url.Parse(origin)
|
||
if err != nil || u.Host == "" {
|
||
logs.Warn(r.Context(), "ws.handler.checkOrigin", "Origin 非法,拒绝",
|
||
zap.String("origin", origin))
|
||
return false
|
||
}
|
||
// 同源(Origin.Host == Request.Host)直接放行
|
||
if strings.EqualFold(u.Host, r.Host) {
|
||
return true
|
||
}
|
||
allowed := h.serverCfg.AllowedOrigins()
|
||
// 未配置白名单:dev 放行,release 拒绝
|
||
if len(allowed) == 0 {
|
||
if h.serverCfg.IsRelease() {
|
||
logs.Warn(r.Context(), "ws.handler.checkOrigin", "release 模式未配置 WSAllowedOrigins,拒绝跨源",
|
||
zap.String("origin", origin))
|
||
return false
|
||
}
|
||
return true
|
||
}
|
||
// 精确匹配白名单
|
||
for _, ao := range allowed {
|
||
if strings.EqualFold(ao, origin) {
|
||
return true
|
||
}
|
||
}
|
||
logs.Warn(r.Context(), "ws.handler.checkOrigin", "Origin 不在白名单,拒绝",
|
||
zap.String("origin", origin))
|
||
return false
|
||
}
|
||
|
||
// SetOfflinePusher 设置离线消息推送器(由 IM 模块在初始化时注入)
|
||
func (h *Handler) SetOfflinePusher(pusher OfflineMessagePusher) {
|
||
h.offlinePusher = pusher
|
||
}
|
||
|
||
// SetNotifyConnectHook 设置通知未读补偿钩子(由 notify 模块在初始化时注入)
|
||
func (h *Handler) SetNotifyConnectHook(hook NotifyConnectHook) {
|
||
h.notifyConnectHook = hook
|
||
}
|
||
|
||
// SetMeetingDisconnectHook 设置会议 WS 断线钩子(由 meeting 模块在初始化时注入)
|
||
// Phase 2e-2 Task 8:WS 最终下线时触发 host 宽限期 / 媒体资源清理
|
||
func (h *Handler) SetMeetingDisconnectHook(hook MeetingDisconnectHook) {
|
||
h.meetingDisconnectHook = hook
|
||
}
|
||
|
||
// Upgrade 处理 WebSocket 升级请求
|
||
// GET /ws?token=xxx → JWT 认证 → 升级连接 → 注册 Hub → 订阅 Redis 频道
|
||
func (h *Handler) Upgrade(c *gin.Context) {
|
||
funcName := "ws.handler.Upgrade"
|
||
|
||
token := c.Query("token")
|
||
if token == "" {
|
||
utils.ResponseUnauthorized(c, "缺少认证 Token")
|
||
return
|
||
}
|
||
|
||
claims, err := utils.ParseToken(h.jwtCfg, token)
|
||
if err != nil {
|
||
logs.Warn(nil, funcName, "WebSocket Token 验证失败", zap.Error(err))
|
||
utils.ResponseUnauthorized(c, "Token 无效或已过期")
|
||
return
|
||
}
|
||
|
||
clientType := claims.ClientType
|
||
if clientType == "" {
|
||
clientType = "frontend"
|
||
}
|
||
if h.tokenValidator != nil && !h.tokenValidator.ValidateAccessToken(c.Request.Context(), claims.UserID, clientType, token) {
|
||
logs.Warn(nil, funcName, "WebSocket Token 已失效(Redis 校验)",
|
||
zap.Int64("user_id", claims.UserID))
|
||
utils.ResponseUnauthorized(c, "认证已失效,请重新登录")
|
||
return
|
||
}
|
||
|
||
conn, err := h.upgrader.Upgrade(c.Writer, c.Request, nil)
|
||
if err != nil {
|
||
logs.Error(nil, funcName, "WebSocket 升级失败",
|
||
zap.Int64("user_id", claims.UserID), zap.Error(err))
|
||
return
|
||
}
|
||
|
||
client := ws.NewClient(h.hub, conn, claims.UserID)
|
||
client.SetOnDisconnect(func(userID int64) {
|
||
if client.IsClosedByHub() && h.hub.IsOnline(userID) {
|
||
logs.Info(nil, "ws.handler.onDisconnect", "连接被新连接替换,跳过下线清理",
|
||
zap.Int64("user_id", userID))
|
||
return
|
||
}
|
||
h.pubsub.Unsubscribe(userID)
|
||
h.onlineService.UserOffline(context.Background(), userID)
|
||
// Task 8:通知会议模块处理 host 宽限期 / 媒体资源清理
|
||
// 钩子为可选(测试或 meeting 模块未加载场景安全跳过)
|
||
if h.meetingDisconnectHook != nil {
|
||
go h.meetingDisconnectHook.OnWSDisconnect(context.Background(), userID)
|
||
}
|
||
})
|
||
h.hub.Register(client)
|
||
h.pubsub.Subscribe(claims.UserID)
|
||
h.onlineService.UserOnline(c.Request.Context(), claims.UserID, c.ClientIP())
|
||
|
||
logs.Info(nil, funcName, "WebSocket 连接建立",
|
||
zap.Int64("user_id", claims.UserID),
|
||
zap.String("ip", c.ClientIP()))
|
||
|
||
go client.WritePump()
|
||
go client.ReadPump(h.createReadHandler(claims.UserID))
|
||
|
||
if h.offlinePusher != nil {
|
||
go h.offlinePusher.PushOfflineMessages(context.Background(), claims.UserID)
|
||
}
|
||
|
||
if h.notifyConnectHook != nil {
|
||
go h.notifyConnectHook.PushUnreadTotalOnConnect(context.Background(), claims.UserID)
|
||
}
|
||
}
|
||
|
||
// createReadHandler 创建带生命周期管理的消息处理函数
|
||
// 优先查 Hub 事件路由表(业务模块注册的处理器),未命中再走内置 fallback
|
||
func (h *Handler) createReadHandler(userID int64) ws.MessageHandler {
|
||
return func(client *ws.Client, msg *ws.Message) {
|
||
funcName := "ws.handler.onMessage"
|
||
logs.Debug(nil, funcName, "收到 WebSocket 消息",
|
||
zap.Int64("user_id", client.UserID),
|
||
zap.String("event", msg.Event),
|
||
zap.Int64("seq", msg.Seq))
|
||
|
||
// 优先查事件路由表(IM、Meeting 等模块注册的处理器)
|
||
if h.hub.DispatchEvent(client, msg) {
|
||
return
|
||
}
|
||
|
||
// 内置事件 fallback
|
||
switch msg.Event {
|
||
case "heartbeat":
|
||
h.onlineService.HeartbeatRenew(context.Background(), userID)
|
||
resp := ws.NewResponse(msg.Event, msg.Seq, 0, "pong", nil)
|
||
data, err := ws.MarshalResponse(resp)
|
||
if err != nil {
|
||
logs.Error(nil, funcName, "序列化心跳响应失败", zap.Error(err))
|
||
return
|
||
}
|
||
client.Send(data)
|
||
default:
|
||
logs.Warn(nil, funcName, "未知事件类型",
|
||
zap.String("event", msg.Event),
|
||
zap.Int64("user_id", client.UserID))
|
||
resp := ws.NewResponse(msg.Event, msg.Seq, -1, "未知事件", nil)
|
||
data, err := ws.MarshalResponse(resp)
|
||
if err != nil {
|
||
logs.Error(nil, funcName, "序列化响应失败", zap.Error(err))
|
||
return
|
||
}
|
||
client.Send(data)
|
||
}
|
||
}
|
||
}
|
||
|
||
// GetHub 返回 Hub 实例(供在线状态等模块访问)
|
||
func (h *Handler) GetHub() *ws.Hub {
|
||
return h.hub
|
||
}
|
||
|
||
// GetPubSub 返回 PubSub 实例(供业务模块发送推送)
|
||
func (h *Handler) GetPubSub() *ws.PubSub {
|
||
return h.pubsub
|
||
}
|