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
This commit is contained in:
@@ -160,6 +160,27 @@ func (d *MeetingRoomDAO) ListByHost(ctx context.Context, hostID int64, status in
|
||||
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 上限用于分批处理,避免一次捞太多
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"github.com/echochat/backend/app/meeting/controller"
|
||||
"github.com/echochat/backend/app/meeting/dao"
|
||||
"github.com/echochat/backend/app/meeting/service"
|
||||
"github.com/echochat/backend/app/meeting/task"
|
||||
"github.com/google/wire"
|
||||
)
|
||||
|
||||
@@ -20,10 +21,12 @@ var MeetingSet = wire.NewSet(
|
||||
dao.NewMeetingParticipantDAO,
|
||||
dao.NewMeetingChatDAO,
|
||||
service.NewMeetingBroadcaster,
|
||||
service.NewMeetingLifecycleService, // Task 8 新增:会议生命周期状态机
|
||||
service.NewMeetingService,
|
||||
service.NewMeetingSignalService,
|
||||
controller.NewMeetingController,
|
||||
controller.NewMeetingWSHandler,
|
||||
task.NewMeetingCleanupTask, // Task 8 新增:生命周期兜底定时任务
|
||||
|
||||
// MediaOrchestrator 使用真实 HTTP 实现(Task 7 落地)
|
||||
// 注入 *config.Config,通过 cfg.MediaServer 获取 base_url / internal_token / timeout 等配置
|
||||
|
||||
@@ -85,9 +85,21 @@ func NewHTTPMediaOrchestrator(cfg *config.Config) *HTTPMediaOrchestrator {
|
||||
|
||||
// CreateRouter 调用 POST /internal/v1/routers
|
||||
// 成功后在本地缓存 roomCode → routerID 映射供 CloseRouter 使用
|
||||
// Task 8 幂等防御(双层,与业务层 JoinRoom 不再主动调用形成互补):
|
||||
// - 命中缓存直接返回已缓存 routerID,不发起 HTTP 请求
|
||||
// - 仅在缓存缺失时真正调 Node;避免业务误多次调用产生 Node 侧重复 Router
|
||||
func (h *HTTPMediaOrchestrator) CreateRouter(ctx context.Context, roomCode string) (string, error) {
|
||||
funcName := "service.http_media_orchestrator.CreateRouter"
|
||||
|
||||
if cached, ok := h.roomRouterIDs.Load(roomCode); ok {
|
||||
if rid, _ := cached.(string); rid != "" {
|
||||
logs.Debug(ctx, funcName, "命中本地 Router 缓存,跳过 Node 调用",
|
||||
zap.String("room_code", roomCode),
|
||||
zap.String("router_id", rid))
|
||||
return rid, nil
|
||||
}
|
||||
}
|
||||
|
||||
reqBody := map[string]any{"roomCode": roomCode}
|
||||
var resp struct {
|
||||
RouterID string `json:"routerId"`
|
||||
@@ -104,8 +116,17 @@ func (h *HTTPMediaOrchestrator) CreateRouter(ctx context.Context, roomCode strin
|
||||
return "", err
|
||||
}
|
||||
|
||||
// 记忆映射:一个 roomCode 仅对应一个 Router;若此前存在旧 Router ID(极少见),以新值覆盖
|
||||
h.roomRouterIDs.Store(roomCode, resp.RouterID)
|
||||
// LoadOrStore 防止并发首次创建时重复覆盖:若他人已写入先到的值就用已存在的
|
||||
actual, loaded := h.roomRouterIDs.LoadOrStore(roomCode, resp.RouterID)
|
||||
if loaded {
|
||||
if existing, _ := actual.(string); existing != "" && existing != resp.RouterID {
|
||||
logs.Warn(ctx, funcName, "并发创建 Router 出现不一致,保留先到值",
|
||||
zap.String("room_code", roomCode),
|
||||
zap.String("existing", existing),
|
||||
zap.String("new", resp.RouterID))
|
||||
return existing, nil
|
||||
}
|
||||
}
|
||||
|
||||
logs.Info(ctx, funcName, "Node Router 创建成功",
|
||||
zap.String("room_code", roomCode),
|
||||
@@ -113,6 +134,21 @@ func (h *HTTPMediaOrchestrator) CreateRouter(ctx context.Context, roomCode strin
|
||||
return resp.RouterID, nil
|
||||
}
|
||||
|
||||
// ResolveRouterID 从本地缓存读取 roomCode 对应的 routerID(不触发 HTTP)
|
||||
// 返回 (id, true) 表示命中;返回 (_, false) 表示缓存缺失(通常意味着该房间尚未创建 Router)
|
||||
// 供 MeetingService.JoinRoom 等需复用已有 Router 的场景使用,避免重复调 Node
|
||||
func (h *HTTPMediaOrchestrator) ResolveRouterID(roomCode string) (string, bool) {
|
||||
v, ok := h.roomRouterIDs.Load(roomCode)
|
||||
if !ok {
|
||||
return "", false
|
||||
}
|
||||
rid, _ := v.(string)
|
||||
if rid == "" {
|
||||
return "", false
|
||||
}
|
||||
return rid, true
|
||||
}
|
||||
|
||||
// CloseRouter 调用 DELETE /internal/v1/routers/:routerId
|
||||
// 入参是 roomCode(符合设计 §6.6 NodeClient 契约),内部反查本地缓存拿到 routerID
|
||||
// 若本地缓存不存在映射(可能是 go-service 重启后状态丢失),直接返回 nil:
|
||||
|
||||
@@ -86,9 +86,14 @@ type CreateConsumerReq struct {
|
||||
// - 所有方法均幂等:重复调用不报错,符合 WS 信令重试语义
|
||||
type MediaOrchestrator interface {
|
||||
// CreateRouter 为会议房间创建 mediasoup Router
|
||||
// Task 8 起实现需自带 sync.Map 幂等缓存:roomCode 已有 Router 时直接返回缓存值
|
||||
CreateRouter(ctx context.Context, roomCode string) (routerID string, err error)
|
||||
// CloseRouter 关闭房间对应 Router 及其下所有资源(幂等)
|
||||
CloseRouter(ctx context.Context, roomCode string) error
|
||||
// ResolveRouterID 从本地缓存查询 roomCode 对应的 routerID(不触发 HTTP 调用)
|
||||
// Task 8 引入:用于 JoinRoom 等"仅需复用"的场景,避免重复调 CreateRouter
|
||||
// 返回 (id, true) 命中;(_, false) 缺失(通常说明房间未创建 Router 或服务重启后未 reschedule)
|
||||
ResolveRouterID(roomCode string) (string, bool)
|
||||
|
||||
// CreateTransport 为用户创建 send/recv WebRTC Transport
|
||||
CreateTransport(ctx context.Context, req *CreateTransportReq) (*TransportInfo, error)
|
||||
@@ -127,6 +132,11 @@ func (n *NoopMediaOrchestrator) CloseRouter(_ context.Context, _ string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// ResolveRouterID 占位实现:始终命中,返回 noop-router-{code}
|
||||
func (n *NoopMediaOrchestrator) ResolveRouterID(roomCode string) (string, bool) {
|
||||
return "noop-router-" + roomCode, true
|
||||
}
|
||||
|
||||
// CreateTransport 占位:返回以 "noop-transport-" 为前缀的伪造 ID,带最小合法 JSON 结构
|
||||
func (n *NoopMediaOrchestrator) CreateTransport(_ context.Context, req *CreateTransportReq) (*TransportInfo, error) {
|
||||
return &TransportInfo{
|
||||
|
||||
@@ -0,0 +1,475 @@
|
||||
// Package service 提供 meeting 模块的业务逻辑
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/echochat/backend/app/constants"
|
||||
"github.com/echochat/backend/app/meeting/dao"
|
||||
"github.com/echochat/backend/config"
|
||||
"github.com/echochat/backend/pkg/logs"
|
||||
"github.com/redis/go-redis/v9"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
// Redis key 前缀(设计文档 §5.4 生命周期 keys;Task 8 落地)
|
||||
const (
|
||||
redisKeyHostGracePrefix = "echo:meeting:host_grace:"
|
||||
redisKeyEmptyTTLPrefix = "echo:meeting:empty_ttl:"
|
||||
)
|
||||
|
||||
// hostGracePayload host 宽限期 Redis value 结构
|
||||
type hostGracePayload struct {
|
||||
HostID int64 `json:"host_id"`
|
||||
StartedAt int64 `json:"started_at"` // 秒级 unix
|
||||
GraceUntil int64 `json:"grace_until"` // 秒级 unix
|
||||
}
|
||||
|
||||
// MeetingLifecycleService 会议生命周期状态机
|
||||
// 负责设计 §6.5 的 5 类状态跃迁副作用(host 宽限期、自动转让、空房 TTL、TTL 过期销毁)
|
||||
// 采用 time.AfterFunc 本地 timer + Redis key TTL 双保险:
|
||||
// - 本地 timer 负责低延迟触发;服务重启由 RescheduleFromRedis 按 Redis 剩余 PTTL 重新装载
|
||||
// - 后台定时任务 MeetingCleanupTask 每 N 秒扫描 Redis 兜底防御(timer 丢失 / 多实例部署重复触发场景)
|
||||
//
|
||||
// 并发保护:HandleHostGraceExpired / HandleEmptyRoomExpired 入口先 DEL Redis key
|
||||
// 通过 DEL 返回值确认唯一处理权(0 = 已被其他 goroutine/节点清走,直接 return 避免重复操作)
|
||||
type MeetingLifecycleService struct {
|
||||
roomDAO *dao.MeetingRoomDAO
|
||||
participantDAO *dao.MeetingParticipantDAO
|
||||
db dbExecutor
|
||||
redis *redis.Client
|
||||
broadcaster *MeetingBroadcaster
|
||||
media MediaOrchestrator
|
||||
cfg config.MeetingConfig
|
||||
|
||||
graceTimers sync.Map // roomCode -> *time.Timer(host 宽限期)
|
||||
emptyTTLTimers sync.Map // roomCode -> *time.Timer(空房 TTL)
|
||||
}
|
||||
|
||||
// dbExecutor MeetingLifecycleService 对事务 DB 的最小依赖抽象
|
||||
// 仅用于 HandleEmptyRoomExpired 的 MarkEnded 调用链路,避免强绑 *gorm.DB
|
||||
type dbExecutor interface{}
|
||||
|
||||
// NewMeetingLifecycleService 构造生命周期服务
|
||||
// media 由上游通过 wire 注入 HTTPMediaOrchestrator
|
||||
func NewMeetingLifecycleService(
|
||||
roomDAO *dao.MeetingRoomDAO,
|
||||
participantDAO *dao.MeetingParticipantDAO,
|
||||
redis *redis.Client,
|
||||
broadcaster *MeetingBroadcaster,
|
||||
media MediaOrchestrator,
|
||||
cfg *config.Config,
|
||||
) *MeetingLifecycleService {
|
||||
meetingCfg := cfg.Meeting
|
||||
// 零值兜底:防止 yaml 未配置时走 0 秒 TTL
|
||||
if meetingCfg.HostGraceSeconds <= 0 {
|
||||
meetingCfg.HostGraceSeconds = constants.MeetingHostGraceSeconds
|
||||
}
|
||||
if meetingCfg.EmptyRoomTTLSeconds <= 0 {
|
||||
meetingCfg.EmptyRoomTTLSeconds = constants.MeetingEmptyRoomTTLSeconds
|
||||
}
|
||||
if meetingCfg.CleanupIntervalSeconds <= 0 {
|
||||
meetingCfg.CleanupIntervalSeconds = 30
|
||||
}
|
||||
if meetingCfg.StaleRoomHours <= 0 {
|
||||
meetingCfg.StaleRoomHours = 4
|
||||
}
|
||||
return &MeetingLifecycleService{
|
||||
roomDAO: roomDAO,
|
||||
participantDAO: participantDAO,
|
||||
redis: redis,
|
||||
broadcaster: broadcaster,
|
||||
media: media,
|
||||
cfg: meetingCfg,
|
||||
}
|
||||
}
|
||||
|
||||
// Config 返回生效的生命周期配置(供定时任务读取扫描周期等)
|
||||
func (s *MeetingLifecycleService) Config() config.MeetingConfig {
|
||||
return s.cfg
|
||||
}
|
||||
|
||||
// HostGraceKey 返回 host 宽限期 Redis key
|
||||
func HostGraceKey(roomCode string) string {
|
||||
return redisKeyHostGracePrefix + roomCode
|
||||
}
|
||||
|
||||
// EmptyTTLKey 返回空房 TTL Redis key
|
||||
func EmptyTTLKey(roomCode string) string {
|
||||
return redisKeyEmptyTTLPrefix + roomCode
|
||||
}
|
||||
|
||||
// OnHostDisconnect host WS 断线钩子
|
||||
// 写入 host_grace key(TTL HostGraceSeconds)并启动本地 AfterFunc
|
||||
// 若已有同 roomCode 的 grace 记录(重复触发),幂等刷新 TTL + 替换 timer
|
||||
func (s *MeetingLifecycleService) OnHostDisconnect(ctx context.Context, roomCode string, hostID int64) {
|
||||
funcName := "service.meeting_lifecycle_service.OnHostDisconnect"
|
||||
|
||||
now := time.Now()
|
||||
payload := hostGracePayload{
|
||||
HostID: hostID,
|
||||
StartedAt: now.Unix(),
|
||||
GraceUntil: now.Add(time.Duration(s.cfg.HostGraceSeconds) * time.Second).Unix(),
|
||||
}
|
||||
raw, _ := json.Marshal(payload)
|
||||
|
||||
graceDur := time.Duration(s.cfg.HostGraceSeconds) * time.Second
|
||||
// Redis TTL 比本地 timer 多留一段 buffer(max(cleanup*2, 30s))
|
||||
// 保证本地 timer 先触发 → DEL 命中;即便本地 timer 丢失,cleanup 也能在 buffer 内兜底扫到
|
||||
redisTTL := graceDur + s.ttlBuffer()
|
||||
if err := s.redis.Set(ctx, HostGraceKey(roomCode), string(raw), redisTTL).Err(); err != nil {
|
||||
logs.Warn(ctx, funcName, "写入 host_grace key 失败",
|
||||
zap.String("room_code", roomCode), zap.Int64("host_id", hostID), zap.Error(err))
|
||||
return
|
||||
}
|
||||
|
||||
s.scheduleHostGraceTimer(roomCode, graceDur)
|
||||
logs.Info(ctx, funcName, "host 进入宽限期",
|
||||
zap.String("room_code", roomCode),
|
||||
zap.Int64("host_id", hostID),
|
||||
zap.Int("grace_seconds", s.cfg.HostGraceSeconds))
|
||||
}
|
||||
|
||||
// OnHostReconnect host 重连钩子
|
||||
// 清除 host_grace key 并取消本地 timer;用于 host 掉线后 TTL 内重新出现在 WS 层
|
||||
// 幂等:key 不存在也不报错
|
||||
func (s *MeetingLifecycleService) OnHostReconnect(ctx context.Context, roomCode string, hostID int64) {
|
||||
funcName := "service.meeting_lifecycle_service.OnHostReconnect"
|
||||
|
||||
deleted, err := s.redis.Del(ctx, HostGraceKey(roomCode)).Result()
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "清除 host_grace key 失败",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
}
|
||||
s.cancelTimer(&s.graceTimers, roomCode)
|
||||
if deleted > 0 {
|
||||
logs.Info(ctx, funcName, "host 宽限期内重连,撤销转让计划",
|
||||
zap.String("room_code", roomCode), zap.Int64("host_id", hostID))
|
||||
}
|
||||
}
|
||||
|
||||
// HandleHostGraceExpired host 宽限期过期处理
|
||||
// 触发时机:本地 timer 到期 或 cleanup task 扫到 key TTL <=0
|
||||
// 流程:
|
||||
// 1. DEL host_grace key → 返回 0 表示已被其他路径处理,直接 return 防止重复转让
|
||||
// 2. 加载 room + 活跃成员;若 room 已 Ended / host 已变更,视为幂等完成
|
||||
// 3. 活跃成员存在:选最早加入者转 host(事务),广播 host.changed,auto_reason=host_grace_expired
|
||||
// 4. 无活跃成员:走 OnAllMembersLeft 路径(置 empty_ttl)
|
||||
func (s *MeetingLifecycleService) HandleHostGraceExpired(ctx context.Context, roomCode string) {
|
||||
funcName := "service.meeting_lifecycle_service.HandleHostGraceExpired"
|
||||
|
||||
deleted, err := s.redis.Del(ctx, HostGraceKey(roomCode)).Result()
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "DEL host_grace key 失败(跳过本次处理)",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
return
|
||||
}
|
||||
s.cancelTimer(&s.graceTimers, roomCode)
|
||||
if deleted == 0 {
|
||||
logs.Debug(ctx, funcName, "host_grace key 不存在(已被重连/其他节点清走)",
|
||||
zap.String("room_code", roomCode))
|
||||
return
|
||||
}
|
||||
|
||||
room, err := s.roomDAO.GetByCode(ctx, roomCode)
|
||||
if err != nil || room == nil {
|
||||
logs.Warn(ctx, funcName, "加载会议失败或已不存在",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
return
|
||||
}
|
||||
if room.Status == constants.MeetingStatusEnded {
|
||||
logs.Info(ctx, funcName, "会议已结束,跳过 host 宽限期过期处理",
|
||||
zap.String("room_code", roomCode))
|
||||
return
|
||||
}
|
||||
|
||||
actives, err := s.participantDAO.ListActiveByRoom(ctx, room.ID)
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "拉取活跃成员失败",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
return
|
||||
}
|
||||
// 排除原 host(若 left_at IS NULL 但 WS 掉线未落库)
|
||||
candidates := make([]int64, 0, len(actives))
|
||||
for _, p := range actives {
|
||||
if p.UserID == room.HostID {
|
||||
continue
|
||||
}
|
||||
candidates = append(candidates, p.UserID)
|
||||
}
|
||||
|
||||
if len(candidates) == 0 {
|
||||
logs.Info(ctx, funcName, "host 宽限期过期且无其他活跃成员,转空房 TTL",
|
||||
zap.String("room_code", roomCode), zap.Int64("old_host_id", room.HostID))
|
||||
// 清掉 host 自己的记录(若 WS 掉线未落库,由兜底 LeaveRoom 处理)
|
||||
_, _ = s.participantDAO.LeaveRoom(ctx, room.ID, room.HostID, constants.MeetingLeftReasonDisconnect)
|
||||
s.OnAllMembersLeft(ctx, roomCode)
|
||||
return
|
||||
}
|
||||
|
||||
newHostID := candidates[0]
|
||||
if err := s.participantDAO.TransferHost(ctx, room.ID, room.HostID, newHostID); err != nil {
|
||||
logs.Error(ctx, funcName, "host 自动转让 TransferHost 失败",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
return
|
||||
}
|
||||
if err := s.roomDAO.UpdateHost(ctx, room.ID, newHostID); err != nil {
|
||||
logs.Warn(ctx, funcName, "UpdateHost 失败(已发生参与者 role 变更)",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
}
|
||||
// 老 host 彻底离会(宽限期内未重连视为网络断开)
|
||||
_, _ = s.participantDAO.LeaveRoom(ctx, room.ID, room.HostID, constants.MeetingLeftReasonDisconnect)
|
||||
|
||||
s.broadcaster.BroadcastToMeeting(ctx, room.ID, constants.MeetingWSEventHostChanged, map[string]interface{}{
|
||||
"room_code": roomCode,
|
||||
"old_host_id": room.HostID,
|
||||
"new_host_id": newHostID,
|
||||
"auto_reason": "host_grace_expired",
|
||||
})
|
||||
logs.Info(ctx, funcName, "host 宽限期过期自动转让完成",
|
||||
zap.String("room_code", roomCode),
|
||||
zap.Int64("old_host_id", room.HostID),
|
||||
zap.Int64("new_host_id", newHostID))
|
||||
}
|
||||
|
||||
// OnAllMembersLeft 全员退出钩子
|
||||
// 设置 empty_ttl key + 启动本地 AfterFunc;期间有人重新加入可由 CancelEmptyTTL 复活房间
|
||||
func (s *MeetingLifecycleService) OnAllMembersLeft(ctx context.Context, roomCode string) {
|
||||
funcName := "service.meeting_lifecycle_service.OnAllMembersLeft"
|
||||
|
||||
ttl := time.Duration(s.cfg.EmptyRoomTTLSeconds) * time.Second
|
||||
// Redis TTL 同样加 buffer 避免与本地 timer 同时到期导致 DEL 无法命中
|
||||
redisTTL := ttl + s.ttlBuffer()
|
||||
if err := s.redis.Set(ctx, EmptyTTLKey(roomCode), "1", redisTTL).Err(); err != nil {
|
||||
logs.Warn(ctx, funcName, "写入 empty_ttl key 失败",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
return
|
||||
}
|
||||
s.scheduleEmptyTTLTimer(roomCode, ttl)
|
||||
logs.Info(ctx, funcName, "空房 TTL 已启动",
|
||||
zap.String("room_code", roomCode), zap.Int("ttl_seconds", s.cfg.EmptyRoomTTLSeconds))
|
||||
}
|
||||
|
||||
// CancelEmptyTTL 取消空房 TTL(有新成员加入时调用,房间复活)
|
||||
// 幂等:key 不存在也不报错
|
||||
func (s *MeetingLifecycleService) CancelEmptyTTL(ctx context.Context, roomCode string) {
|
||||
funcName := "service.meeting_lifecycle_service.CancelEmptyTTL"
|
||||
|
||||
deleted, err := s.redis.Del(ctx, EmptyTTLKey(roomCode)).Result()
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "清除 empty_ttl key 失败",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
}
|
||||
s.cancelTimer(&s.emptyTTLTimers, roomCode)
|
||||
if deleted > 0 {
|
||||
logs.Info(ctx, funcName, "空房 TTL 已取消(有人加入)",
|
||||
zap.String("room_code", roomCode))
|
||||
}
|
||||
}
|
||||
|
||||
// HandleEmptyRoomExpired 空房 TTL 过期处理
|
||||
// 流程:DEL empty_ttl key → 确认没被新成员清走 → MarkEnded(reason=empty_ttl) + CloseRouter 幂等
|
||||
func (s *MeetingLifecycleService) HandleEmptyRoomExpired(ctx context.Context, roomCode string) {
|
||||
funcName := "service.meeting_lifecycle_service.HandleEmptyRoomExpired"
|
||||
|
||||
deleted, err := s.redis.Del(ctx, EmptyTTLKey(roomCode)).Result()
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "DEL empty_ttl key 失败",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
return
|
||||
}
|
||||
s.cancelTimer(&s.emptyTTLTimers, roomCode)
|
||||
if deleted == 0 {
|
||||
logs.Debug(ctx, funcName, "empty_ttl key 已消失(有人重入 或 已处理过)",
|
||||
zap.String("room_code", roomCode))
|
||||
return
|
||||
}
|
||||
|
||||
room, err := s.roomDAO.GetByCode(ctx, roomCode)
|
||||
if err != nil || room == nil {
|
||||
logs.Warn(ctx, funcName, "加载会议失败或已不存在",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
return
|
||||
}
|
||||
if room.Status == constants.MeetingStatusEnded {
|
||||
return
|
||||
}
|
||||
|
||||
// 二次校验:DEL 成功与 DB 校验之间可能有新成员加入(并发场景)
|
||||
activeCount, _ := s.participantDAO.CountActiveByRoom(ctx, room.ID)
|
||||
if activeCount > 0 {
|
||||
logs.Info(ctx, funcName, "检测到活跃成员,放弃销毁",
|
||||
zap.String("room_code", roomCode), zap.Int64("active_count", activeCount))
|
||||
return
|
||||
}
|
||||
|
||||
if _, mErr := s.roomDAO.MarkEnded(ctx, room.ID, constants.MeetingEndedReasonEmptyTTL, time.Now()); mErr != nil {
|
||||
logs.Error(ctx, funcName, "MarkEnded 失败",
|
||||
zap.String("room_code", roomCode), zap.Error(mErr))
|
||||
return
|
||||
}
|
||||
if err := s.media.CloseRouter(ctx, roomCode); err != nil {
|
||||
logs.Warn(ctx, funcName, "CloseRouter 失败(不影响 DB 状态)",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
}
|
||||
|
||||
s.broadcaster.BroadcastToMeeting(ctx, room.ID, constants.MeetingWSEventRoomEnded, map[string]interface{}{
|
||||
"room_code": roomCode,
|
||||
"reason": constants.MeetingEndedReasonEmptyTTL,
|
||||
})
|
||||
logs.Info(ctx, funcName, "空房 TTL 过期,会议已销毁",
|
||||
zap.String("room_code", roomCode), zap.Int64("room_id", room.ID))
|
||||
}
|
||||
|
||||
// RescheduleFromRedis 服务启动时扫描 Redis 已存在的 grace / empty_ttl key,按剩余 PTTL 补装本地 timer
|
||||
// 服务重启后调用一次,避免重启期间 WS 断线事件的本地 timer 丢失导致宽限期/TTL 永远不触发
|
||||
func (s *MeetingLifecycleService) RescheduleFromRedis(ctx context.Context) {
|
||||
funcName := "service.meeting_lifecycle_service.RescheduleFromRedis"
|
||||
|
||||
s.rescheduleByPattern(ctx, funcName, redisKeyHostGracePrefix+"*", func(roomCode string, dur time.Duration) {
|
||||
s.scheduleHostGraceTimer(roomCode, dur)
|
||||
})
|
||||
s.rescheduleByPattern(ctx, funcName, redisKeyEmptyTTLPrefix+"*", func(roomCode string, dur time.Duration) {
|
||||
s.scheduleEmptyTTLTimer(roomCode, dur)
|
||||
})
|
||||
}
|
||||
|
||||
// ScanExpired 扫描指定 key pattern,返回 TTL <=0 (或即将过期)对应的 roomCode 列表
|
||||
// 供定时任务兜底调用:若本地 timer 已丢失,此时 Redis key 可能仍存在但 TTL 为负
|
||||
// 本实现采用 SCAN + PTTL,单次扫描 limit 控制内存
|
||||
func (s *MeetingLifecycleService) ScanExpired(ctx context.Context, keyPrefix string, limit int64) []string {
|
||||
funcName := "service.meeting_lifecycle_service.ScanExpired"
|
||||
|
||||
expired := make([]string, 0, limit)
|
||||
var cursor uint64
|
||||
for {
|
||||
keys, next, err := s.redis.Scan(ctx, cursor, keyPrefix+"*", limit).Result()
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "SCAN 失败", zap.String("pattern", keyPrefix), zap.Error(err))
|
||||
return expired
|
||||
}
|
||||
for _, k := range keys {
|
||||
pttl, err := s.redis.PTTL(ctx, k).Result()
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
// PTTL<=0(redis 过期延迟删除场景);PTTL==-2 已不存在
|
||||
if pttl > 0 {
|
||||
continue
|
||||
}
|
||||
roomCode := keyToRoomCode(k, keyPrefix)
|
||||
if roomCode != "" {
|
||||
expired = append(expired, roomCode)
|
||||
}
|
||||
}
|
||||
if next == 0 {
|
||||
break
|
||||
}
|
||||
cursor = next
|
||||
}
|
||||
return expired
|
||||
}
|
||||
|
||||
// HostGracePrefix / EmptyTTLPrefix 暴露 key 前缀供定时任务扫描使用
|
||||
func (s *MeetingLifecycleService) HostGracePrefix() string { return redisKeyHostGracePrefix }
|
||||
func (s *MeetingLifecycleService) EmptyTTLPrefix() string { return redisKeyEmptyTTLPrefix }
|
||||
|
||||
// ========== 内部辅助 ==========
|
||||
|
||||
func (s *MeetingLifecycleService) scheduleHostGraceTimer(roomCode string, d time.Duration) {
|
||||
s.cancelTimer(&s.graceTimers, roomCode)
|
||||
timer := time.AfterFunc(d, func() {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
s.HandleHostGraceExpired(ctx, roomCode)
|
||||
})
|
||||
s.graceTimers.Store(roomCode, timer)
|
||||
}
|
||||
|
||||
func (s *MeetingLifecycleService) scheduleEmptyTTLTimer(roomCode string, d time.Duration) {
|
||||
s.cancelTimer(&s.emptyTTLTimers, roomCode)
|
||||
timer := time.AfterFunc(d, func() {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
s.HandleEmptyRoomExpired(ctx, roomCode)
|
||||
})
|
||||
s.emptyTTLTimers.Store(roomCode, timer)
|
||||
}
|
||||
|
||||
// ttlBuffer 返回 Redis key 相对本地 timer 的额外 buffer 时长
|
||||
// 用于保证本地 timer 先触发(DEL 命中),避免 Redis 自动过期后 DEL 返回 0 被误判为重复触发
|
||||
// 规则:max(CleanupIntervalSeconds * 2, 30s),确保单机 cleanup 至少能扫到一轮兜底
|
||||
func (s *MeetingLifecycleService) ttlBuffer() time.Duration {
|
||||
buf := time.Duration(s.cfg.CleanupIntervalSeconds*2) * time.Second
|
||||
if buf < 30*time.Second {
|
||||
buf = 30 * time.Second
|
||||
}
|
||||
return buf
|
||||
}
|
||||
|
||||
// cancelTimer 取消本地 timer(幂等,不存在也安全)
|
||||
func (s *MeetingLifecycleService) cancelTimer(store *sync.Map, roomCode string) {
|
||||
if raw, ok := store.LoadAndDelete(roomCode); ok {
|
||||
if t, ok := raw.(*time.Timer); ok {
|
||||
t.Stop()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// rescheduleByPattern 通用"按 pattern 扫描 Redis + 按 PTTL 重建本地 timer"逻辑
|
||||
func (s *MeetingLifecycleService) rescheduleByPattern(
|
||||
ctx context.Context,
|
||||
funcName string,
|
||||
pattern string,
|
||||
schedule func(roomCode string, dur time.Duration),
|
||||
) {
|
||||
var cursor uint64
|
||||
total := 0
|
||||
for {
|
||||
keys, next, err := s.redis.Scan(ctx, cursor, pattern, 100).Result()
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "SCAN 失败", zap.String("pattern", pattern), zap.Error(err))
|
||||
return
|
||||
}
|
||||
for _, k := range keys {
|
||||
pttl, err := s.redis.PTTL(ctx, k).Result()
|
||||
if err != nil || pttl <= 0 {
|
||||
continue
|
||||
}
|
||||
prefix := pattern[:len(pattern)-1] // 去掉末尾 *
|
||||
roomCode := keyToRoomCode(k, prefix)
|
||||
if roomCode == "" {
|
||||
continue
|
||||
}
|
||||
schedule(roomCode, pttl)
|
||||
total++
|
||||
}
|
||||
if next == 0 {
|
||||
break
|
||||
}
|
||||
cursor = next
|
||||
}
|
||||
if total > 0 {
|
||||
logs.Info(ctx, funcName, "服务启动重建 timer 完成",
|
||||
zap.String("pattern", pattern), zap.Int("count", total))
|
||||
}
|
||||
}
|
||||
|
||||
// keyToRoomCode 从完整 Redis key 中去掉前缀得到 roomCode
|
||||
// 返回空字符串表示格式异常
|
||||
func keyToRoomCode(fullKey, prefix string) string {
|
||||
if len(fullKey) <= len(prefix) {
|
||||
return ""
|
||||
}
|
||||
return fullKey[len(prefix):]
|
||||
}
|
||||
|
||||
// formatErr 构造统一带上下文的错误(用于少数需要 return 的场景,避免裸字符串拼接)
|
||||
// nolint:unused
|
||||
func formatErr(funcName, msg string, err error) error {
|
||||
return fmt.Errorf("%s: %s: %w", funcName, msg, err)
|
||||
}
|
||||
@@ -65,11 +65,13 @@ type MeetingService struct {
|
||||
userResolver UserInfoResolver
|
||||
onlineChecker OnlineChecker
|
||||
mediaOrchestrator MediaOrchestrator
|
||||
lifecycleSvc *MeetingLifecycleService
|
||||
}
|
||||
|
||||
// NewMeetingService 创建 MeetingService 实例
|
||||
// 依赖通过构造函数注入;接口依赖由上游 Wire 绑定到具体实现
|
||||
// Task 7 (2026-04-21) 起 mediaOrchestrator 默认绑定 HTTPMediaOrchestrator
|
||||
// Task 8 (2026-04-21) 起新增 lifecycleSvc 依赖:JoinRoom/LeaveRoom 的空房/复活场景由 lifecycleSvc 托管
|
||||
func NewMeetingService(
|
||||
roomDAO *dao.MeetingRoomDAO,
|
||||
participantDAO *dao.MeetingParticipantDAO,
|
||||
@@ -81,6 +83,7 @@ func NewMeetingService(
|
||||
userResolver UserInfoResolver,
|
||||
onlineChecker OnlineChecker,
|
||||
mediaOrchestrator MediaOrchestrator,
|
||||
lifecycleSvc *MeetingLifecycleService,
|
||||
) *MeetingService {
|
||||
return &MeetingService{
|
||||
roomDAO: roomDAO,
|
||||
@@ -93,6 +96,7 @@ func NewMeetingService(
|
||||
userResolver: userResolver,
|
||||
onlineChecker: onlineChecker,
|
||||
mediaOrchestrator: mediaOrchestrator,
|
||||
lifecycleSvc: lifecycleSvc,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -210,9 +214,10 @@ func (s *MeetingService) CreateRoom(ctx context.Context, hostID int64, req *dto.
|
||||
return nil, nil, "", err
|
||||
}
|
||||
|
||||
// Task 7 起 CreateRouter 对接真实 mediasoup;Task 8 起仅在会议首次创建时调一次(房间级资源)
|
||||
routerID, mediaErr := s.mediaOrchestrator.CreateRouter(ctx, code)
|
||||
if mediaErr != nil {
|
||||
logs.Warn(ctx, funcName, "mediasoup Router 创建失败(Task 7 前为占位实现,不影响流程)",
|
||||
logs.Warn(ctx, funcName, "mediasoup Router 创建失败",
|
||||
zap.String("room_code", code), zap.Error(mediaErr))
|
||||
routerID = ""
|
||||
}
|
||||
@@ -309,10 +314,18 @@ func (s *MeetingService) JoinRoom(ctx context.Context, userID int64, code, passw
|
||||
return nil, nil, "", err
|
||||
}
|
||||
|
||||
routerID, mediaErr := s.mediaOrchestrator.CreateRouter(ctx, code)
|
||||
if mediaErr != nil {
|
||||
logs.Warn(ctx, funcName, "mediasoup Router 创建失败(占位实现)", zap.String("room_code", code), zap.Error(mediaErr))
|
||||
routerID = ""
|
||||
// Task 8:空房 TTL 复活
|
||||
// 若该房间正处于 empty_ttl 阶段(全员离开后的 5 分钟窗口),新成员加入立即取消销毁
|
||||
if s.lifecycleSvc != nil {
|
||||
s.lifecycleSvc.CancelEmptyTTL(ctx, code)
|
||||
}
|
||||
|
||||
// Task 8:JoinRoom 不再主动调 CreateRouter(Router 在 CreateRoom 时创建、由 HTTPMediaOrchestrator 本地缓存)
|
||||
// 从缓存读取 routerID;缺失时(极少见:服务重启后未重建缓存)保持为空,不阻塞加入流程
|
||||
routerID, _ := s.mediaOrchestrator.ResolveRouterID(code)
|
||||
if routerID == "" {
|
||||
logs.Debug(ctx, funcName, "RouterID 缓存缺失(非致命,可能服务重启)",
|
||||
zap.String("room_code", code))
|
||||
}
|
||||
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventMemberJoined, map[string]interface{}{
|
||||
@@ -378,9 +391,10 @@ func (s *MeetingService) LeaveRoom(ctx context.Context, userID int64, code strin
|
||||
}
|
||||
}
|
||||
|
||||
if len(actives) == 0 {
|
||||
_, _ = s.roomDAO.MarkEnded(ctx, room.ID, constants.MeetingEndedReasonEmptyTTL, time.Now())
|
||||
_ = s.mediaOrchestrator.CloseRouter(ctx, code)
|
||||
// Task 8:全员离会改走生命周期状态机(空房 TTL),不再立即销毁房间
|
||||
// 在 TTL 窗口内如有新成员加入,房间会被 CancelEmptyTTL 复活;TTL 到期由 HandleEmptyRoomExpired 兜底销毁
|
||||
if len(actives) == 0 && s.lifecycleSvc != nil {
|
||||
s.lifecycleSvc.OnAllMembersLeft(ctx, code)
|
||||
}
|
||||
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{
|
||||
|
||||
@@ -31,15 +31,18 @@ type MeetingSignalService struct {
|
||||
|
||||
broadcaster *MeetingBroadcaster
|
||||
mediaOrchestrator MediaOrchestrator
|
||||
lifecycleSvc *MeetingLifecycleService
|
||||
}
|
||||
|
||||
// NewMeetingSignalService 构造 WS 信令服务
|
||||
// Task 8 注入 lifecycleSvc:host 掉线 / 重连的生命周期联动
|
||||
func NewMeetingSignalService(
|
||||
roomDAO *dao.MeetingRoomDAO,
|
||||
participantDAO *dao.MeetingParticipantDAO,
|
||||
redis *redis.Client,
|
||||
broadcaster *MeetingBroadcaster,
|
||||
mediaOrchestrator MediaOrchestrator,
|
||||
lifecycleSvc *MeetingLifecycleService,
|
||||
) *MeetingSignalService {
|
||||
return &MeetingSignalService{
|
||||
roomDAO: roomDAO,
|
||||
@@ -47,6 +50,7 @@ func NewMeetingSignalService(
|
||||
redis: redis,
|
||||
broadcaster: broadcaster,
|
||||
mediaOrchestrator: mediaOrchestrator,
|
||||
lifecycleSvc: lifecycleSvc,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -116,6 +120,12 @@ func (s *MeetingSignalService) OnRoomJoin(ctx context.Context, userID int64, roo
|
||||
key := resourceTrackKey(roomCode, userID)
|
||||
_ = s.redis.Expire(ctx, key, resourceTTL).Err()
|
||||
|
||||
// Task 8:若本次 WS 在线的用户是会议 host,尝试撤销 host 宽限期
|
||||
// 本方法幂等:若此前无 host_grace key(正常场景)DEL 返回 0,不产生副作用
|
||||
if s.lifecycleSvc != nil && room.HostID == userID {
|
||||
s.lifecycleSvc.OnHostReconnect(ctx, roomCode, userID)
|
||||
}
|
||||
|
||||
logs.Info(ctx, "service.meeting_signal_service.OnRoomJoin", "用户宣告 WS 在线",
|
||||
zap.String("room_code", roomCode),
|
||||
zap.Int64("user_id", userID),
|
||||
@@ -123,6 +133,47 @@ func (s *MeetingSignalService) OnRoomJoin(ctx context.Context, userID int64, roo
|
||||
return nil
|
||||
}
|
||||
|
||||
// OnWSDisconnect WS 断线钩子,实现 ws.MeetingDisconnectHook 接口(Task 8)
|
||||
// 由 ws.handler 的 SetOnDisconnect 回调触发;场景:用户的最后一条 WS 连接被移除
|
||||
// 职责:
|
||||
// 1. 若该用户当前存在活跃 meeting_participants 记录:清理其媒体资源 + 若是 host 启动 host 宽限期
|
||||
// 2. 普通成员:仅清媒体资源(设计 Q1=A:不动 meeting_participants,长期不活跃由 4h 兜底清理)
|
||||
//
|
||||
// 容错:所有异常仅 Warn 日志,不返回 error;保证 ws.handler 主断连流程不受影响
|
||||
func (s *MeetingSignalService) OnWSDisconnect(ctx context.Context, userID int64) {
|
||||
funcName := "service.meeting_signal_service.OnWSDisconnect"
|
||||
|
||||
active, err := s.participantDAO.FindActiveByUser(ctx, userID)
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "查询活跃会议失败", zap.Int64("user_id", userID), zap.Error(err))
|
||||
return
|
||||
}
|
||||
if active == nil {
|
||||
// 用户当前不在任何会议
|
||||
return
|
||||
}
|
||||
room, err := s.roomDAO.GetByID(ctx, active.RoomID)
|
||||
if err != nil || room == nil {
|
||||
logs.Warn(ctx, funcName, "加载 room 失败或已不存在",
|
||||
zap.Int64("room_id", active.RoomID), zap.Error(err))
|
||||
return
|
||||
}
|
||||
if room.Status == constants.MeetingStatusEnded {
|
||||
return
|
||||
}
|
||||
|
||||
s.cleanupUserResources(ctx, room.RoomCode, userID)
|
||||
|
||||
if s.lifecycleSvc != nil && room.HostID == userID {
|
||||
s.lifecycleSvc.OnHostDisconnect(ctx, room.RoomCode, userID)
|
||||
}
|
||||
|
||||
logs.Info(ctx, funcName, "WS 断线已处理",
|
||||
zap.String("room_code", room.RoomCode),
|
||||
zap.Int64("user_id", userID),
|
||||
zap.Bool("is_host", room.HostID == userID))
|
||||
}
|
||||
|
||||
// OnRoomLeave 处理 meeting.room.leave 事件
|
||||
// 语义:WS 层面的主动离会(等价 REST leave 但不强制要求落库事务;
|
||||
// 当前实现:仅清理该用户在本会议的所有媒体资源(batch close producer/consumer)+ 广播 meeting.member.left
|
||||
|
||||
191
backend/go-service/app/meeting/task/meeting_cleanup_task.go
Normal file
191
backend/go-service/app/meeting/task/meeting_cleanup_task.go
Normal file
@@ -0,0 +1,191 @@
|
||||
// Package task 提供 meeting 模块的后台定时任务
|
||||
package task
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"github.com/echochat/backend/app/constants"
|
||||
"github.com/echochat/backend/app/meeting/dao"
|
||||
"github.com/echochat/backend/app/meeting/service"
|
||||
"github.com/echochat/backend/pkg/logs"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
// MeetingCleanupTask 会议生命周期兜底清理任务(Phase 2e-2 Task 8)
|
||||
// 职责(设计 §11.4 兜底策略 + Task 8 落地):
|
||||
// 1. 启动时一次性 RescheduleFromRedis,按 Redis 残留 grace/empty_ttl key 重建本地 timer
|
||||
// 2. 每 cleanupInterval 秒扫描 host_grace / empty_ttl,命中"Redis 仍在但本地 timer 可能已丢失"的场景
|
||||
// 兜底触发 HandleHostGraceExpired / HandleEmptyRoomExpired;真正过期判断由 Lifecycle 内部 DEL 互斥
|
||||
// 3. 每 10 分钟扫一次 stale active rooms(活跃状态但超过 StaleRoomHours 无成员动作)→ 强制 MarkEnded(system_error)
|
||||
// 4. 每 24 小时扫 Ended 超过 chat 保留时长的房间 → DeleteByRoomIDs 清理 meeting_chats
|
||||
//
|
||||
// 所有周期都由独立 ticker 驱动,相互不阻塞;退出走 stopCh 统一关闭
|
||||
type MeetingCleanupTask struct {
|
||||
lifecycleSvc *service.MeetingLifecycleService
|
||||
roomDAO *dao.MeetingRoomDAO
|
||||
chatDAO *dao.MeetingChatDAO
|
||||
|
||||
cleanupInterval time.Duration // 扫 grace/empty_ttl 的周期
|
||||
staleInterval time.Duration // 扫 stale active 的周期
|
||||
chatInterval time.Duration // 扫 chat 清理的周期
|
||||
staleHours int // 超过此小时的活跃房间视为 stale
|
||||
|
||||
stopCh chan struct{}
|
||||
}
|
||||
|
||||
// NewMeetingCleanupTask 构造定时任务
|
||||
// 扫描参数从 lifecycleSvc.Config() 读取,避免两处配置分离
|
||||
func NewMeetingCleanupTask(
|
||||
lifecycleSvc *service.MeetingLifecycleService,
|
||||
roomDAO *dao.MeetingRoomDAO,
|
||||
chatDAO *dao.MeetingChatDAO,
|
||||
) *MeetingCleanupTask {
|
||||
cfg := lifecycleSvc.Config()
|
||||
interval := time.Duration(cfg.CleanupIntervalSeconds) * time.Second
|
||||
if interval <= 0 {
|
||||
interval = 30 * time.Second
|
||||
}
|
||||
staleHours := cfg.StaleRoomHours
|
||||
if staleHours <= 0 {
|
||||
staleHours = 4
|
||||
}
|
||||
return &MeetingCleanupTask{
|
||||
lifecycleSvc: lifecycleSvc,
|
||||
roomDAO: roomDAO,
|
||||
chatDAO: chatDAO,
|
||||
cleanupInterval: interval,
|
||||
staleInterval: 10 * time.Minute,
|
||||
chatInterval: 24 * time.Hour,
|
||||
staleHours: staleHours,
|
||||
stopCh: make(chan struct{}),
|
||||
}
|
||||
}
|
||||
|
||||
// Start 启动任务(非阻塞)
|
||||
// 应在 main 进程启动后调用,进程退出时调用 Stop
|
||||
func (t *MeetingCleanupTask) Start() {
|
||||
funcName := "task.meeting_cleanup_task.Start"
|
||||
logs.Info(nil, funcName, "启动会议生命周期兜底任务",
|
||||
zap.Duration("cleanup_interval", t.cleanupInterval),
|
||||
zap.Duration("stale_interval", t.staleInterval),
|
||||
zap.Duration("chat_interval", t.chatInterval),
|
||||
zap.Int("stale_hours", t.staleHours))
|
||||
|
||||
go t.runRescheduleOnce()
|
||||
go t.runLoop("grace_empty_ttl", t.cleanupInterval, 5*time.Second, t.scanExpiredKeys)
|
||||
go t.runLoop("stale_active", t.staleInterval, 30*time.Second, t.scanStaleActive)
|
||||
go t.runLoop("chat_cleanup", t.chatInterval, 60*time.Second, t.cleanupExpiredChats)
|
||||
}
|
||||
|
||||
// Stop 停止所有循环
|
||||
func (t *MeetingCleanupTask) Stop() {
|
||||
close(t.stopCh)
|
||||
}
|
||||
|
||||
// runRescheduleOnce 启动后一次性扫描 Redis 残留 key 并补装本地 timer
|
||||
func (t *MeetingCleanupTask) runRescheduleOnce() {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
t.lifecycleSvc.RescheduleFromRedis(ctx)
|
||||
}
|
||||
|
||||
// runLoop 通用"固定间隔 + 预热延迟 + stopCh 退出"循环
|
||||
// warmup 用于错峰避免启动风暴,单次业务执行由 do 函数自行控制超时
|
||||
func (t *MeetingCleanupTask) runLoop(name string, interval, warmup time.Duration, do func()) {
|
||||
funcName := "task.meeting_cleanup_task.runLoop." + name
|
||||
|
||||
timer := time.NewTimer(warmup)
|
||||
select {
|
||||
case <-t.stopCh:
|
||||
timer.Stop()
|
||||
return
|
||||
case <-timer.C:
|
||||
}
|
||||
|
||||
do()
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-t.stopCh:
|
||||
logs.Info(nil, funcName, "任务循环退出")
|
||||
return
|
||||
case <-ticker.C:
|
||||
do()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// scanExpiredKeys 扫 host_grace 与 empty_ttl 两个 key 集合
|
||||
// Redis 原生 TTL 过期时 Go 的 PTTL 会返回 -2(不存在),此时无需兜底
|
||||
// 少数场景(AfterFunc 丢失 / 多实例部署)会出现 Redis 仍有 key 但无本地 timer
|
||||
// 本扫描对 "TTL≤0 但仍存在" 的 key 主动调 HandleXxxExpired(内部 DEL 互斥)
|
||||
func (t *MeetingCleanupTask) scanExpiredKeys() {
|
||||
funcName := "task.meeting_cleanup_task.scanExpiredKeys"
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
|
||||
hostExpired := t.lifecycleSvc.ScanExpired(ctx, t.lifecycleSvc.HostGracePrefix(), 100)
|
||||
for _, code := range hostExpired {
|
||||
t.lifecycleSvc.HandleHostGraceExpired(ctx, code)
|
||||
}
|
||||
emptyExpired := t.lifecycleSvc.ScanExpired(ctx, t.lifecycleSvc.EmptyTTLPrefix(), 100)
|
||||
for _, code := range emptyExpired {
|
||||
t.lifecycleSvc.HandleEmptyRoomExpired(ctx, code)
|
||||
}
|
||||
if len(hostExpired)+len(emptyExpired) > 0 {
|
||||
logs.Info(ctx, funcName, "兜底处理过期 key 完成",
|
||||
zap.Int("host_grace", len(hostExpired)),
|
||||
zap.Int("empty_ttl", len(emptyExpired)))
|
||||
}
|
||||
}
|
||||
|
||||
// scanStaleActive 扫超 staleHours 无任何成员动作的活跃会议,强制结束
|
||||
// reason=system_error:明确标记非正常终止,便于运维排查 Node worker 掉链等异常
|
||||
func (t *MeetingCleanupTask) scanStaleActive() {
|
||||
funcName := "task.meeting_cleanup_task.scanStaleActive"
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
|
||||
defer cancel()
|
||||
|
||||
ids, err := t.roomDAO.ListStaleActive(ctx, t.staleHours, 100)
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "扫描 stale active 失败", zap.Error(err))
|
||||
return
|
||||
}
|
||||
if len(ids) == 0 {
|
||||
return
|
||||
}
|
||||
now := time.Now()
|
||||
for _, id := range ids {
|
||||
if _, err := t.roomDAO.MarkEnded(ctx, id, constants.MeetingEndedReasonSystemError, now); err != nil {
|
||||
logs.Warn(ctx, funcName, "兜底 MarkEnded 失败",
|
||||
zap.Int64("room_id", id), zap.Error(err))
|
||||
}
|
||||
}
|
||||
logs.Info(ctx, funcName, "stale 会议兜底结束", zap.Int("count", len(ids)))
|
||||
}
|
||||
|
||||
// cleanupExpiredChats 扫超过 MeetingChatRetentionHours 的 Ended 会议,批量删除聊天记录
|
||||
func (t *MeetingCleanupTask) cleanupExpiredChats() {
|
||||
funcName := "task.meeting_cleanup_task.cleanupExpiredChats"
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)
|
||||
defer cancel()
|
||||
|
||||
ids, err := t.roomDAO.ListExpiredForCleanup(ctx, constants.MeetingChatRetentionHours, 500)
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "扫描过期会议失败", zap.Error(err))
|
||||
return
|
||||
}
|
||||
if len(ids) == 0 {
|
||||
return
|
||||
}
|
||||
deleted, err := t.chatDAO.DeleteByRoomIDs(ctx, ids)
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "批量删除聊天记录失败", zap.Error(err))
|
||||
return
|
||||
}
|
||||
logs.Info(ctx, funcName, "会议聊天清理完成",
|
||||
zap.Int("room_count", len(ids)),
|
||||
zap.Int64("chat_deleted", deleted))
|
||||
}
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
imHandler "github.com/echochat/backend/app/im/handler"
|
||||
meetingController "github.com/echochat/backend/app/meeting/controller"
|
||||
meetingService "github.com/echochat/backend/app/meeting/service"
|
||||
meetingTask "github.com/echochat/backend/app/meeting/task"
|
||||
notifyController "github.com/echochat/backend/app/notify/controller"
|
||||
notifyService "github.com/echochat/backend/app/notify/service"
|
||||
notifyTask "github.com/echochat/backend/app/notify/task"
|
||||
@@ -56,8 +57,10 @@ type App struct {
|
||||
NotifyCleanupTask *notifyTask.CleanupTask // 通知清理定时任务
|
||||
MeetingService *meetingService.MeetingService // 会议 REST 业务服务(Task 5 落地)
|
||||
MeetingSignalService *meetingService.MeetingSignalService // 会议 WS 信令业务服务(Task 6 落地)
|
||||
MeetingLifecycleSvc *meetingService.MeetingLifecycleService // 会议生命周期状态机(Task 8 落地)
|
||||
MeetingController *meetingController.MeetingController // 会议 REST 控制器
|
||||
MeetingWSHandler *meetingController.MeetingWSHandler // 会议 WS 事件 Handler(构造时自动注册路由到 Hub)
|
||||
MeetingCleanupTask *meetingTask.MeetingCleanupTask // 会议生命周期兜底定时任务(Task 8 落地)
|
||||
}
|
||||
|
||||
// NewApp 创建应用实例
|
||||
@@ -89,11 +92,14 @@ func NewApp(
|
||||
notifyCleanup *notifyTask.CleanupTask,
|
||||
meetingSvc *meetingService.MeetingService,
|
||||
meetingSignalSvc *meetingService.MeetingSignalService,
|
||||
meetingLifecycleSvc *meetingService.MeetingLifecycleService,
|
||||
meetingCtrl *meetingController.MeetingController,
|
||||
meetingWSHandler *meetingController.MeetingWSHandler,
|
||||
meetingCleanup *meetingTask.MeetingCleanupTask,
|
||||
) *App {
|
||||
wsHandler.SetOfflinePusher(offlinePusher)
|
||||
wsHandler.SetNotifyConnectHook(notifySvc)
|
||||
wsHandler.SetMeetingDisconnectHook(meetingSignalSvc)
|
||||
|
||||
return &App{
|
||||
Config: cfg,
|
||||
@@ -123,8 +129,10 @@ func NewApp(
|
||||
NotifyCleanupTask: notifyCleanup,
|
||||
MeetingService: meetingSvc,
|
||||
MeetingSignalService: meetingSignalSvc,
|
||||
MeetingLifecycleSvc: meetingLifecycleSvc,
|
||||
MeetingController: meetingCtrl,
|
||||
MeetingWSHandler: meetingWSHandler,
|
||||
MeetingCleanupTask: meetingCleanup,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -54,6 +54,8 @@ func InitializeApp(cfg *config.Config) (*App, error) {
|
||||
wire.Bind(new(meetingService.NotifyPusher), new(*notifyService.NotifyService)),
|
||||
wire.Bind(new(meetingService.UserInfoResolver), new(*contactDAO.FriendshipDAO)),
|
||||
wire.Bind(new(meetingService.OnlineChecker), new(*wsApp.OnlineService)),
|
||||
// Task 8:WS 断线钩子 → MeetingSignalService
|
||||
wire.Bind(new(wsApp.MeetingDisconnectHook), new(*meetingService.MeetingSignalService)),
|
||||
)
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
@@ -25,6 +25,7 @@ import (
|
||||
controller7 "github.com/echochat/backend/app/meeting/controller"
|
||||
dao6 "github.com/echochat/backend/app/meeting/dao"
|
||||
service7 "github.com/echochat/backend/app/meeting/service"
|
||||
task2 "github.com/echochat/backend/app/meeting/task"
|
||||
controller6 "github.com/echochat/backend/app/notify/controller"
|
||||
dao5 "github.com/echochat/backend/app/notify/dao"
|
||||
service3 "github.com/echochat/backend/app/notify/service"
|
||||
@@ -103,10 +104,12 @@ func InitializeApp(cfg *config.Config) (*App, error) {
|
||||
meetingChatDAO := dao6.NewMeetingChatDAO(gormDB)
|
||||
meetingBroadcaster := service7.NewMeetingBroadcaster(meetingParticipantDAO, pubSub)
|
||||
httpMediaOrchestrator := service7.NewHTTPMediaOrchestrator(cfg)
|
||||
meetingService := service7.NewMeetingService(meetingRoomDAO, meetingParticipantDAO, meetingChatDAO, gormDB, client, meetingBroadcaster, notifyService, friendshipDAO, onlineService, httpMediaOrchestrator)
|
||||
meetingSignalService := service7.NewMeetingSignalService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, httpMediaOrchestrator)
|
||||
meetingLifecycleService := service7.NewMeetingLifecycleService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, httpMediaOrchestrator, cfg)
|
||||
meetingService := service7.NewMeetingService(meetingRoomDAO, meetingParticipantDAO, meetingChatDAO, gormDB, client, meetingBroadcaster, notifyService, friendshipDAO, onlineService, httpMediaOrchestrator, meetingLifecycleService)
|
||||
meetingSignalService := service7.NewMeetingSignalService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, httpMediaOrchestrator, meetingLifecycleService)
|
||||
meetingController := controller7.NewMeetingController(meetingService)
|
||||
meetingWSHandler := controller7.NewMeetingWSHandler(meetingSignalService, hub)
|
||||
app := NewApp(cfg, gormDB, client, minioClient, authService, authController, adminAuthController, userManageController, onlineController, contactManageController, groupManageController, messageManageController, handler, hub, pubSub, onlineService, contactController, imController, eventHandler, offlinePusher, fileController, groupController, notifyService, notificationController, cleanupTask, meetingService, meetingSignalService, meetingController, meetingWSHandler)
|
||||
meetingCleanupTask := task2.NewMeetingCleanupTask(meetingLifecycleService, meetingRoomDAO, meetingChatDAO)
|
||||
app := NewApp(cfg, gormDB, client, minioClient, authService, authController, adminAuthController, userManageController, onlineController, contactManageController, groupManageController, messageManageController, handler, hub, pubSub, onlineService, contactController, imController, eventHandler, offlinePusher, fileController, groupController, notifyService, notificationController, cleanupTask, meetingService, meetingSignalService, meetingLifecycleService, meetingController, meetingWSHandler, meetingCleanupTask)
|
||||
return app, nil
|
||||
}
|
||||
|
||||
@@ -35,6 +35,15 @@ type NotifyConnectHook interface {
|
||||
PushUnreadTotalOnConnect(ctx context.Context, userID int64)
|
||||
}
|
||||
|
||||
// MeetingDisconnectHook 会议 WS 断线钩子接口(Phase 2e-2 Task 8)
|
||||
// 由 meeting.service.MeetingSignalService 隐式实现
|
||||
// 触发时机:用户最后一条 WS 连接被移除(最终下线)
|
||||
// 职责:清理该用户在会议中的媒体资源;若为 host 则启动 host 宽限期
|
||||
// 抽象为接口避免 ws 包反向依赖 meeting 包引发循环引用
|
||||
type MeetingDisconnectHook interface {
|
||||
OnWSDisconnect(ctx context.Context, userID int64)
|
||||
}
|
||||
|
||||
var upgrader = websocket.Upgrader{
|
||||
ReadBufferSize: 1024,
|
||||
WriteBufferSize: 1024,
|
||||
@@ -45,13 +54,14 @@ var upgrader = websocket.Upgrader{
|
||||
|
||||
// Handler WebSocket 连接处理器
|
||||
type Handler struct {
|
||||
hub *ws.Hub
|
||||
pubsub *ws.PubSub
|
||||
jwtCfg *config.JWTConfig
|
||||
onlineService *OnlineService
|
||||
tokenValidator TokenValidator
|
||||
offlinePusher OfflineMessagePusher
|
||||
notifyConnectHook NotifyConnectHook
|
||||
hub *ws.Hub
|
||||
pubsub *ws.PubSub
|
||||
jwtCfg *config.JWTConfig
|
||||
onlineService *OnlineService
|
||||
tokenValidator TokenValidator
|
||||
offlinePusher OfflineMessagePusher
|
||||
notifyConnectHook NotifyConnectHook
|
||||
meetingDisconnectHook MeetingDisconnectHook // Task 8 注入
|
||||
}
|
||||
|
||||
// NewHandler 创建 WebSocket Handler 实例
|
||||
@@ -75,6 +85,12 @@ func (h *Handler) SetNotifyConnectHook(hook NotifyConnectHook) {
|
||||
h.notifyConnectHook = hook
|
||||
}
|
||||
|
||||
// SetMeetingDisconnectHook 设置会议 WS 断线钩子(由 meeting 模块在初始化时注入)
|
||||
// Phase 2e-2 Task 8:WS 最终下线时触发 host 宽限期 / 媒体资源清理
|
||||
func (h *Handler) SetMeetingDisconnectHook(hook MeetingDisconnectHook) {
|
||||
h.meetingDisconnectHook = hook
|
||||
}
|
||||
|
||||
// Upgrade 处理 WebSocket 升级请求
|
||||
// GET /ws?token=xxx → JWT 认证 → 升级连接 → 注册 Hub → 订阅 Redis 频道
|
||||
func (h *Handler) Upgrade(c *gin.Context) {
|
||||
@@ -120,6 +136,11 @@ func (h *Handler) Upgrade(c *gin.Context) {
|
||||
}
|
||||
h.pubsub.Unsubscribe(userID)
|
||||
h.onlineService.UserOffline(context.Background(), userID)
|
||||
// Task 8:通知会议模块处理 host 宽限期 / 媒体资源清理
|
||||
// 钩子为可选(测试或 meeting 模块未加载场景安全跳过)
|
||||
if h.meetingDisconnectHook != nil {
|
||||
go h.meetingDisconnectHook.OnWSDisconnect(context.Background(), userID)
|
||||
}
|
||||
})
|
||||
h.hub.Register(client)
|
||||
h.pubsub.Subscribe(claims.UserID)
|
||||
|
||||
@@ -82,6 +82,11 @@ func main() {
|
||||
app.NotifyCleanupTask.Start()
|
||||
defer app.NotifyCleanupTask.Stop()
|
||||
|
||||
// 7.2 启动 meeting 模块生命周期兜底任务(Phase 2e-2 Task 8)
|
||||
// 内部会先 RescheduleFromRedis 恢复服务重启前残留的 host_grace / empty_ttl timer
|
||||
app.MeetingCleanupTask.Start()
|
||||
defer app.MeetingCleanupTask.Stop()
|
||||
|
||||
// 8. 启动 HTTP 服务(优雅关闭)
|
||||
addr := fmt.Sprintf(":%d", cfg.Server.Port)
|
||||
srv := &http.Server{
|
||||
|
||||
@@ -63,3 +63,11 @@ media_server:
|
||||
timeout_ms: 5000 # 创建类接口超时(毫秒)
|
||||
close_timeout_ms: 2000 # 关闭类接口超时(毫秒)
|
||||
close_retry: 2 # 关闭类接口失败重试次数(指数退避 200/500ms)
|
||||
|
||||
# 会议模块生命周期参数(Phase 2e-2 Task 8 会议状态机)
|
||||
# E2E 测试可通过环境变量 ECHOCHAT_MEETING_HOST_GRACE_SECONDS 等覆盖为更短值以加速触发
|
||||
meeting:
|
||||
host_grace_seconds: 120 # 主持人掉线宽限期(秒),期间保留 host 身份
|
||||
empty_room_ttl_seconds: 300 # 空房自动销毁 TTL(秒),期间有人加入可复活
|
||||
cleanup_interval_seconds: 30 # 兜底扫描周期(秒),定时任务检查过期 grace/empty_ttl key
|
||||
stale_room_hours: 4 # 活跃超过此小时且无活跃成员视为 stale,强制结束
|
||||
|
||||
@@ -55,3 +55,10 @@ media_server:
|
||||
timeout_ms: 5000
|
||||
close_timeout_ms: 2000
|
||||
close_retry: 2
|
||||
|
||||
# 会议模块生命周期参数(Phase 2e-2 Task 8)
|
||||
meeting:
|
||||
host_grace_seconds: 120
|
||||
empty_room_ttl_seconds: 300
|
||||
cleanup_interval_seconds: 30
|
||||
stale_room_hours: 4
|
||||
|
||||
@@ -18,6 +18,16 @@ type Config struct {
|
||||
Log LogConfig `mapstructure:"log"`
|
||||
Minio MinioConfig `mapstructure:"minio"`
|
||||
MediaServer MediaServerConfig `mapstructure:"media_server"`
|
||||
Meeting MeetingConfig `mapstructure:"meeting"`
|
||||
}
|
||||
|
||||
// MeetingConfig 会议模块生命周期参数(Phase 2e-2 Task 8)
|
||||
// 所有字段均允许通过环境变量 ECHOCHAT_MEETING_* 覆盖,便于 E2E 测试以短时长加速场景触发
|
||||
type MeetingConfig struct {
|
||||
HostGraceSeconds int `mapstructure:"host_grace_seconds"` // host 掉线宽限期秒数,默认 120
|
||||
EmptyRoomTTLSeconds int `mapstructure:"empty_room_ttl_seconds"` // 空房自动销毁 TTL 秒数,默认 300
|
||||
CleanupIntervalSeconds int `mapstructure:"cleanup_interval_seconds"` // 兜底扫描周期秒数,默认 30
|
||||
StaleRoomHours int `mapstructure:"stale_room_hours"` // 活跃超过此小时且无成员视为 stale,默认 4
|
||||
}
|
||||
|
||||
// MediaServerConfig Node media-server 接入配置(Phase 2e-2 Task 7)
|
||||
|
||||
Reference in New Issue
Block a user