Files
EchoChat/backend/go-service/app/meeting/service/meeting_signal_service.go
bujinyuan f5ae095033 fix(meeting): Task 16 会议销毁路径资源清理专项
前端修复:
- _onRoomEnded 保留结束页数据,新增 exitEndedRoom() 动作
  room.vue onUnload 在 ENDED 状态调用,释放 currentRoom / participants / chatMessages 等 pinia state
- endMeeting() API 成功后立即触发 _onRoomEnded 本地兜底,不再依赖 room.ended 广播
  避免 API 成功但 WS 抖动导致 engine / timer / AudioContext 泄漏
- 新增 _pendingBroadcastTimers 集合收纳 _broadcastSelfState 的 setTimeout 句柄
  _cleanupMedia 统一 clearTimeout,避免会议结束后 stale 重试逻辑排队

后端修复:
- MeetingLifecycleService.OnRoomEnded 统一生命周期收尾
  取消 graceTimers / emptyTTLTimers + Pipeline DEL host_grace / empty_ttl / 两个 handling 锁 key
- EndRoom 事务提交后遍历 activesBefore 调 cleanupRoomRedisResidual
  批量 DEL 每个前任活跃成员的 resourceTrackKey + memberStateKey
  补调 lifecycleSvc.OnRoomEnded 清理自身 timer + keys
- HandleEmptyRoomExpired MarkEnded 后补调 cleanupRoomRedisResidual + 清 host_grace 残留

验证:go vet + go build + npm run build:h5 全部通过

Made-with: Cursor
2026-04-23 17:12:38 +08:00

736 lines
29 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 service 提供 meeting 模块的业务逻辑
package service
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/echochat/backend/app/constants"
"github.com/echochat/backend/app/meeting/dao"
"github.com/echochat/backend/app/meeting/model"
"github.com/echochat/backend/pkg/logs"
"github.com/redis/go-redis/v9"
"go.uber.org/zap"
)
// MeetingSignalService WS 信令事件业务处理
// Task 6 落地:处理设计 §6.3 的 11 个 meeting.* 事件,分 3 组:
// - 房间组meeting.room.join/leave绑定 WS 连接 ↔ roomCode辅助断连清理
// - 成员组meeting.member.state.changed麦/视频状态变更 + host 静音他人
// - 媒体组meeting.transport.*/produce.*/consume.*/producer.close
// 对接 MediaOrchestrator实现 mediasoup signaling 桥接
//
// 所有方法统一返回 (ackData, error)error 非 nil 表示业务失败,由 Handler 映射为 ACK code=-1
// 广播副作用(如 meeting.member.producer.new在方法内部通过 broadcaster 发出,不进 ACK
type MeetingSignalService struct {
roomDAO *dao.MeetingRoomDAO
participantDAO *dao.MeetingParticipantDAO
redis *redis.Client
broadcaster *MeetingBroadcaster
mediaOrchestrator MediaOrchestrator
lifecycleSvc *MeetingLifecycleService
}
// NewMeetingSignalService 构造 WS 信令服务
// Task 8 注入 lifecycleSvchost 掉线 / 重连的生命周期联动
func NewMeetingSignalService(
roomDAO *dao.MeetingRoomDAO,
participantDAO *dao.MeetingParticipantDAO,
redis *redis.Client,
broadcaster *MeetingBroadcaster,
mediaOrchestrator MediaOrchestrator,
lifecycleSvc *MeetingLifecycleService,
) *MeetingSignalService {
return &MeetingSignalService{
roomDAO: roomDAO,
participantDAO: participantDAO,
redis: redis,
broadcaster: broadcaster,
mediaOrchestrator: mediaOrchestrator,
lifecycleSvc: lifecycleSvc,
}
}
// Redis 资源追踪 key设计 §九 - 断线清理)
// Set 内元素格式:"transport:{id}" / "producer:{id}" / "consumer:{id}"
func resourceTrackKey(roomCode string, userID int64) string {
return fmt.Sprintf("echo:meeting:resource:%s:%d", roomCode, userID)
}
// memberStateKey Redis Hash保存用户当前的音视频开关状态
// fieldsaudio_enabled / video_enabled值为 "true" / "false"
// 用途:后入者加入房间时,后端从这里读取每个已有成员的最新状态定向补推 state.changed
// 解决"后入者的成员列表里其他人图标一直灰色"的问题
func memberStateKey(roomCode string, userID int64) string {
return fmt.Sprintf("echo:meeting:member_state:%s:%d", roomCode, userID)
}
// resourceTTL 单个用户资源追踪集合 TTL
// 设计:会议期间维持可达即可;若用户长期不活跃由断线清理接管
const resourceTTL = time.Hour
// trackResource 记录用户在会议中持有的媒体资源 ID
func (s *MeetingSignalService) trackResource(ctx context.Context, roomCode string, userID int64, kind, id string) {
key := resourceTrackKey(roomCode, userID)
member := kind + ":" + id
if err := s.redis.SAdd(ctx, key, member).Err(); err != nil {
logs.Warn(ctx, "service.meeting_signal_service.trackResource", "追踪媒体资源失败",
zap.String("key", key), zap.String("member", member), zap.Error(err))
return
}
_ = s.redis.Expire(ctx, key, resourceTTL).Err()
}
// untrackResource 从集合中移除资源 ID关闭 producer/consumer 时)
func (s *MeetingSignalService) untrackResource(ctx context.Context, roomCode string, userID int64, kind, id string) {
key := resourceTrackKey(roomCode, userID)
member := kind + ":" + id
_ = s.redis.SRem(ctx, key, member).Err()
}
// assertOwnsResource 校验指定资源kind: transport/producer/consumer是否由 userID 在该 roomCode 中创建
// 背景WS 信令入口原本只校验"用户是否在会议中",但客户端上报的 transport_id / producer_id / consumer_id 完全可伪造,
// 造成任意成员可关闭他人 producer、把 consumer 挂到他人 recv transport 等横向越权P0-1 / P0-2 审计)
// 判据trackResource 记录的 Redis Set 是最权威的归属来源kind:id 写入 / untrack 时移除)
// 返回:
// - nil归属合法
// - ErrResourceNotOwnedSet 里不存在该元素(越权尝试)
// - Redis 错误时同样返回 ErrResourceNotOwned 并记录 Warn按"fail-closed"保守拒绝,避免服务抖动打开权限口子
func (s *MeetingSignalService) assertOwnsResource(ctx context.Context, roomCode string, userID int64, kind, id string) error {
if id == "" {
return fmt.Errorf("%s_id 不能为空", kind)
}
key := resourceTrackKey(roomCode, userID)
member := kind + ":" + id
ok, err := s.redis.SIsMember(ctx, key, member).Result()
if err != nil {
logs.Warn(ctx, "service.meeting_signal_service.assertOwnsResource", "查询资源归属失败",
zap.String("key", key), zap.String("member", member), zap.Error(err))
return ErrResourceNotOwned
}
if !ok {
logs.Warn(ctx, "service.meeting_signal_service.assertOwnsResource", "资源归属不匹配,拒绝越权操作",
zap.String("room_code", roomCode), zap.Int64("user_id", userID),
zap.String("kind", kind), zap.String("id", id))
return ErrResourceNotOwned
}
return nil
}
// updateMemberState 将某用户的音视频开关状态持久化到 Redis Hash
// 非 nil 的字段才写入audio/video 任意一个 nil 都不碰它,避免误覆盖另一维状态
// 容错:任何一步失败仅 Warn 日志,不阻断业务流程(前端已发 state.changed 作为权威广播)
func (s *MeetingSignalService) updateMemberState(ctx context.Context, roomCode string, userID int64, audio, video *bool) {
if audio == nil && video == nil {
return
}
key := memberStateKey(roomCode, userID)
values := make([]interface{}, 0, 4)
if audio != nil {
values = append(values, "audio_enabled", boolToStr(*audio))
}
if video != nil {
values = append(values, "video_enabled", boolToStr(*video))
}
if err := s.redis.HSet(ctx, key, values...).Err(); err != nil {
logs.Warn(ctx, "service.meeting_signal_service.updateMemberState", "写入成员音视频状态失败",
zap.String("key", key), zap.Error(err))
return
}
_ = s.redis.Expire(ctx, key, resourceTTL).Err()
}
// loadMemberState 读取某用户的音视频开关状态Hash 不存在时两个字段都返回 (nil, nil)
func (s *MeetingSignalService) loadMemberState(ctx context.Context, roomCode string, userID int64) (audio *bool, video *bool) {
key := memberStateKey(roomCode, userID)
fields, err := s.redis.HGetAll(ctx, key).Result()
if err != nil || len(fields) == 0 {
return nil, nil
}
if v, ok := fields["audio_enabled"]; ok {
b := v == "true"
audio = &b
}
if v, ok := fields["video_enabled"]; ok {
b := v == "true"
video = &b
}
return audio, video
}
func boolToStr(b bool) string {
if b {
return "true"
}
return "false"
}
// loadRoomAndParticipant 通用前置校验:拉取房间 + 确认用户是活跃参会者
// 所有信令事件在进入业务前都要过这一关;返回的 *MeetingRoom 供后续广播使用 roomID
func (s *MeetingSignalService) loadRoomAndParticipant(ctx context.Context, roomCode string, userID int64) (*model.MeetingRoom, error) {
room, err := s.roomDAO.GetByCode(ctx, roomCode)
if err != nil {
return nil, err
}
if room == nil {
return nil, ErrMeetingNotFound
}
if room.Status == constants.MeetingStatusEnded {
return nil, ErrMeetingEnded
}
p, err := s.participantDAO.GetByRoomAndUser(ctx, room.ID, userID)
if err != nil {
return nil, err
}
if p == nil || !p.IsActive() {
return nil, ErrNotInMeeting
}
return room, nil
}
// ========== 房间组2 个 C→S==========
// OnRoomJoin 处理 meeting.room.join 事件
// 语义:客户端 REST 加入会议成功后,通过 WS 宣告在线;服务端记录 userID ↔ roomCode 映射
// 副作用Task 15 增强):补推"房间内其他用户已有的 producer 列表"给刚 join 的用户,
// 解决 Mediasoup SFU "后入者错过历史 producer.new"的经典问题
func (s *MeetingSignalService) OnRoomJoin(ctx context.Context, userID int64, roomCode string) error {
room, err := s.loadRoomAndParticipant(ctx, roomCode, userID)
if err != nil {
return err
}
// 触发资源追踪 key 续期(空集合 TTL 续期无副作用)
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),
zap.Int64("room_id", room.ID))
// P1-8pushExistingRoomState 同步化Task 16
// 旧版异步 go 导致 REST `meeting.member.joined` 与 WS `pushExistingRoomState` 并行,
// 极端场景下新人先收到未知 user_id 的 state.changed 再收到 REST participant名称头像短暂缺失
// 前端本就在等 room.join 的 ACK这里同步在 ACK 返回前完成补推,可消除并行窗口。
// 单次调用仅 O(当前房间 producer 数量) Redis 读 + 若干 WS 发送P99 <30ms不影响 ACK 体验。
s.pushExistingRoomState(ctx, room.ID, room.RoomCode, userID)
return nil
}
// pushExistingRoomState 向刚加入者定向推送房间的历史媒体状态,分两件事:
// 1. producer 列表:复用 meeting.member.producer.new 事件语义,前端自动发起 consume
// 2. 成员 audio/video 开关:复用 meeting.member.state.changed 事件语义,前端更新成员面板图标
//
// 容错:任何单步失败仅 Warn 日志,不中断流程
func (s *MeetingSignalService) pushExistingRoomState(ctx context.Context, roomID int64, roomCode string, userID int64) {
s.pushExistingProducers(ctx, roomID, roomCode, userID)
s.pushExistingMemberStates(ctx, roomID, roomCode, userID)
}
// pushExistingProducers 向刚加入者定向推送房间里其他用户已产生的 producer 列表
// 复用 meeting.member.producer.new 事件语义,前端 _onProducerNew handler 无需改动
// 容错:任何单步失败仅 Warn 日志,不中断流程
func (s *MeetingSignalService) pushExistingProducers(ctx context.Context, roomID int64, roomCode string, userID int64) {
funcName := "service.meeting_signal_service.pushExistingProducers"
actives, err := s.participantDAO.ListActiveByRoom(ctx, roomID)
if err != nil {
logs.Warn(ctx, funcName, "列出房间活跃参会者失败",
zap.Int64("room_id", roomID), zap.Error(err))
return
}
pushCount := 0
for i := range actives {
other := &actives[i]
if other.UserID == userID {
continue
}
otherKey := resourceTrackKey(roomCode, other.UserID)
members, err := s.redis.SMembers(ctx, otherKey).Result()
if err != nil {
logs.Warn(ctx, funcName, "读取他人资源追踪集合失败",
zap.String("key", otherKey), zap.Error(err))
continue
}
for _, m := range members {
// 格式:"kind:id";仅关心 producer
idx := -1
for j, c := range m {
if c == ':' {
idx = j
break
}
}
if idx < 0 {
continue
}
kind, id := m[:idx], m[idx+1:]
if kind != "producer" || id == "" {
continue
}
if err := s.broadcaster.PublishToUser(ctx, userID, constants.MeetingWSEventMemberProducerNew, map[string]interface{}{
"room_code": roomCode,
"user_id": other.UserID,
"producer_id": id,
"existing": true,
}); err != nil {
logs.Warn(ctx, funcName, "定向推送 existing producer 失败",
zap.Int64("to_user", userID),
zap.Int64("owner_user", other.UserID),
zap.String("producer_id", id),
zap.Error(err))
continue
}
pushCount++
}
}
if pushCount > 0 {
logs.Info(ctx, funcName, "补推已有 producer 完成",
zap.String("room_code", roomCode),
zap.Int64("to_user", userID),
zap.Int("push_count", pushCount))
}
}
// pushExistingMemberStates 向刚加入者定向推送房间里其他活跃成员当前的音视频开关状态
// 数据源updateMemberState 在每次 OnMemberStateChanged 成功后写入的 Redis Hash
// 前端 _onMemberStateChanged 处理逻辑已存在,收到本事件后会刷新 MemberPanel 图标色彩
// 容错:读取 Hash 失败、该用户暂无状态记录都直接跳过
func (s *MeetingSignalService) pushExistingMemberStates(ctx context.Context, roomID int64, roomCode string, userID int64) {
funcName := "service.meeting_signal_service.pushExistingMemberStates"
actives, err := s.participantDAO.ListActiveByRoom(ctx, roomID)
if err != nil {
logs.Warn(ctx, funcName, "列出房间活跃参会者失败",
zap.Int64("room_id", roomID), zap.Error(err))
return
}
pushCount := 0
for i := range actives {
other := &actives[i]
if other.UserID == userID {
continue
}
audio, video := s.loadMemberState(ctx, roomCode, other.UserID)
if audio == nil && video == nil {
continue
}
data := map[string]interface{}{
"room_code": roomCode,
"user_id": other.UserID,
"changed_by": other.UserID,
"existing": true,
}
if audio != nil {
data["audio_enabled"] = *audio
}
if video != nil {
data["video_enabled"] = *video
}
if err := s.broadcaster.PublishToUser(ctx, userID, constants.MeetingWSEventMemberStateChange, data); err != nil {
logs.Warn(ctx, funcName, "定向推送 existing state.changed 失败",
zap.Int64("to_user", userID),
zap.Int64("owner_user", other.UserID),
zap.Error(err))
continue
}
pushCount++
}
if pushCount > 0 {
logs.Info(ctx, funcName, "补推已有成员状态完成",
zap.String("room_code", roomCode),
zap.Int64("to_user", userID),
zap.Int("push_count", pushCount))
}
}
// 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
// 参会者表的 LeaveRoom 逻辑仍由 REST API 负责(避免 WS 并发引起 left_at 重复写入)
func (s *MeetingSignalService) OnRoomLeave(ctx context.Context, userID int64, roomCode string) error {
room, err := s.loadRoomAndParticipant(ctx, roomCode, userID)
if err != nil {
return err
}
s.cleanupUserResources(ctx, roomCode, userID)
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{
"room_code": roomCode,
"user_id": userID,
"reason": "ws_disconnect",
}, userID)
return nil
}
// ========== 成员组1 个双向)==========
// MemberStateChangePayload meeting.member.state.changed 请求载荷
// host 可通过 target_user_id 静音他人 / 关其摄像头;非 host 传该字段将被拒绝
type MemberStateChangePayload struct {
RoomCode string `json:"room_code"`
TargetUserID int64 `json:"target_user_id,omitempty"` // 可选host 强制他人状态
AudioEnabled *bool `json:"audio_enabled,omitempty"` // nil 表示不改
VideoEnabled *bool `json:"video_enabled,omitempty"`
}
// OnMemberStateChanged 处理 meeting.member.state.changed 事件
// 权限:操作自己无限制;操作他人必须是 host
// 行为:广播 meeting.member.state.changed 给房间其他成员(发起者自己不收到回显)
func (s *MeetingSignalService) OnMemberStateChanged(ctx context.Context, fromUserID int64, payload *MemberStateChangePayload) error {
room, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, fromUserID)
if err != nil {
return err
}
targetID := fromUserID
if payload.TargetUserID != 0 && payload.TargetUserID != fromUserID {
if room.HostID != fromUserID {
return ErrNotMeetingHost
}
targetP, err := s.participantDAO.GetByRoomAndUser(ctx, room.ID, payload.TargetUserID)
if err != nil {
return err
}
if targetP == nil || !targetP.IsActive() {
return ErrTransferTargetInvalid
}
targetID = payload.TargetUserID
}
data := map[string]interface{}{
"room_code": payload.RoomCode,
"user_id": targetID,
"changed_by": fromUserID,
}
if payload.AudioEnabled != nil {
data["audio_enabled"] = *payload.AudioEnabled
}
if payload.VideoEnabled != nil {
data["video_enabled"] = *payload.VideoEnabled
}
// 持久化到 Redis Hash作为后入者补推 state.changed 的数据源
// 注意:写 Hash 之所以同步执行而非放进 goroutine是因为同一用户并发开关操作需要顺序一致
s.updateMemberState(ctx, payload.RoomCode, targetID, payload.AudioEnabled, payload.VideoEnabled)
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberStateChange, data, fromUserID)
return nil
}
// ========== 媒体组5 个mediasoup signaling==========
// TransportCreatePayload meeting.transport.create 请求载荷
type TransportCreatePayload struct {
RoomCode string `json:"room_code"`
Direction string `json:"direction"` // "send" | "recv"
}
// OnTransportCreate 处理 meeting.transport.create 事件
func (s *MeetingSignalService) OnTransportCreate(ctx context.Context, userID int64, payload *TransportCreatePayload) (*TransportInfo, error) {
if payload.Direction != "send" && payload.Direction != "recv" {
return nil, fmt.Errorf("direction 非法,必须是 send 或 recv")
}
if _, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID); err != nil {
return nil, err
}
info, err := s.mediaOrchestrator.CreateTransport(ctx, &CreateTransportReq{
RoomCode: payload.RoomCode,
UserID: userID,
Direction: payload.Direction,
})
if err != nil {
return nil, err
}
s.trackResource(ctx, payload.RoomCode, userID, "transport", info.ID)
return info, nil
}
// TransportConnectPayload meeting.transport.connect 请求载荷
type TransportConnectPayload struct {
RoomCode string `json:"room_code"`
TransportID string `json:"transport_id"`
DtlsParameters json.RawMessage `json:"dtls_parameters"`
}
// OnTransportConnect 处理 meeting.transport.connect 事件
func (s *MeetingSignalService) OnTransportConnect(ctx context.Context, userID int64, payload *TransportConnectPayload) error {
if _, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID); err != nil {
return err
}
if err := s.assertOwnsResource(ctx, payload.RoomCode, userID, "transport", payload.TransportID); err != nil {
return err
}
return s.mediaOrchestrator.ConnectTransport(ctx, payload.TransportID, payload.DtlsParameters)
}
// ProduceStartPayload meeting.produce.start 请求载荷
type ProduceStartPayload struct {
RoomCode string `json:"room_code"`
TransportID string `json:"transport_id"`
Kind string `json:"kind"` // "audio" | "video"
RtpParameters json.RawMessage `json:"rtp_parameters"`
}
// ProduceStartResult 返回给客户端的 producerID
type ProduceStartResult struct {
ProducerID string `json:"producer_id"`
}
// OnProduceStart 处理 meeting.produce.start 事件
// 成功后广播 meeting.member.producer.new 给房间内其他成员,驱动对端自动创建 Consumer
func (s *MeetingSignalService) OnProduceStart(ctx context.Context, userID int64, payload *ProduceStartPayload) (*ProduceStartResult, error) {
if payload.Kind != "audio" && payload.Kind != "video" {
return nil, fmt.Errorf("kind 非法,必须是 audio 或 video")
}
room, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID)
if err != nil {
return nil, err
}
if err := s.assertOwnsResource(ctx, payload.RoomCode, userID, "transport", payload.TransportID); err != nil {
return nil, err
}
producerID, err := s.mediaOrchestrator.CreateProducer(ctx, &CreateProducerReq{
RoomCode: payload.RoomCode,
UserID: userID,
TransportID: payload.TransportID,
Kind: payload.Kind,
RtpParameters: payload.RtpParameters,
})
if err != nil {
return nil, err
}
s.trackResource(ctx, payload.RoomCode, userID, "producer", producerID)
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberProducerNew, map[string]interface{}{
"room_code": payload.RoomCode,
"user_id": userID,
"producer_id": producerID,
"kind": payload.Kind,
}, userID)
return &ProduceStartResult{ProducerID: producerID}, nil
}
// ConsumeStartPayload meeting.consume.start 请求载荷
type ConsumeStartPayload struct {
RoomCode string `json:"room_code"`
TransportID string `json:"transport_id"` // 客户端 recv Transport
ProducerID string `json:"producer_id"` // 要订阅的远端 Producer
RtpCapabilities json.RawMessage `json:"rtp_capabilities"`
}
// OnConsumeStart 处理 meeting.consume.start 事件
// P0-2 修复:新增 transport_id 归属校验,拒绝把 consumer 挂到他人的 recv transport
// producer_id 归属天然不需要校验:消费他人 producer 正是订阅逻辑本身Node 侧会验证 producer 是否存在)
func (s *MeetingSignalService) OnConsumeStart(ctx context.Context, userID int64, payload *ConsumeStartPayload) (*ConsumerInfo, error) {
if payload.ProducerID == "" {
return nil, fmt.Errorf("producer_id 不能为空")
}
if _, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID); err != nil {
return nil, err
}
if err := s.assertOwnsResource(ctx, payload.RoomCode, userID, "transport", payload.TransportID); err != nil {
return nil, err
}
info, err := s.mediaOrchestrator.CreateConsumer(ctx, &CreateConsumerReq{
RoomCode: payload.RoomCode,
UserID: userID,
TransportID: payload.TransportID,
ProducerID: payload.ProducerID,
RtpCapabilities: payload.RtpCapabilities,
})
if err != nil {
return nil, err
}
s.trackResource(ctx, payload.RoomCode, userID, "consumer", info.ID)
return info, nil
}
// ConsumeResumePayload meeting.consume.resume 请求载荷Task 9
// 前端 recv Transport 与 track 就绪后发送,用于把 Node 侧 paused Consumer 切到 active
type ConsumeResumePayload struct {
RoomCode string `json:"room_code"`
ConsumerID string `json:"consumer_id"`
}
// OnConsumeResume 处理 meeting.consume.resume 事件Task 9
// 语义:告知 media-server 把指定 Consumer 从 paused 切到 active
// 权限:仅当 userID 是会议活跃成员且 consumerID 归属该用户时允许
// 幂等Node 对已 active Consumer 再次 resume 不报错Consumer 不存在则 ACK 返回友好错误
func (s *MeetingSignalService) OnConsumeResume(ctx context.Context, userID int64, payload *ConsumeResumePayload) error {
if _, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID); err != nil {
return err
}
if err := s.assertOwnsResource(ctx, payload.RoomCode, userID, "consumer", payload.ConsumerID); err != nil {
return err
}
if err := s.mediaOrchestrator.ResumeConsumer(ctx, payload.ConsumerID); err != nil {
logs.Warn(ctx, "service.meeting_signal_service.OnConsumeResume", "恢复 Consumer 失败",
zap.String("room_code", payload.RoomCode),
zap.Int64("user_id", userID),
zap.String("consumer_id", payload.ConsumerID),
zap.Error(err))
return err
}
return nil
}
// ProducerClosePayload meeting.producer.close 请求载荷
type ProducerClosePayload struct {
RoomCode string `json:"room_code"`
ProducerID string `json:"producer_id"`
}
// OnProducerClose 处理 meeting.producer.close 事件
// 成功后广播给房间内其他成员(与 mediasoup 的 producerclose 级联动作平级)
// P0-1 修复:新增 producer 归属校验,拒绝一名参会人关闭他人的 producer横向越权审计 P0-1
func (s *MeetingSignalService) OnProducerClose(ctx context.Context, userID int64, payload *ProducerClosePayload) error {
room, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID)
if err != nil {
return err
}
if err := s.assertOwnsResource(ctx, payload.RoomCode, userID, "producer", payload.ProducerID); err != nil {
return err
}
if err := s.mediaOrchestrator.CloseProducer(ctx, payload.ProducerID); err != nil {
logs.Warn(ctx, "service.meeting_signal_service.OnProducerClose", "关闭 Producer 失败",
zap.String("producer_id", payload.ProducerID), zap.Error(err))
}
s.untrackResource(ctx, payload.RoomCode, userID, "producer", payload.ProducerID)
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberProducerNew, map[string]interface{}{
"room_code": payload.RoomCode,
"user_id": userID,
"producer_id": payload.ProducerID,
"closed": true,
}, userID)
return nil
}
// ========== 资源清理 ==========
// cleanupUserResources 批量关闭指定用户在某会议的所有媒体资源
// WS 断开、主动离会、被踢时使用;依赖 Redis 集合中追踪的资源 ID
func (s *MeetingSignalService) cleanupUserResources(ctx context.Context, roomCode string, userID int64) {
funcName := "service.meeting_signal_service.cleanupUserResources"
key := resourceTrackKey(roomCode, userID)
members, err := s.redis.SMembers(ctx, key).Result()
if err != nil {
logs.Warn(ctx, funcName, "读取资源追踪集合失败", zap.String("key", key), zap.Error(err))
return
}
for _, m := range members {
// 格式:"kind:id"
idx := -1
for i, c := range m {
if c == ':' {
idx = i
break
}
}
if idx < 0 {
continue
}
kind, id := m[:idx], m[idx+1:]
switch kind {
case "producer":
_ = s.mediaOrchestrator.CloseProducer(ctx, id)
case "consumer":
_ = s.mediaOrchestrator.CloseConsumer(ctx, id)
// transport 关闭一般由 Router 级联;这里不单独处理
}
}
_ = s.redis.Del(ctx, key).Err()
// 一并清理成员音视频状态,避免对方下次入会时读到旧 host 的僵尸 AV 状态
_ = s.redis.Del(ctx, memberStateKey(roomCode, userID)).Err()
if len(members) > 0 {
logs.Info(ctx, funcName, "清理用户媒体资源",
zap.String("room_code", roomCode),
zap.Int64("user_id", userID),
zap.Int("resource_count", len(members)))
}
}
// cleanupRoomRedisResidual 批量清理指定 roomCode 下一批用户的 Redis 资源追踪集合 + 音视频状态 Hash
// Task 16 资源清理专项:
// - EndRoom / HandleEmptyRoomExpired 等会议销毁路径调用,避免正常 WS 断开钩子尚未触发就结束会议导致的残留
// - userIDs 为销毁前活跃成员快照;无活跃成员可传 nil/空切片,此时仅走上层 lifecycle 清理
// - 使用 Pipeline 减少 Redis 往返;失败仅 Warn不阻断上层销毁流程
// - package 私有:允许 meeting_service / meeting_lifecycle_service 等同包文件直接调用
func cleanupRoomRedisResidual(ctx context.Context, rdb *redis.Client, roomCode string, userIDs []int64) {
funcName := "service.meeting_signal_service.cleanupRoomRedisResidual"
if rdb == nil || len(userIDs) == 0 {
return
}
pipe := rdb.Pipeline()
for _, uid := range userIDs {
pipe.Del(ctx, resourceTrackKey(roomCode, uid))
pipe.Del(ctx, memberStateKey(roomCode, uid))
}
if _, err := pipe.Exec(ctx); err != nil {
logs.Warn(ctx, funcName, "批量清理用户 Redis 残留失败(忽略,继续)",
zap.String("room_code", roomCode), zap.Int("user_count", len(userIDs)), zap.Error(err))
return
}
logs.Debug(ctx, funcName, "会议销毁,用户 Redis 残留已清理",
zap.String("room_code", roomCode), zap.Int("user_count", len(userIDs)))
}