feat(phase2e-2): WebSocket 信令协议 13 事件全量落地(Task 6)
将 meeting.* 事件族从 Task 5 的 PublishToUser 循环升级为完整 WS 信令协议:
抽离统一广播层、扩容 MediaOrchestrator 接口、实现 8 个 C→S 事件业务逻辑
+ Redis 资源追踪 + host 权限校验,端到端 18/18 PASS。
核心产出:
- 新建 MeetingBroadcaster(统一广播层,REST/WS 共用)
- 新建 MeetingSignalService 8 C→S 事件 + cleanupUserResources
- 新建 MeetingWSHandler 薄层 controller
- MediaOrchestrator 扩容 9 方法 + NoopMediaOrchestrator 占位(Task 7 替换)
- C→S 白名单机制防恶意伪造广播事件
- Redis Set 资源追踪防 mediasoup 端资源泄漏
文档同步:
- docs/api/frontend/meeting.md 追加 §WebSocket 信令协议(Task 6)200 行
- docs/progress/CURRENT_STATUS.md + project-context.mdc + 实施计划 Task 6 ✅
Made-with: Cursor
This commit is contained in:
196
backend/go-service/app/meeting/controller/meeting_ws_handler.go
Normal file
196
backend/go-service/app/meeting/controller/meeting_ws_handler.go
Normal file
@@ -0,0 +1,196 @@
|
||||
// Package controller 会议模块 HTTP / WS 入口
|
||||
package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
|
||||
"github.com/echochat/backend/app/constants"
|
||||
"github.com/echochat/backend/app/meeting/service"
|
||||
"github.com/echochat/backend/pkg/logs"
|
||||
"github.com/echochat/backend/pkg/ws"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
// MeetingWSHandler 会议 WS 事件入口
|
||||
// Task 6 落地:注册设计 §6.3 的 11 个 meeting.* 事件到 Hub 路由表
|
||||
// 薄层设计:
|
||||
// - 仅负责 JSON Unmarshal → 调 SignalService 业务方法 → 组装 ACK 响应
|
||||
// - 权限校验全部下沉到 SignalService(service 层统一 assertIsActiveParticipant / assertIsHost)
|
||||
// - 广播副作用由 SignalService 自行发起(ACK 不携带广播数据)
|
||||
type MeetingWSHandler struct {
|
||||
signalSvc *service.MeetingSignalService
|
||||
hub *ws.Hub
|
||||
}
|
||||
|
||||
// NewMeetingWSHandler 构造 WS Handler 并自动注册事件到 Hub
|
||||
// 与 im.handler.EventHandler 保持一致的风格:构造时注册,启动时由 DI 容器持有
|
||||
func NewMeetingWSHandler(signalSvc *service.MeetingSignalService, hub *ws.Hub) *MeetingWSHandler {
|
||||
h := &MeetingWSHandler{
|
||||
signalSvc: signalSvc,
|
||||
hub: hub,
|
||||
}
|
||||
h.registerEvents()
|
||||
return h
|
||||
}
|
||||
|
||||
// registerEvents 将所有 meeting.* C→S 事件注册到 Hub 路由表
|
||||
// S→C 事件(ended/joined/left/host.changed/producer.new 等)不在此处注册,由 broadcaster 发出
|
||||
func (h *MeetingWSHandler) registerEvents() {
|
||||
// 房间组
|
||||
h.hub.RegisterEvent(constants.MeetingWSEventRoomJoin, h.handleRoomJoin)
|
||||
h.hub.RegisterEvent(constants.MeetingWSEventRoomLeave, h.handleRoomLeave)
|
||||
// 成员组
|
||||
h.hub.RegisterEvent(constants.MeetingWSEventMemberStateChange, h.handleMemberStateChange)
|
||||
// 媒体组(5 个)
|
||||
h.hub.RegisterEvent(constants.MeetingWSEventTransportCreate, h.handleTransportCreate)
|
||||
h.hub.RegisterEvent(constants.MeetingWSEventTransportConnect, h.handleTransportConnect)
|
||||
h.hub.RegisterEvent(constants.MeetingWSEventProduceStart, h.handleProduceStart)
|
||||
h.hub.RegisterEvent(constants.MeetingWSEventConsumeStart, h.handleConsumeStart)
|
||||
h.hub.RegisterEvent(constants.MeetingWSEventProducerClose, h.handleProducerClose)
|
||||
}
|
||||
|
||||
// ====== 通用工具 ======
|
||||
|
||||
// simpleRoomPayload 房间组仅需 room_code 的请求体
|
||||
type simpleRoomPayload struct {
|
||||
RoomCode string `json:"room_code"`
|
||||
}
|
||||
|
||||
// sendACK 统一发送 ACK 响应
|
||||
func (h *MeetingWSHandler) sendACK(client *ws.Client, msg *ws.Message, code int, message string, data interface{}) {
|
||||
resp := ws.NewResponse(msg.Event, msg.Seq, code, message, data)
|
||||
bytes, err := ws.MarshalResponse(resp)
|
||||
if err != nil {
|
||||
logs.Error(context.Background(), "controller.meeting_ws_handler.sendACK", "序列化 ACK 失败",
|
||||
zap.String("event", msg.Event), zap.Error(err))
|
||||
return
|
||||
}
|
||||
client.Send(bytes)
|
||||
}
|
||||
|
||||
// unmarshal 通用反序列化 + 错误 ACK
|
||||
func (h *MeetingWSHandler) unmarshal(client *ws.Client, msg *ws.Message, target interface{}) bool {
|
||||
if err := json.Unmarshal(msg.Data, target); err != nil {
|
||||
logs.Warn(nil, "controller.meeting_ws_handler.unmarshal", "反序列化失败",
|
||||
zap.String("event", msg.Event),
|
||||
zap.Int64("user_id", client.UserID),
|
||||
zap.Error(err))
|
||||
h.sendACK(client, msg, -1, "请求参数格式错误", nil)
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// ====== 房间组 ======
|
||||
|
||||
// handleRoomJoin 处理 meeting.room.join
|
||||
func (h *MeetingWSHandler) handleRoomJoin(client *ws.Client, msg *ws.Message) {
|
||||
var payload simpleRoomPayload
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
if err := h.signalSvc.OnRoomJoin(context.Background(), client.UserID, payload.RoomCode); err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
h.sendACK(client, msg, 0, "ok", nil)
|
||||
}
|
||||
|
||||
// handleRoomLeave 处理 meeting.room.leave
|
||||
func (h *MeetingWSHandler) handleRoomLeave(client *ws.Client, msg *ws.Message) {
|
||||
var payload simpleRoomPayload
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
if err := h.signalSvc.OnRoomLeave(context.Background(), client.UserID, payload.RoomCode); err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
h.sendACK(client, msg, 0, "ok", nil)
|
||||
}
|
||||
|
||||
// ====== 成员组 ======
|
||||
|
||||
// handleMemberStateChange 处理 meeting.member.state.changed
|
||||
func (h *MeetingWSHandler) handleMemberStateChange(client *ws.Client, msg *ws.Message) {
|
||||
var payload service.MemberStateChangePayload
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
if err := h.signalSvc.OnMemberStateChanged(context.Background(), client.UserID, &payload); err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
h.sendACK(client, msg, 0, "ok", nil)
|
||||
}
|
||||
|
||||
// ====== 媒体组 ======
|
||||
|
||||
// handleTransportCreate 处理 meeting.transport.create
|
||||
func (h *MeetingWSHandler) handleTransportCreate(client *ws.Client, msg *ws.Message) {
|
||||
var payload service.TransportCreatePayload
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
info, err := h.signalSvc.OnTransportCreate(context.Background(), client.UserID, &payload)
|
||||
if err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
h.sendACK(client, msg, 0, "ok", info)
|
||||
}
|
||||
|
||||
// handleTransportConnect 处理 meeting.transport.connect
|
||||
func (h *MeetingWSHandler) handleTransportConnect(client *ws.Client, msg *ws.Message) {
|
||||
var payload service.TransportConnectPayload
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
if err := h.signalSvc.OnTransportConnect(context.Background(), client.UserID, &payload); err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
h.sendACK(client, msg, 0, "ok", nil)
|
||||
}
|
||||
|
||||
// handleProduceStart 处理 meeting.produce.start
|
||||
func (h *MeetingWSHandler) handleProduceStart(client *ws.Client, msg *ws.Message) {
|
||||
var payload service.ProduceStartPayload
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
result, err := h.signalSvc.OnProduceStart(context.Background(), client.UserID, &payload)
|
||||
if err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
h.sendACK(client, msg, 0, "ok", result)
|
||||
}
|
||||
|
||||
// handleConsumeStart 处理 meeting.consume.start
|
||||
func (h *MeetingWSHandler) handleConsumeStart(client *ws.Client, msg *ws.Message) {
|
||||
var payload service.ConsumeStartPayload
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
info, err := h.signalSvc.OnConsumeStart(context.Background(), client.UserID, &payload)
|
||||
if err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
h.sendACK(client, msg, 0, "ok", info)
|
||||
}
|
||||
|
||||
// handleProducerClose 处理 meeting.producer.close
|
||||
func (h *MeetingWSHandler) handleProducerClose(client *ws.Client, msg *ws.Message) {
|
||||
var payload service.ProducerClosePayload
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
if err := h.signalSvc.OnProducerClose(context.Background(), client.UserID, &payload); err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
h.sendACK(client, msg, 0, "ok", nil)
|
||||
}
|
||||
@@ -9,15 +9,21 @@ import (
|
||||
|
||||
// MeetingSet 会议模块 Wire Provider Set
|
||||
// 对外暴露:
|
||||
// - *service.MeetingService —— 业务服务,供未来 ws handler / media-server 回调使用
|
||||
// - *service.MeetingService —— REST 业务服务(Task 5)
|
||||
// - *service.MeetingBroadcaster —— WS 广播中枢(Task 6 新增,供 SignalService 复用)
|
||||
// - *service.MeetingSignalService —— WS 信令事件业务逻辑(Task 6)
|
||||
// - *controller.MeetingController —— REST API 控制器
|
||||
// - *controller.MeetingWSHandler —— WS 事件 Handler(Task 6,启动时自动注册到 Hub)
|
||||
// 依赖的接口 NotifyPusher / UserInfoResolver / OnlineChecker 由上游 wire.Bind 绑定具体实现
|
||||
var MeetingSet = wire.NewSet(
|
||||
dao.NewMeetingRoomDAO,
|
||||
dao.NewMeetingParticipantDAO,
|
||||
dao.NewMeetingChatDAO,
|
||||
service.NewMeetingBroadcaster,
|
||||
service.NewMeetingService,
|
||||
service.NewMeetingSignalService,
|
||||
controller.NewMeetingController,
|
||||
controller.NewMeetingWSHandler,
|
||||
|
||||
// MediaOrchestrator 目前使用 Noop 实现(Task 7 将替换为 node_client.NodeClient)
|
||||
service.NewNoopMediaOrchestrator,
|
||||
|
||||
@@ -3,6 +3,7 @@ package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
|
||||
authModel "github.com/echochat/backend/app/auth/model"
|
||||
notifyService "github.com/echochat/backend/app/notify/service"
|
||||
@@ -30,22 +31,82 @@ type OnlineChecker interface {
|
||||
IsOnline(ctx context.Context, userID int64) bool
|
||||
}
|
||||
|
||||
// MediaOrchestrator 媒体服务器编排接口
|
||||
// Task 7 落地 Go → Node media-server HTTP Client 后由 node_client.NodeClient 实现
|
||||
// Task 5 阶段使用 NoopMediaOrchestrator 占位,仅返回假的 RouterID,不产生真实媒体资源
|
||||
// 语义:会议房间在 Go 侧落库后,通过该接口驱动 Node 端 mediasoup Router 创建与销毁
|
||||
type MediaOrchestrator interface {
|
||||
// CreateRouter 为会议房间创建 mediasoup Router
|
||||
// 入参:room_code 作为 Node 端聚合键;返回 Router ID 与推荐编解码参数(Task 7 对接时补充)
|
||||
CreateRouter(ctx context.Context, roomCode string) (routerID string, err error)
|
||||
// ====== MediaOrchestrator:Go → Node media-server 统一封装(设计 §6.6)======
|
||||
|
||||
// CloseRouter 关闭会议房间对应的 mediasoup Router 及其下所有 transport/producer/consumer
|
||||
// 幂等:重复关闭不返回错误
|
||||
CloseRouter(ctx context.Context, roomCode string) error
|
||||
// TransportInfo Node 端创建的 Transport 元信息
|
||||
// 字段与 mediasoup-client 的 Device.createSendTransport / createRecvTransport 入参兼容
|
||||
type TransportInfo struct {
|
||||
ID string `json:"id"` // Transport ID
|
||||
IceParameters json.RawMessage `json:"iceParameters"` // mediasoup ICE 参数
|
||||
IceCandidates json.RawMessage `json:"iceCandidates"` // ICE 候选列表
|
||||
DtlsParameters json.RawMessage `json:"dtlsParameters"` // DTLS 指纹参数
|
||||
SctpParameters json.RawMessage `json:"sctpParameters,omitempty"`
|
||||
}
|
||||
|
||||
// NoopMediaOrchestrator 占位实现:Task 5 完成生命周期接口时使用
|
||||
// 返回伪造的 RouterID,所有操作仅写日志不调用 Node
|
||||
// ConsumerInfo Node 端创建的 Consumer 元信息
|
||||
type ConsumerInfo struct {
|
||||
ID string `json:"id"`
|
||||
ProducerID string `json:"producerId"`
|
||||
Kind string `json:"kind"` // "audio" | "video"
|
||||
RtpParameters json.RawMessage `json:"rtpParameters"`
|
||||
Type string `json:"type,omitempty"` // "simple" | "simulcast" | "svc"
|
||||
ProducerPaused bool `json:"producerPaused,omitempty"`
|
||||
}
|
||||
|
||||
// CreateTransportReq 创建 Transport 请求
|
||||
type CreateTransportReq struct {
|
||||
RoomCode string `json:"roomCode"`
|
||||
UserID int64 `json:"userId"`
|
||||
Direction string `json:"direction"` // "send" | "recv"
|
||||
}
|
||||
|
||||
// CreateProducerReq 创建 Producer 请求
|
||||
type CreateProducerReq struct {
|
||||
RoomCode string `json:"roomCode"`
|
||||
UserID int64 `json:"userId"`
|
||||
TransportID string `json:"transportId"`
|
||||
Kind string `json:"kind"` // "audio" | "video"
|
||||
RtpParameters json.RawMessage `json:"rtpParameters"` // 直接转发给 Node,由其做结构校验
|
||||
}
|
||||
|
||||
// CreateConsumerReq 创建 Consumer 请求
|
||||
type CreateConsumerReq struct {
|
||||
RoomCode string `json:"roomCode"`
|
||||
UserID int64 `json:"userId"`
|
||||
TransportID string `json:"transportId"`
|
||||
ProducerID string `json:"producerId"`
|
||||
RtpCapabilities json.RawMessage `json:"rtpCapabilities"`
|
||||
}
|
||||
|
||||
// MediaOrchestrator 媒体服务器编排接口(设计 §6.6 NodeClient)
|
||||
// Task 7 落地 Go → Node media-server HTTP Client 后由 node_client.NodeClient 实现
|
||||
// Task 5/6 阶段使用 NoopMediaOrchestrator 占位:
|
||||
// - CreateRouter 返回 "noop-router-{code}";其他方法返回可解析的占位数据,用于 WS 信令链路自测
|
||||
// - 所有方法均幂等:重复调用不报错,符合 WS 信令重试语义
|
||||
type MediaOrchestrator interface {
|
||||
// CreateRouter 为会议房间创建 mediasoup Router
|
||||
CreateRouter(ctx context.Context, roomCode string) (routerID string, err error)
|
||||
// CloseRouter 关闭房间对应 Router 及其下所有资源(幂等)
|
||||
CloseRouter(ctx context.Context, roomCode string) error
|
||||
|
||||
// CreateTransport 为用户创建 send/recv WebRTC Transport
|
||||
CreateTransport(ctx context.Context, req *CreateTransportReq) (*TransportInfo, error)
|
||||
// ConnectTransport Transport DTLS 握手(幂等:重复 connect 对已连接 transport 视为成功)
|
||||
ConnectTransport(ctx context.Context, transportID string, dtlsParameters json.RawMessage) error
|
||||
|
||||
// CreateProducer 在指定 send Transport 上创建 Producer
|
||||
CreateProducer(ctx context.Context, req *CreateProducerReq) (producerID string, err error)
|
||||
// CloseProducer 关闭指定 Producer(幂等)
|
||||
CloseProducer(ctx context.Context, producerID string) error
|
||||
|
||||
// CreateConsumer 在指定 recv Transport 上创建 Consumer(订阅远端 Producer)
|
||||
CreateConsumer(ctx context.Context, req *CreateConsumerReq) (*ConsumerInfo, error)
|
||||
// CloseConsumer 关闭指定 Consumer(幂等)
|
||||
CloseConsumer(ctx context.Context, consumerID string) error
|
||||
}
|
||||
|
||||
// NoopMediaOrchestrator 占位实现:Task 7 完成前使用
|
||||
// 返回伪造的 ID 与固定占位数据(JSON:空对象 / 空数组),所有操作仅写日志不调用 Node
|
||||
// Task 7 完成后全局 wire 切换到真实 NodeClient 实现
|
||||
type NoopMediaOrchestrator struct{}
|
||||
|
||||
@@ -55,7 +116,6 @@ func NewNoopMediaOrchestrator() *NoopMediaOrchestrator {
|
||||
}
|
||||
|
||||
// CreateRouter 返回以 "noop-router-" 为前缀的伪造 RouterID
|
||||
// 调用方可据此区分真实 / 占位实现,便于调试与切换
|
||||
func (n *NoopMediaOrchestrator) CreateRouter(_ context.Context, roomCode string) (string, error) {
|
||||
return "noop-router-" + roomCode, nil
|
||||
}
|
||||
@@ -64,3 +124,44 @@ func (n *NoopMediaOrchestrator) CreateRouter(_ context.Context, roomCode string)
|
||||
func (n *NoopMediaOrchestrator) CloseRouter(_ context.Context, _ string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// CreateTransport 占位:返回以 "noop-transport-" 为前缀的伪造 ID,带最小合法 JSON 结构
|
||||
func (n *NoopMediaOrchestrator) CreateTransport(_ context.Context, req *CreateTransportReq) (*TransportInfo, error) {
|
||||
return &TransportInfo{
|
||||
ID: "noop-transport-" + req.Direction,
|
||||
IceParameters: json.RawMessage(`{}`),
|
||||
IceCandidates: json.RawMessage(`[]`),
|
||||
DtlsParameters: json.RawMessage(`{}`),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// ConnectTransport 占位:直接返回 nil
|
||||
func (n *NoopMediaOrchestrator) ConnectTransport(_ context.Context, _ string, _ json.RawMessage) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// CreateProducer 占位:返回伪造 ID
|
||||
func (n *NoopMediaOrchestrator) CreateProducer(_ context.Context, req *CreateProducerReq) (string, error) {
|
||||
return "noop-producer-" + req.Kind, nil
|
||||
}
|
||||
|
||||
// CloseProducer 占位:直接返回 nil
|
||||
func (n *NoopMediaOrchestrator) CloseProducer(_ context.Context, _ string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// CreateConsumer 占位:返回伪造的 ConsumerInfo(含目标 producerID 回填)
|
||||
func (n *NoopMediaOrchestrator) CreateConsumer(_ context.Context, req *CreateConsumerReq) (*ConsumerInfo, error) {
|
||||
return &ConsumerInfo{
|
||||
ID: "noop-consumer-" + req.ProducerID,
|
||||
ProducerID: req.ProducerID,
|
||||
Kind: "video",
|
||||
RtpParameters: json.RawMessage(`{}`),
|
||||
Type: "simple",
|
||||
}, nil
|
||||
}
|
||||
|
||||
// CloseConsumer 占位:直接返回 nil
|
||||
func (n *NoopMediaOrchestrator) CloseConsumer(_ context.Context, _ string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -0,0 +1,91 @@
|
||||
// Package service 提供 meeting 模块的业务逻辑
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/echochat/backend/app/meeting/dao"
|
||||
"github.com/echochat/backend/pkg/logs"
|
||||
"github.com/echochat/backend/pkg/ws"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
// MeetingBroadcaster 会议 WS 广播中枢
|
||||
// Task 6 引入,替代原 MeetingService.broadcastToActiveParticipants 内联实现
|
||||
// 语义:
|
||||
// - BroadcastToMeeting(ctx, roomID, event, payload, excludeUserIDs...) 向房间内所有活跃成员广播
|
||||
// - PublishToUser(ctx, userID, event, payload) 定向推送给指定用户(如 meeting.member.kicked 定向通知被踢者)
|
||||
// 通过 PubSub.Publish 跨实例传递;本地 Hub 若订阅了目标用户频道则自动路由到 WS 连接
|
||||
type MeetingBroadcaster struct {
|
||||
participantDAO *dao.MeetingParticipantDAO
|
||||
pubsub *ws.PubSub
|
||||
}
|
||||
|
||||
// NewMeetingBroadcaster 构造广播中枢
|
||||
func NewMeetingBroadcaster(participantDAO *dao.MeetingParticipantDAO, pubsub *ws.PubSub) *MeetingBroadcaster {
|
||||
return &MeetingBroadcaster{
|
||||
participantDAO: participantDAO,
|
||||
pubsub: pubsub,
|
||||
}
|
||||
}
|
||||
|
||||
// BroadcastToMeeting 向房间内所有活跃参会者广播 WS 事件
|
||||
// excludeUserIDs 中的用户将被跳过(通常排除发送者本人,避免"自己收到自己的事件")
|
||||
// 非阻塞:单个用户发送失败只记录 WARN 日志,不中断循环
|
||||
func (b *MeetingBroadcaster) BroadcastToMeeting(ctx context.Context, roomID int64, event string, data interface{}, excludeUserIDs ...int64) {
|
||||
funcName := "service.meeting_broadcaster.BroadcastToMeeting"
|
||||
|
||||
participants, err := b.participantDAO.ListActiveByRoom(ctx, roomID)
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "拉取活跃参会者失败",
|
||||
zap.Int64("room_id", roomID),
|
||||
zap.String("event", event),
|
||||
zap.Error(err))
|
||||
return
|
||||
}
|
||||
if len(participants) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
exclude := make(map[int64]struct{}, len(excludeUserIDs))
|
||||
for _, id := range excludeUserIDs {
|
||||
exclude[id] = struct{}{}
|
||||
}
|
||||
|
||||
msg := ws.NewPushMessage(event, data)
|
||||
pushed := 0
|
||||
for _, p := range participants {
|
||||
if _, skip := exclude[p.UserID]; skip {
|
||||
continue
|
||||
}
|
||||
if err := b.pubsub.PublishToUser(ctx, p.UserID, msg); err != nil {
|
||||
logs.Warn(ctx, funcName, "WS 广播单用户失败",
|
||||
zap.Int64("user_id", p.UserID),
|
||||
zap.String("event", event),
|
||||
zap.Error(err))
|
||||
continue
|
||||
}
|
||||
pushed++
|
||||
}
|
||||
|
||||
logs.Debug(ctx, funcName, "会议事件广播完成",
|
||||
zap.Int64("room_id", roomID),
|
||||
zap.String("event", event),
|
||||
zap.Int("pushed", pushed),
|
||||
zap.Int("total_active", len(participants)))
|
||||
}
|
||||
|
||||
// PublishToUser 定向推送给单个用户
|
||||
// 场景:meeting.member.kicked(被踢者收到的定向通知)、ACK 补偿推送等
|
||||
// 返回 error 让调用方根据业务语义决定是否需要感知失败
|
||||
func (b *MeetingBroadcaster) PublishToUser(ctx context.Context, userID int64, event string, data interface{}) error {
|
||||
msg := ws.NewPushMessage(event, data)
|
||||
if err := b.pubsub.PublishToUser(ctx, userID, msg); err != nil {
|
||||
logs.Warn(ctx, "service.meeting_broadcaster.PublishToUser", "定向 WS 推送失败",
|
||||
zap.Int64("user_id", userID),
|
||||
zap.String("event", event),
|
||||
zap.Error(err))
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -16,7 +16,6 @@ import (
|
||||
notifyService "github.com/echochat/backend/app/notify/service"
|
||||
"github.com/echochat/backend/pkg/logs"
|
||||
"github.com/echochat/backend/pkg/utils"
|
||||
"github.com/echochat/backend/pkg/ws"
|
||||
"github.com/redis/go-redis/v9"
|
||||
"go.uber.org/zap"
|
||||
"gorm.io/gorm"
|
||||
@@ -51,16 +50,17 @@ const (
|
||||
|
||||
// MeetingService 会议业务服务
|
||||
// Task 5 完成:会议生命周期、主持人管理、邀请、会议内聊天 12 个 REST API 全部落地
|
||||
// Task 6 会在此基础上追加 WS 信令事件处理器;Task 7 会把 mediaOrchestrator 的 Noop 实现替换为真实 Node HTTP Client
|
||||
// Task 6 新增:广播能力抽离到 MeetingBroadcaster;后续 WS 信令事件由 MeetingSignalService 处理
|
||||
// Task 7 会把 mediaOrchestrator 的 Noop 实现替换为真实 Node HTTP Client
|
||||
type MeetingService struct {
|
||||
roomDAO *dao.MeetingRoomDAO
|
||||
participantDAO *dao.MeetingParticipantDAO
|
||||
chatDAO *dao.MeetingChatDAO
|
||||
|
||||
db *gorm.DB
|
||||
redis *redis.Client
|
||||
pubsub *ws.PubSub
|
||||
db *gorm.DB
|
||||
redis *redis.Client
|
||||
|
||||
broadcaster *MeetingBroadcaster
|
||||
notifyPusher NotifyPusher
|
||||
userResolver UserInfoResolver
|
||||
onlineChecker OnlineChecker
|
||||
@@ -68,14 +68,15 @@ type MeetingService struct {
|
||||
}
|
||||
|
||||
// NewMeetingService 创建 MeetingService 实例
|
||||
// 依赖通过构造函数注入;接口依赖由上游 Wire 绑定到具体实现(Task 7 之前 mediaOrchestrator 使用 NoopMediaOrchestrator)
|
||||
// 依赖通过构造函数注入;接口依赖由上游 Wire 绑定到具体实现
|
||||
// Task 7 之前 mediaOrchestrator 使用 NoopMediaOrchestrator
|
||||
func NewMeetingService(
|
||||
roomDAO *dao.MeetingRoomDAO,
|
||||
participantDAO *dao.MeetingParticipantDAO,
|
||||
chatDAO *dao.MeetingChatDAO,
|
||||
db *gorm.DB,
|
||||
redis *redis.Client,
|
||||
pubsub *ws.PubSub,
|
||||
broadcaster *MeetingBroadcaster,
|
||||
notifyPusher NotifyPusher,
|
||||
userResolver UserInfoResolver,
|
||||
onlineChecker OnlineChecker,
|
||||
@@ -87,7 +88,7 @@ func NewMeetingService(
|
||||
chatDAO: chatDAO,
|
||||
db: db,
|
||||
redis: redis,
|
||||
pubsub: pubsub,
|
||||
broadcaster: broadcaster,
|
||||
notifyPusher: notifyPusher,
|
||||
userResolver: userResolver,
|
||||
onlineChecker: onlineChecker,
|
||||
@@ -139,28 +140,10 @@ func (s *MeetingService) generateUniqueRoomCode(ctx context.Context) (string, er
|
||||
}
|
||||
|
||||
// broadcastToActiveParticipants 向房间内所有活跃参会者广播 WS 事件
|
||||
// 调用 PubSub.Publish 支持多实例;可选排除发送者自身(excludeUserIDs)
|
||||
// 该辅助作为 Task 6 BroadcastToMeeting 的临时等效实现,签名保持兼容便于后续替换
|
||||
// Task 6 起实现已迁移至 MeetingBroadcaster.BroadcastToMeeting,本方法作为兼容壳保留
|
||||
// 以减少调用侧改动;未来可逐步替换为直接调用 s.broadcaster.BroadcastToMeeting
|
||||
func (s *MeetingService) broadcastToActiveParticipants(ctx context.Context, roomID int64, event string, data interface{}, excludeUserIDs ...int64) {
|
||||
funcName := "service.meeting_service.broadcastToActiveParticipants"
|
||||
participants, err := s.participantDAO.ListActiveByRoom(ctx, roomID)
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "拉取活跃参会者失败", zap.Int64("room_id", roomID), zap.Error(err))
|
||||
return
|
||||
}
|
||||
exclude := make(map[int64]struct{}, len(excludeUserIDs))
|
||||
for _, id := range excludeUserIDs {
|
||||
exclude[id] = struct{}{}
|
||||
}
|
||||
msg := ws.NewPushMessage(event, data)
|
||||
for _, p := range participants {
|
||||
if _, skip := exclude[p.UserID]; skip {
|
||||
continue
|
||||
}
|
||||
if err := s.pubsub.PublishToUser(ctx, p.UserID, msg); err != nil {
|
||||
logs.Warn(ctx, funcName, "WS 广播失败", zap.Int64("user_id", p.UserID), zap.String("event", event), zap.Error(err))
|
||||
}
|
||||
}
|
||||
s.broadcaster.BroadcastToMeeting(ctx, roomID, event, data, excludeUserIDs...)
|
||||
}
|
||||
|
||||
// ====== 会议生命周期 ======
|
||||
@@ -386,7 +369,7 @@ func (s *MeetingService) LeaveRoom(ctx context.Context, userID int64, code strin
|
||||
if uErr := s.roomDAO.UpdateHost(ctx, room.ID, newHost.UserID); uErr != nil {
|
||||
logs.Warn(ctx, funcName, "UpdateHost 失败", zap.Error(uErr))
|
||||
}
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventRoomHostChange, map[string]interface{}{
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventHostChanged, map[string]interface{}{
|
||||
"room_code": code,
|
||||
"old_host_id": userID,
|
||||
"new_host_id": newHost.UserID,
|
||||
@@ -441,15 +424,14 @@ func (s *MeetingService) EndRoom(ctx context.Context, userID int64, code string)
|
||||
}
|
||||
|
||||
payload := map[string]interface{}{
|
||||
"room_code": code,
|
||||
"room_code": code,
|
||||
"ended_reason": constants.MeetingEndedReasonHostEnded,
|
||||
"ended_at": now.Format("2006-01-02 15:04:05"),
|
||||
"ended_at": now.Format("2006-01-02 15:04:05"),
|
||||
}
|
||||
msg := ws.NewPushMessage(constants.MeetingWSEventRoomEnded, payload)
|
||||
// activesBefore 已经是结束前的活跃成员快照,此时 participantDAO.ListActiveByRoom 查会返回空集合
|
||||
// 因此使用 broadcaster.PublishToUser 逐人定向推送(而非 BroadcastToMeeting 基于当前状态查库)
|
||||
for _, p := range activesBefore {
|
||||
if err := s.pubsub.PublishToUser(ctx, p.UserID, msg); err != nil {
|
||||
logs.Warn(ctx, funcName, "meeting.room.ended 广播失败", zap.Int64("user_id", p.UserID), zap.Error(err))
|
||||
}
|
||||
_ = s.broadcaster.PublishToUser(ctx, p.UserID, constants.MeetingWSEventRoomEnded, payload)
|
||||
}
|
||||
|
||||
if err := s.mediaOrchestrator.CloseRouter(ctx, code); err != nil {
|
||||
@@ -498,7 +480,7 @@ func (s *MeetingService) TransferHost(ctx context.Context, operatorID int64, cod
|
||||
return err
|
||||
}
|
||||
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventRoomHostChange, map[string]interface{}{
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventHostChanged, map[string]interface{}{
|
||||
"room_code": code,
|
||||
"old_host_id": operatorID,
|
||||
"new_host_id": target.UserID,
|
||||
@@ -542,12 +524,11 @@ func (s *MeetingService) KickMember(ctx context.Context, operatorID int64, code
|
||||
return ErrNotInMeeting
|
||||
}
|
||||
|
||||
kickMsg := ws.NewPushMessage(constants.MeetingWSEventMemberKicked, map[string]interface{}{
|
||||
_ = s.broadcaster.PublishToUser(ctx, targetUserID, constants.MeetingWSEventMemberKicked, map[string]interface{}{
|
||||
"room_code": code,
|
||||
"user_id": targetUserID,
|
||||
"by": operatorID,
|
||||
})
|
||||
_ = s.pubsub.PublishToUser(ctx, targetUserID, kickMsg)
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{
|
||||
"room_code": code,
|
||||
"user_id": targetUserID,
|
||||
|
||||
393
backend/go-service/app/meeting/service/meeting_signal_service.go
Normal file
393
backend/go-service/app/meeting/service/meeting_signal_service.go
Normal file
@@ -0,0 +1,393 @@
|
||||
// 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
|
||||
}
|
||||
|
||||
// NewMeetingSignalService 构造 WS 信令服务
|
||||
func NewMeetingSignalService(
|
||||
roomDAO *dao.MeetingRoomDAO,
|
||||
participantDAO *dao.MeetingParticipantDAO,
|
||||
redis *redis.Client,
|
||||
broadcaster *MeetingBroadcaster,
|
||||
mediaOrchestrator MediaOrchestrator,
|
||||
) *MeetingSignalService {
|
||||
return &MeetingSignalService{
|
||||
roomDAO: roomDAO,
|
||||
participantDAO: participantDAO,
|
||||
redis: redis,
|
||||
broadcaster: broadcaster,
|
||||
mediaOrchestrator: mediaOrchestrator,
|
||||
}
|
||||
}
|
||||
|
||||
// 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)
|
||||
}
|
||||
|
||||
// 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()
|
||||
}
|
||||
|
||||
// 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 映射
|
||||
// 仅做存在性校验 + 心跳意义上的资源 key 刷新,不产生副作用
|
||||
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()
|
||||
|
||||
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))
|
||||
return nil
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
|
||||
// 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()
|
||||
|
||||
if len(members) > 0 {
|
||||
logs.Info(ctx, funcName, "清理用户媒体资源",
|
||||
zap.String("room_code", roomCode),
|
||||
zap.Int64("user_id", userID),
|
||||
zap.Int("resource_count", len(members)))
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user