Files
EchoChat/backend/go-service/app/meeting/dao/meeting_room_dao.go
bujinyuan 3b83c79036 feat(phase2e-2): 落地会议生命周期状态机 + Router 幂等双层防御(Task 8)
核心交付:
- 新建 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
2026-04-21 18:21:04 +08:00

204 lines
7.4 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 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
}