Files
EchoChat/backend/go-service/app/ws/handler.go
bujinyuan c35097a0d8 fix(meeting): Task 16 修复 code-reviewer 审计 P2 七项 + Nit 七项
覆盖 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 处 + 历次 commit cdaa39d / ea2bf96 / f5ae095 / 5ed14c2)
- 推迟登记表(P2-4 / P2-5 / 端口收敛 / appData 校验 / RFC3339 时间格式,共 5 项)

Made-with: Cursor
2026-04-23 17:45:03 +08:00

262 lines
8.9 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// 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 NitCheckOrigin 白名单需要
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 8WS 最终下线时触发 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
}