按 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
220 lines
7.8 KiB
Go
220 lines
7.8 KiB
Go
// 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 响应
|
||
// - 权限校验全部下沉到 SignalService(service 层统一 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-5:WS 层无 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.resume(Task 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)
|
||
}
|