按 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
230 lines
8.4 KiB
Go
230 lines
8.4 KiB
Go
// Package dao 提供 meeting 模块的数据库访问操作
|
||
package dao
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"time"
|
||
|
||
"github.com/echochat/backend/app/constants"
|
||
"github.com/echochat/backend/app/meeting/model"
|
||
"github.com/echochat/backend/pkg/logs"
|
||
"go.uber.org/zap"
|
||
"gorm.io/gorm"
|
||
"gorm.io/gorm/clause"
|
||
)
|
||
|
||
// MeetingRoomDAO 会议房间数据访问对象
|
||
type MeetingRoomDAO struct {
|
||
db *gorm.DB
|
||
}
|
||
|
||
// NewMeetingRoomDAO 创建 MeetingRoomDAO 实例
|
||
func NewMeetingRoomDAO(db *gorm.DB) *MeetingRoomDAO {
|
||
return &MeetingRoomDAO{db: db}
|
||
}
|
||
|
||
// WithTx 返回绑定到指定事务句柄的新 DAO 实例
|
||
// 用法:service 层使用 s.db.Transaction(func(tx *gorm.DB) error { return s.roomDAO.WithTx(tx).MarkEnded(...) })
|
||
// 所有 DAO 方法都通过 d.db.WithContext(ctx) 派生查询,替换 d.db 为 tx 即可让现有方法参与上游事务。
|
||
// Task 16 P1-4:EndRoom 行锁 + 事务化 新增
|
||
func (d *MeetingRoomDAO) WithTx(tx *gorm.DB) *MeetingRoomDAO {
|
||
return &MeetingRoomDAO{db: tx}
|
||
}
|
||
|
||
// GetByIDForUpdate 事务内 SELECT ... FOR UPDATE 行锁查询
|
||
// 仅在上游已开启事务(调用方通过 WithTx(tx) 注入)时使用;否则等同于普通查询不生效。
|
||
// Task 16 P1-4:EndRoom 行锁 新增
|
||
func (d *MeetingRoomDAO) GetByIDForUpdate(ctx context.Context, id int64) (*model.MeetingRoom, error) {
|
||
var room model.MeetingRoom
|
||
err := d.db.WithContext(ctx).
|
||
Clauses(clause.Locking{Strength: "UPDATE"}).
|
||
First(&room, id).Error
|
||
if err != nil {
|
||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||
return nil, nil
|
||
}
|
||
return nil, err
|
||
}
|
||
return &room, nil
|
||
}
|
||
|
||
// Create 新建一个会议房间
|
||
// 由调用方预先生成 RoomCode 并确认唯一,本方法不做冲突重试
|
||
func (d *MeetingRoomDAO) Create(ctx context.Context, room *model.MeetingRoom) error {
|
||
funcName := "dao.meeting_room_dao.Create"
|
||
logs.Info(ctx, funcName, "创建会议房间",
|
||
zap.String("room_code", room.RoomCode), zap.Int64("host_id", room.HostID))
|
||
|
||
err := d.db.WithContext(ctx).Create(room).Error
|
||
if err != nil {
|
||
logs.Error(ctx, funcName, "创建会议房间失败",
|
||
zap.String("room_code", room.RoomCode), zap.Error(err))
|
||
}
|
||
return err
|
||
}
|
||
|
||
// GetByID 按主键查询
|
||
// 记录不存在时返回 (nil, nil),便于上层直接用 room == nil 判空并返回业务错误
|
||
func (d *MeetingRoomDAO) GetByID(ctx context.Context, id int64) (*model.MeetingRoom, error) {
|
||
var room model.MeetingRoom
|
||
err := d.db.WithContext(ctx).First(&room, id).Error
|
||
if err != nil {
|
||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||
return nil, nil
|
||
}
|
||
return nil, err
|
||
}
|
||
return &room, nil
|
||
}
|
||
|
||
// GetByCode 按会议号查询(入会流程的主要入口)
|
||
// 记录不存在时返回 (nil, nil),上层应通过 room == nil 判定并返回业务错误 ErrMeetingNotFound
|
||
func (d *MeetingRoomDAO) GetByCode(ctx context.Context, code string) (*model.MeetingRoom, error) {
|
||
funcName := "dao.meeting_room_dao.GetByCode"
|
||
|
||
var room model.MeetingRoom
|
||
err := d.db.WithContext(ctx).Where("room_code = ?", code).First(&room).Error
|
||
if err != nil {
|
||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||
return nil, nil
|
||
}
|
||
logs.Error(ctx, funcName, "按会议号查询房间失败",
|
||
zap.String("room_code", code), zap.Error(err))
|
||
return nil, err
|
||
}
|
||
return &room, nil
|
||
}
|
||
|
||
// ExistsCode 会议号是否已被占用(用于创建会议时的 code 冲突重试)
|
||
func (d *MeetingRoomDAO) ExistsCode(ctx context.Context, code string) (bool, error) {
|
||
var count int64
|
||
err := d.db.WithContext(ctx).
|
||
Model(&model.MeetingRoom{}).
|
||
Where("room_code = ?", code).
|
||
Count(&count).Error
|
||
return count > 0, err
|
||
}
|
||
|
||
// MarkStarted 记录会议实际开始时间(host 首次加入时调用)
|
||
// 仅当 status=pending 且 started_at IS NULL 时才写入,避免重复覆盖
|
||
func (d *MeetingRoomDAO) MarkStarted(ctx context.Context, id int64, startedAt time.Time) (int64, error) {
|
||
funcName := "dao.meeting_room_dao.MarkStarted"
|
||
res := d.db.WithContext(ctx).
|
||
Model(&model.MeetingRoom{}).
|
||
Where("id = ? AND status = ? AND started_at IS NULL", id, constants.MeetingStatusPending).
|
||
Updates(map[string]interface{}{
|
||
"status": constants.MeetingStatusActive,
|
||
"started_at": startedAt,
|
||
})
|
||
if res.Error != nil {
|
||
logs.Error(ctx, funcName, "标记会议开始失败",
|
||
zap.Int64("room_id", id), zap.Error(res.Error))
|
||
}
|
||
return res.RowsAffected, res.Error
|
||
}
|
||
|
||
// MarkEnded 记录会议结束(status=2 + ended_at + ended_reason)
|
||
// 乐观锁:只对 status != ended 的行生效,避免重复结束覆盖原始 reason
|
||
func (d *MeetingRoomDAO) MarkEnded(ctx context.Context, id int64, reason string, endedAt time.Time) (int64, error) {
|
||
funcName := "dao.meeting_room_dao.MarkEnded"
|
||
res := d.db.WithContext(ctx).
|
||
Model(&model.MeetingRoom{}).
|
||
Where("id = ? AND status != ?", id, constants.MeetingStatusEnded).
|
||
Updates(map[string]interface{}{
|
||
"status": constants.MeetingStatusEnded,
|
||
"ended_at": endedAt,
|
||
"ended_reason": reason,
|
||
})
|
||
if res.Error != nil {
|
||
logs.Error(ctx, funcName, "标记会议结束失败",
|
||
zap.Int64("room_id", id), zap.String("reason", reason), zap.Error(res.Error))
|
||
}
|
||
return res.RowsAffected, res.Error
|
||
}
|
||
|
||
// UpdateHost 更新主持人(仅修改 meeting_rooms.host_id 字段,meeting_participants.role 由 ParticipantDAO.TransferHost 在同一事务中处理)
|
||
func (d *MeetingRoomDAO) UpdateHost(ctx context.Context, id, newHostID int64) error {
|
||
return d.db.WithContext(ctx).
|
||
Model(&model.MeetingRoom{}).
|
||
Where("id = ?", id).
|
||
Update("host_id", newHostID).Error
|
||
}
|
||
|
||
// UpdateSettings 更新房间级配置(settings 字段整体替换)
|
||
// 调用方需保证 settingsJSON 是合法 JSON 字符串
|
||
func (d *MeetingRoomDAO) UpdateSettings(ctx context.Context, id int64, settingsJSON string) error {
|
||
return d.db.WithContext(ctx).
|
||
Model(&model.MeetingRoom{}).
|
||
Where("id = ?", id).
|
||
Update("settings", settingsJSON).Error
|
||
}
|
||
|
||
// ListByHost 按主持人查询会议列表(支持 status 过滤 + 分页)
|
||
// status 传 -1 表示不限
|
||
// 返回按 created_at DESC 排序,走 idx_meeting_rooms_host_status 索引
|
||
func (d *MeetingRoomDAO) ListByHost(ctx context.Context, hostID int64, status int, offset, limit int) ([]model.MeetingRoom, int64, error) {
|
||
funcName := "dao.meeting_room_dao.ListByHost"
|
||
q := d.db.WithContext(ctx).
|
||
Model(&model.MeetingRoom{}).
|
||
Where("host_id = ?", hostID)
|
||
if status >= 0 {
|
||
q = q.Where("status = ?", status)
|
||
}
|
||
|
||
var total int64
|
||
if err := q.Count(&total).Error; err != nil {
|
||
logs.Error(ctx, funcName, "计数失败", zap.Int64("host_id", hostID), zap.Error(err))
|
||
return nil, 0, err
|
||
}
|
||
|
||
var rooms []model.MeetingRoom
|
||
err := q.Order("created_at DESC").Offset(offset).Limit(limit).Find(&rooms).Error
|
||
if err != nil {
|
||
logs.Error(ctx, funcName, "查询失败", zap.Int64("host_id", hostID), zap.Error(err))
|
||
}
|
||
return rooms, total, err
|
||
}
|
||
|
||
// ListStaleActive 返回 status != Ended 且 started_at(或 created_at 若未开始)早于 hoursAgo 小时的房间 ID
|
||
// 用于定时任务兜底清理"创建 / Active 但长期无任何成员活动的僵尸房间"(设计 §11.4 的 4 小时兜底)
|
||
// 排序:最早 started_at 在前,优先清理最老
|
||
func (d *MeetingRoomDAO) ListStaleActive(ctx context.Context, hoursAgo int, limit int) ([]int64, error) {
|
||
funcName := "dao.meeting_room_dao.ListStaleActive"
|
||
cutoff := time.Now().Add(-time.Duration(hoursAgo) * time.Hour)
|
||
|
||
var ids []int64
|
||
err := d.db.WithContext(ctx).
|
||
Model(&model.MeetingRoom{}).
|
||
Where("status != ? AND COALESCE(started_at, created_at) < ?",
|
||
constants.MeetingStatusEnded, cutoff).
|
||
Order("COALESCE(started_at, created_at) ASC").
|
||
Limit(limit).
|
||
Pluck("id", &ids).Error
|
||
if err != nil {
|
||
logs.Error(ctx, funcName, "查询 stale 活跃房间失败", zap.Error(err))
|
||
}
|
||
return ids, err
|
||
}
|
||
|
||
// ListExpiredForCleanup 返回已结束且超过指定小时数的房间 ID 列表
|
||
// 供定时任务清理 meeting_chats 使用;room 本身保留归档不删
|
||
// 返回 limit 上限用于分批处理,避免一次捞太多
|
||
func (d *MeetingRoomDAO) ListExpiredForCleanup(ctx context.Context, hoursAgo int, limit int) ([]int64, error) {
|
||
funcName := "dao.meeting_room_dao.ListExpiredForCleanup"
|
||
cutoff := time.Now().Add(-time.Duration(hoursAgo) * time.Hour)
|
||
|
||
var ids []int64
|
||
err := d.db.WithContext(ctx).
|
||
Model(&model.MeetingRoom{}).
|
||
Where("status = ? AND ended_at IS NOT NULL AND ended_at < ?",
|
||
constants.MeetingStatusEnded, cutoff).
|
||
Order("ended_at ASC").
|
||
Limit(limit).
|
||
Pluck("id", &ids).Error
|
||
if err != nil {
|
||
logs.Error(ctx, funcName, "查询过期房间失败", zap.Error(err))
|
||
}
|
||
return ids, err
|
||
}
|