Files
EchoChat/backend/go-service/app/meeting/service/meeting_signal_service.go
bujinyuan 790996feaf feat(meeting): Task 15 UI 打磨 + 主持人四件套 + 媒体层稳定性回归修复
- Task 15 UI:6 项原创特色(说话者流光 / 柔性网格 / 自视频浮窗 / 静音氛围色 / 入会滑入 / NetworkBadge 3 条波浪)+ 说话者双源探测(RTP audioLevel + WebAudio RMS)+ 主持人四件套(静音/开麦/转让/踢出)
- 新增 SelfVideoFloat 浮窗组件 + 4 份 design-system 页面文档(home / preview / room / invite)
- 媒体回归补丁(手工联调触发):
  * 后端 signal_service 在 OnRoomJoin 追加 pushExistingRoomState → 向新加入者补发历史 producers + 历史成员 audio/video 状态
  * 新增 Redis Hash memberStateKey 持久化成员 audio/video enabled,OnMemberStateChanged 落盘、cleanupUserResources 清理
  * 前端 _broadcastSelfState 在本地音视频开关末尾同步 state.changed;_afterJoined 先 getRoom 再 room.join 修复后入者成员状态时序;_onMemberStateChanged 占位兜底
  * _cleanupRemoteProducer 改为按 producerId 精准清理(修复 "关音频误关视频")
  * _onRoomEnded + createAndEnter/joinAndEnter 强化 reset(修复重入会看不到自己画面)
  * mediasoup-client ensureSend/RecvTransport 引入 in-flight Promise 锁防并发重复建连
- 文档同步:CURRENT_STATUS.md / Task 15 plan / phase2e-2 implementation plan

Made-with: Cursor
2026-04-23 16:04:46 +08:00

672 lines
25 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()
}
// 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))
// 定向补推已有 producer + 成员状态:异步执行,不阻塞 ACK
go s.pushExistingRoomState(context.Background(), 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(context.Background(), 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(context.Background(), 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 payload.TransportID == "" {
return fmt.Errorf("transport_id 不能为空")
}
if _, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID); 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")
}
if payload.TransportID == "" {
return nil, fmt.Errorf("transport_id 不能为空")
}
room, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID)
if 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(context.Background(), 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 事件
func (s *MeetingSignalService) OnConsumeStart(ctx context.Context, userID int64, payload *ConsumeStartPayload) (*ConsumerInfo, error) {
if payload.TransportID == "" || payload.ProducerID == "" {
return nil, fmt.Errorf("transport_id 与 producer_id 均不能为空")
}
if _, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID); 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 payload.ConsumerID == "" {
return fmt.Errorf("consumer_id 不能为空")
}
if _, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID); 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 级联动作平级)
func (s *MeetingSignalService) OnProducerClose(ctx context.Context, userID int64, payload *ProducerClosePayload) error {
if payload.ProducerID == "" {
return fmt.Errorf("producer_id 不能为空")
}
room, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID)
if 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(context.Background(), 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)))
}
}