Files
EchoChat/backend/go-service/app/meeting/controller/meeting_ws_handler.go
bujinyuan ea2bf96c2f fix(meeting): Task 16 修复 code-reviewer 审计 P1 八项
按 docs/reviews/2026-04-23-phase2e-2-code-review.md 清单落地 P1 批次:

- P1-1 主持人转让原子性:MeetingParticipantDAO.TransferHost 内部事务
  合并 meeting_rooms.host_id 更新,service/lifecycle 移除冗余 UpdateHost
- P1-2 ListChatMessages / ListMyMeetings 回归 DAO:新增
  MeetingChatDAO.ListByRoomBefore 反向游标,service 移除 s.db 直查
- P1-3 ListMyMeetings N+1 放大优化:新增
  MeetingParticipantDAO.ListJoinedRoomsByUser JOIN + DISTINCT 单次 SQL
- P1-4 EndRoom 行锁 + 事务:s.db.Transaction 包裹
  SELECT FOR UPDATE → 快照 → MarkEnded → LeaveAllActive;DAO 补
  WithTx(*gorm.DB) 和 GetByIDForUpdate
- P1-5 后台 goroutine trace_id 保留:新增 logs.DetachContext(ctx);
  替换 meeting_service / meeting_signal_service 共 10 处 context.Background();
  meeting_ws_handler 每条 WS 消息分配独立 trace_id
- P1-6 _broadcastSelfState 重试:最多 2 次指数退避(700ms→2100ms),
  检测 localAudioEnabled/localVideoEnabled 已被后续动作覆盖时放弃旧 patch
- P1-7 HandleHostGraceExpired / HandleEmptyRoomExpired 处理锁:
  独立 host_grace_handling:<code> / empty_ttl_handling:<code> SETNX + 60s TTL,
  消除 Redis 自然过期 + 本地 timer 到点 + DEL 返回 0 的盲区
- P1-8 pushExistingRoomState 同步化:OnRoomJoin 返回前完成补推,
  消除与 REST member.joined 并行造成的前端状态闪烁

并附带修复 .gitignore 中 logs/ 规则误伤 pkg/logs/ 代码目录的历史遗留问题,
将 logger.go / trace.go 正式纳入版本控制。

验证:backend go build + go vet 通过;frontend npm run build:h5 通过。

Made-with: Cursor
2026-04-23 16:53:30 +08:00

220 lines
7.8 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 controller 会议模块 HTTP / WS 入口
package controller
import (
"context"
"encoding/json"
"github.com/echochat/backend/app/constants"
"github.com/echochat/backend/app/meeting/service"
"github.com/echochat/backend/pkg/logs"
"github.com/echochat/backend/pkg/ws"
"go.uber.org/zap"
)
// MeetingWSHandler 会议 WS 事件入口
// Task 6 落地:注册设计 §6.3 的 11 个 meeting.* 事件到 Hub 路由表
// 薄层设计:
// - 仅负责 JSON Unmarshal → 调 SignalService 业务方法 → 组装 ACK 响应
// - 权限校验全部下沉到 SignalServiceservice 层统一 assertIsActiveParticipant / assertIsHost
// - 广播副作用由 SignalService 自行发起ACK 不携带广播数据)
type MeetingWSHandler struct {
signalSvc *service.MeetingSignalService
hub *ws.Hub
}
// NewMeetingWSHandler 构造 WS Handler 并自动注册事件到 Hub
// 与 im.handler.EventHandler 保持一致的风格:构造时注册,启动时由 DI 容器持有
func NewMeetingWSHandler(signalSvc *service.MeetingSignalService, hub *ws.Hub) *MeetingWSHandler {
h := &MeetingWSHandler{
signalSvc: signalSvc,
hub: hub,
}
h.registerEvents()
return h
}
// registerEvents 将所有 meeting.* C→S 事件注册到 Hub 路由表
// S→C 事件ended/joined/left/host.changed/producer.new 等)不在此处注册,由 broadcaster 发出
func (h *MeetingWSHandler) registerEvents() {
// 房间组
h.hub.RegisterEvent(constants.MeetingWSEventRoomJoin, h.handleRoomJoin)
h.hub.RegisterEvent(constants.MeetingWSEventRoomLeave, h.handleRoomLeave)
// 成员组
h.hub.RegisterEvent(constants.MeetingWSEventMemberStateChange, h.handleMemberStateChange)
// 媒体组5 个)
h.hub.RegisterEvent(constants.MeetingWSEventTransportCreate, h.handleTransportCreate)
h.hub.RegisterEvent(constants.MeetingWSEventTransportConnect, h.handleTransportConnect)
h.hub.RegisterEvent(constants.MeetingWSEventProduceStart, h.handleProduceStart)
h.hub.RegisterEvent(constants.MeetingWSEventConsumeStart, h.handleConsumeStart)
h.hub.RegisterEvent(constants.MeetingWSEventConsumeResume, h.handleConsumeResume)
h.hub.RegisterEvent(constants.MeetingWSEventProducerClose, h.handleProducerClose)
}
// ====== 通用工具 ======
// simpleRoomPayload 房间组仅需 room_code 的请求体
type simpleRoomPayload struct {
RoomCode string `json:"room_code"`
}
// newEventContext 为每个 WS 事件生成独立 ctx携带新 trace_id
// Task 16 P1-5WS 层无 HTTP 中间件注入 trace_id旧版 handler 一律 context.Background()
// 导致 signalSvc → Redis / DB / broadcaster 的整条链路日志无法通过 trace_id 串联。
// 现改为每条 WS 消息生成独立 UUID 作为 trace_id相当于 RPC 语义;配合 service 层
// logs.DetachContext 用于后台 goroutine可形成完整"WS 消息 → 业务处理 → 异步广播"的日志链路。
func (h *MeetingWSHandler) newEventContext() context.Context {
return logs.WithTraceID(context.Background(), logs.GenerateTraceID())
}
// sendACK 统一发送 ACK 响应
func (h *MeetingWSHandler) sendACK(client *ws.Client, msg *ws.Message, code int, message string, data interface{}) {
resp := ws.NewResponse(msg.Event, msg.Seq, code, message, data)
bytes, err := ws.MarshalResponse(resp)
if err != nil {
logs.Error(context.Background(), "controller.meeting_ws_handler.sendACK", "序列化 ACK 失败",
zap.String("event", msg.Event), zap.Error(err))
return
}
client.Send(bytes)
}
// unmarshal 通用反序列化 + 错误 ACK
func (h *MeetingWSHandler) unmarshal(client *ws.Client, msg *ws.Message, target interface{}) bool {
if err := json.Unmarshal(msg.Data, target); err != nil {
logs.Warn(nil, "controller.meeting_ws_handler.unmarshal", "反序列化失败",
zap.String("event", msg.Event),
zap.Int64("user_id", client.UserID),
zap.Error(err))
h.sendACK(client, msg, -1, "请求参数格式错误", nil)
return false
}
return true
}
// ====== 房间组 ======
// handleRoomJoin 处理 meeting.room.join
func (h *MeetingWSHandler) handleRoomJoin(client *ws.Client, msg *ws.Message) {
var payload simpleRoomPayload
if !h.unmarshal(client, msg, &payload) {
return
}
if err := h.signalSvc.OnRoomJoin(h.newEventContext(), client.UserID, payload.RoomCode); err != nil {
h.sendACK(client, msg, -1, err.Error(), nil)
return
}
h.sendACK(client, msg, 0, "ok", nil)
}
// handleRoomLeave 处理 meeting.room.leave
func (h *MeetingWSHandler) handleRoomLeave(client *ws.Client, msg *ws.Message) {
var payload simpleRoomPayload
if !h.unmarshal(client, msg, &payload) {
return
}
if err := h.signalSvc.OnRoomLeave(h.newEventContext(), client.UserID, payload.RoomCode); err != nil {
h.sendACK(client, msg, -1, err.Error(), nil)
return
}
h.sendACK(client, msg, 0, "ok", nil)
}
// ====== 成员组 ======
// handleMemberStateChange 处理 meeting.member.state.changed
func (h *MeetingWSHandler) handleMemberStateChange(client *ws.Client, msg *ws.Message) {
var payload service.MemberStateChangePayload
if !h.unmarshal(client, msg, &payload) {
return
}
if err := h.signalSvc.OnMemberStateChanged(h.newEventContext(), client.UserID, &payload); err != nil {
h.sendACK(client, msg, -1, err.Error(), nil)
return
}
h.sendACK(client, msg, 0, "ok", nil)
}
// ====== 媒体组 ======
// handleTransportCreate 处理 meeting.transport.create
func (h *MeetingWSHandler) handleTransportCreate(client *ws.Client, msg *ws.Message) {
var payload service.TransportCreatePayload
if !h.unmarshal(client, msg, &payload) {
return
}
info, err := h.signalSvc.OnTransportCreate(h.newEventContext(), client.UserID, &payload)
if err != nil {
h.sendACK(client, msg, -1, err.Error(), nil)
return
}
h.sendACK(client, msg, 0, "ok", info)
}
// handleTransportConnect 处理 meeting.transport.connect
func (h *MeetingWSHandler) handleTransportConnect(client *ws.Client, msg *ws.Message) {
var payload service.TransportConnectPayload
if !h.unmarshal(client, msg, &payload) {
return
}
if err := h.signalSvc.OnTransportConnect(h.newEventContext(), client.UserID, &payload); err != nil {
h.sendACK(client, msg, -1, err.Error(), nil)
return
}
h.sendACK(client, msg, 0, "ok", nil)
}
// handleProduceStart 处理 meeting.produce.start
func (h *MeetingWSHandler) handleProduceStart(client *ws.Client, msg *ws.Message) {
var payload service.ProduceStartPayload
if !h.unmarshal(client, msg, &payload) {
return
}
result, err := h.signalSvc.OnProduceStart(h.newEventContext(), client.UserID, &payload)
if err != nil {
h.sendACK(client, msg, -1, err.Error(), nil)
return
}
h.sendACK(client, msg, 0, "ok", result)
}
// handleConsumeStart 处理 meeting.consume.start
func (h *MeetingWSHandler) handleConsumeStart(client *ws.Client, msg *ws.Message) {
var payload service.ConsumeStartPayload
if !h.unmarshal(client, msg, &payload) {
return
}
info, err := h.signalSvc.OnConsumeStart(h.newEventContext(), client.UserID, &payload)
if err != nil {
h.sendACK(client, msg, -1, err.Error(), nil)
return
}
h.sendACK(client, msg, 0, "ok", info)
}
// handleConsumeResume 处理 meeting.consume.resumeTask 9
func (h *MeetingWSHandler) handleConsumeResume(client *ws.Client, msg *ws.Message) {
var payload service.ConsumeResumePayload
if !h.unmarshal(client, msg, &payload) {
return
}
if err := h.signalSvc.OnConsumeResume(h.newEventContext(), client.UserID, &payload); err != nil {
h.sendACK(client, msg, -1, err.Error(), nil)
return
}
h.sendACK(client, msg, 0, "ok", nil)
}
// handleProducerClose 处理 meeting.producer.close
func (h *MeetingWSHandler) handleProducerClose(client *ws.Client, msg *ws.Message) {
var payload service.ProducerClosePayload
if !h.unmarshal(client, msg, &payload) {
return
}
if err := h.signalSvc.OnProducerClose(h.newEventContext(), client.UserID, &payload); err != nil {
h.sendACK(client, msg, -1, err.Error(), nil)
return
}
h.sendACK(client, msg, 0, "ok", nil)
}