核心交付:
- 新建 MeetingLifecycleService(6 钩子 + sync.Map 本地 timer + Redis key 双保险 + RescheduleFromRedis)
- 新建 MeetingCleanupTask(启动重建 timer + 每 N 秒扫 host_grace/empty_ttl 兜底 + 4h stale active 回收)
- MediaOrchestrator 新增 ResolveRouterID;HTTPMediaOrchestrator.CreateRouter 入口 sync.Map 幂等防御
- 业务层 JoinRoom 移除 CreateRouter 调用改走 CancelEmptyTTL + ResolveRouterID;LeaveRoom 空房分支改调 OnAllMembersLeft 不再立即销毁
- MeetingSignalService 新增 OnWSDisconnect 实现 ws.MeetingDisconnectHook;OnRoomJoin 追加 host 重连钩子
- ws.handler 定义 MeetingDisconnectHook 接口 + SetMeetingDisconnectHook,解耦 ws→meeting 反向依赖
- config 新增 MeetingConfig{HostGrace=120, EmptyRoomTTL=300, CleanupInterval=30, StaleRoomHours=4}
关键设计决策:
- Redis key TTL = 业务时长 + max(CleanupIntervalSeconds*2, 30s) buffer:避免本地 timer 与
Redis 自动过期同步到期导致 DEL 返回 0 被误判为"已被其他路径处理"而跳过业务逻辑
- Router 幂等双层防御(决策 q2_router_dedup=a2_both):业务层不重复调 + HTTP 层 sync.Map 命中直接返回
- 普通成员 WS 断开仅清 media 资源不动 participant 表(决策 q1_nonhost_disconnect=a1_keep_current)
E2E 验证:docs/verify/meeting_t8_verify.mjs PASS=20 FAIL=0,覆盖 5 场景:
- S1 host 宽限期过期自动转让(meeting.host.changed + DB host_id 更新)
- S2 宽限期内重连保留身份
- S3 empty_ttl 期内新成员加入复活房间
- S4 empty_ttl 过期 → 房间 Ended + 新 join 被拒
- S5 CreateRoom +1 Router / JoinRoom 不再创建新 Router(通过 media-server /internal/info stats.routers 断言)
media-server:/internal/info 响应追加 stats.routers + routers[] 供 E2E 断言 Router 幂等
文档同步:
- docs/progress/CURRENT_STATUS.md 头部 + 新增 Task 8 交付条目
- docs/plans/2026-04-21-phase2e-2-implementation.plan.md Task 8 标记完成 + 实际产出/决策/验证
- docs/api/frontend/meeting.md 补充 host.changed.auto_reason / room.ended.reason=system_error / 空房 TTL 复活语义 + Task 8 验证记录
- docs/architecture/system-architecture.md meeting 模块职责补充"会议生命周期状态机"
- .cursor/rules/project-context.mdc 追加 Task 8 条目并更新 Phase 2e-2 进度(Task 0-8 ✅)
Made-with: Cursor
204 lines
7.4 KiB
Go
204 lines
7.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"
|
||
)
|
||
|
||
// MeetingRoomDAO 会议房间数据访问对象
|
||
type MeetingRoomDAO struct {
|
||
db *gorm.DB
|
||
}
|
||
|
||
// NewMeetingRoomDAO 创建 MeetingRoomDAO 实例
|
||
func NewMeetingRoomDAO(db *gorm.DB) *MeetingRoomDAO {
|
||
return &MeetingRoomDAO{db: db}
|
||
}
|
||
|
||
// 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
|
||
}
|