视频会议保存

This commit is contained in:
duoaohui
2026-05-18 19:42:55 +08:00
parent 315884685f
commit ac8ec4c903
41 changed files with 5069 additions and 39 deletions

View File

@@ -314,14 +314,26 @@ func (h *HTTPMediaOrchestrator) CloseTransport(ctx context.Context, transportID
func (h *HTTPMediaOrchestrator) CreateProducer(ctx context.Context, req *CreateProducerReq) (string, error) {
funcName := "service.http_media_orchestrator.CreateProducer"
// 合并客户端 AppData 与服务端注入字段
// 服务端注入的 userId/roomCode 优先级最高,避免被客户端伪造覆盖
appData := map[string]any{}
if len(req.AppData) > 0 {
if err := json.Unmarshal(req.AppData, &appData); err != nil {
logs.Warn(ctx, funcName, "AppData 解析失败,已忽略",
zap.String("room_code", req.RoomCode),
zap.Int64("user_id", req.UserID),
zap.Error(err))
appData = map[string]any{}
}
}
appData["userId"] = req.UserID
appData["roomCode"] = req.RoomCode
reqBody := map[string]any{
"transportId": req.TransportID,
"kind": req.Kind,
"rtpParameters": req.RtpParameters,
"appData": map[string]any{
"userId": req.UserID,
"roomCode": req.RoomCode,
},
"appData": appData,
}
var resp struct {
ID string `json:"id"`
@@ -423,6 +435,67 @@ func (h *HTTPMediaOrchestrator) CloseConsumer(ctx context.Context, consumerID st
return err
}
// StartRecording 调用 POST /internal/v1/recordingsPhase B 引入)
//
// 注意:录制启动会触发 ffmpeg spawn + 端口分配 + 250ms 等待启动,正常耗时约 500ms~1.5s
// 因此走 TimeoutMS默认 10s而非 CloseTimeoutMS。失败不重试避免重复 spawn ffmpeg
//
// 错误映射:
// - Node 404router 不存在)→ ErrMediaResourceNotFound
// - Node 409producer 不存在/不可消费)→ ErrMediaServerError业务层转为 "录制启动失败"
func (h *HTTPMediaOrchestrator) StartRecording(ctx context.Context, req *StartRecordingReq) (*StartRecordingResp, error) {
funcName := "service.http_media_orchestrator.StartRecording"
// Node 端实际入参字段名与 zod schema 对齐
body := map[string]any{
"routerId": req.RouterID,
"roomCode": req.RoomCode,
"producerIds": req.ProducerIDs,
}
var resp StartRecordingResp
if err := h.doRequest(ctx, requestOptions{
method: http.MethodPost,
path: "/internal/v1/recordings",
body: body,
timeoutMS: h.cfg.TimeoutMS,
funcName: funcName,
logFields: []zap.Field{
zap.String("room_code", req.RoomCode),
zap.String("router_id", req.RouterID),
zap.Int("producer_count", len(req.ProducerIDs)),
},
}, &resp); err != nil {
return nil, err
}
return &resp, nil
}
// StopRecording 调用 DELETE /internal/v1/recordings/:idPhase B 引入)
//
// 时间预算ffmpeg 优雅退出最多 STOP_GRACE_MS=5s + media-server 自身处理 ~200ms
// 因此 timeout 需要明显 > 5s用 TimeoutMS默认 10s足够
//
// 404recording 不存在 / 已停止)映射为 ErrMediaResourceNotFound调用方一般视为幂等成功
// 不走 doCloseRequest 的重试链:录制 stop 不能重试(重试可能误伤新启动的同 id 录制)
func (h *HTTPMediaOrchestrator) StopRecording(ctx context.Context, recordingID string) (*StopRecordingResp, error) {
funcName := "service.http_media_orchestrator.StopRecording"
var resp StopRecordingResp
if err := h.doRequest(ctx, requestOptions{
method: http.MethodDelete,
path: fmt.Sprintf("/internal/v1/recordings/%s", recordingID),
body: nil,
timeoutMS: h.cfg.TimeoutMS,
funcName: funcName,
logFields: []zap.Field{
zap.String("recording_id", recordingID),
},
}, &resp); err != nil {
return nil, err
}
return &resp, nil
}
// ====== 内部工具 ======
// routerIDByRoomCode 从本地缓存反查 routerID缺失时返回 ErrMediaResourceNotFound

View File

@@ -67,6 +67,10 @@ type CreateProducerReq struct {
TransportID string `json:"transportId"`
Kind string `json:"kind"` // "audio" | "video"
RtpParameters json.RawMessage `json:"rtpParameters"` // 直接转发给 Node由其做结构校验
// AppData 客户端自定义元信息,例如 {"screen": true} 表示该 Producer 是屏幕共享流
// Phase 3 屏幕共享引入:与摄像头/麦克风 video Producer 同走 sendTransport
// 通过 appData.screen 区分,便于服务端识别并触发 meeting.screen.* 广播
AppData json.RawMessage `json:"appData,omitempty"`
}
// CreateConsumerReq 创建 Consumer 请求
@@ -78,6 +82,38 @@ type CreateConsumerReq struct {
RtpCapabilities json.RawMessage `json:"rtpCapabilities"`
}
// ====== 会议录制 RPCPhase B 引入)======
// StartRecordingReq 启动录制请求
//
// ProducerIDs 指定要录的若干 Producer由业务层决定Phase B1 通常是 host 的
// audio + video 两个Phase B3 扩展为屏幕共享 + 多人音频后会更多)
type StartRecordingReq struct {
RouterID string `json:"routerId"`
RoomCode string `json:"roomCode"`
ProducerIDs []string `json:"producerIds"`
}
// StartRecordingResp media-server 启动成功后的响应
type StartRecordingResp struct {
RecordingID string `json:"recordingId"`
OutputPath string `json:"outputPath"` // media-server 本地 mp4 路径,用于停止后回拉上传
Tracks []struct {
ProducerID string `json:"producerId"`
Kind string `json:"kind"`
} `json:"tracks"`
}
// StopRecordingResp media-server 停止录制后返回的元数据
type StopRecordingResp struct {
RecordingID string `json:"recordingId"`
OutputPath string `json:"outputPath"`
SizeBytes int64 `json:"sizeBytes"`
DurationMS int64 `json:"durationMs"`
ExitCode *int `json:"exitCode,omitempty"` // ffmpeg 退出码nil = 进程未启动到退出阶段
FailureReason string `json:"failureReason,omitempty"` // 非空表示录制最终标记为 failed
}
// MediaOrchestrator 媒体服务器编排接口(设计 §6.6 NodeClient
// Task 7 (2026-04-21) 起默认实现为 HTTPMediaOrchestrator通过 X-Internal-Token
// 调用 Node media-server 的 /internal/v1/* REST API
@@ -122,6 +158,16 @@ type MediaOrchestrator interface {
ResumeConsumer(ctx context.Context, consumerID string) error
// CloseConsumer 关闭指定 Consumer幂等
CloseConsumer(ctx context.Context, consumerID string) error
// ====== 录制Phase B 引入)======
// StartRecording 在 media-server 上启动一份会议录制
// 失败语义404 RouterNotFound / 409 producerId 不可消费等返回 ErrMediaServerError
// 成功返回 recordingIdhex32+ media-server 本地 mp4 输出路径
StartRecording(ctx context.Context, req *StartRecordingReq) (*StartRecordingResp, error)
// StopRecording 停止录制
// 幂等媒体侧已结束404返回 ErrMediaResourceNotFound调用方一般视为"已停止"
StopRecording(ctx context.Context, recordingID string) (*StopRecordingResp, error)
}
// NoopMediaOrchestrator 本地调试占位实现Task 7 后默认不再使用)
@@ -205,3 +251,18 @@ func (n *NoopMediaOrchestrator) ResumeConsumer(_ context.Context, _ string) erro
func (n *NoopMediaOrchestrator) CloseConsumer(_ context.Context, _ string) error {
return nil
}
// StartRecording 占位:返回伪造 recordingId 与空 outputPath
func (n *NoopMediaOrchestrator) StartRecording(_ context.Context, req *StartRecordingReq) (*StartRecordingResp, error) {
return &StartRecordingResp{
RecordingID: "noop-recording-" + req.RoomCode,
OutputPath: "",
}, nil
}
// StopRecording 占位:返回零大小元数据
func (n *NoopMediaOrchestrator) StopRecording(_ context.Context, recordingID string) (*StopRecordingResp, error) {
return &StopRecordingResp{
RecordingID: recordingID,
}, nil
}

View File

@@ -0,0 +1,469 @@
package service
import (
"context"
"errors"
"fmt"
"os"
"path/filepath"
"strings"
"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/config"
"github.com/echochat/backend/pkg/logs"
"github.com/minio/minio-go/v7"
"github.com/redis/go-redis/v9"
"go.uber.org/zap"
)
// MeetingRecordingService 会议录制业务服务Phase B 引入)
//
// 职责边界:
// - 鉴权(仅 host 可启停)+ 状态机DB 表 meeting_recordings
// - 选定 host 的活跃 Producer 列表 → 委托给 MediaOrchestrator.StartRecording
// - 录制结束后从 media-server 本地拉取 mp4 → 上传到 MinIO → 标记 ready
// - 通过 MeetingBroadcaster 把 started / stopped 事件广播给会议内全员
//
// 设计取舍Phase B1
// - StopRecording 内同步执行 ffmpeg 停止 + MinIO 上传:录制时长一般 ≤ 30 分钟、单文件 ≤ 数百 MB
// 在 HTTP 处理协程内同步上传可接受;后续若需支持长会议再切换为后台 worker
// - 文件路径规则recordings/{roomCode}/{recordingID}-{remoteID}.mp4
// 便于按会议聚合 + 即使 DB 与对象存储错位也能从对象 Key 反查
// - 复用全局 MinIO Bucket与图片/语音同 bucketpublic-read后续可拆分专用 bucket
type MeetingRecordingService struct {
roomDAO *dao.MeetingRoomDAO
participantDAO *dao.MeetingParticipantDAO
recordingDAO *dao.MeetingRecordingDAO
redis *redis.Client
broadcaster *MeetingBroadcaster
mediaOrchestrator MediaOrchestrator
minioClient *minio.Client
minioCfg *config.MinioConfig
}
// NewMeetingRecordingService Wire Provider
func NewMeetingRecordingService(
roomDAO *dao.MeetingRoomDAO,
participantDAO *dao.MeetingParticipantDAO,
recordingDAO *dao.MeetingRecordingDAO,
rdb *redis.Client,
broadcaster *MeetingBroadcaster,
mediaOrchestrator MediaOrchestrator,
minioClient *minio.Client,
minioCfg *config.MinioConfig,
) *MeetingRecordingService {
return &MeetingRecordingService{
roomDAO: roomDAO,
participantDAO: participantDAO,
recordingDAO: recordingDAO,
redis: rdb,
broadcaster: broadcaster,
mediaOrchestrator: mediaOrchestrator,
minioClient: minioClient,
minioCfg: minioCfg,
}
}
// 录制相关领域错误(与 MeetingService 错误并列,使用同一 errors.Is 链)
var (
// ErrRecordingAlreadyActive 同一会议已有活跃录制recording / uploading禁止再次启动
ErrRecordingAlreadyActive = errors.New("当前会议已有正在进行的录制")
// ErrRecordingNoProducers host 在 Redis 资源追踪集合内找不到任何 producer
// 通常意味着 host 尚未开启麦克风/摄像头。规则:必须至少有 1 个 producer 才允许启动录制
ErrRecordingNoProducers = errors.New("无可录制的音视频流,请先开启麦克风或摄像头")
// ErrRecordingNotFound 指定 recordingID 不存在 / 已删除
ErrRecordingNotFound = errors.New("录制记录不存在")
// ErrRecordingNotStoppable 录制不在 recording 状态,无法停止
ErrRecordingNotStoppable = errors.New("录制已结束或正在收尾,无需重复操作")
)
// StartRecording 启动录制host
//
// 流程:
// 1. 鉴权room 存在、Active、调用方=host
// 2. 防重room 内不能已有活跃录制
// 3. 收集 host 自己持有的全部 Producer ID来自 resourceTrackKey 的 set 成员)
// 4. mediaOrchestrator.StartRecording → 拿到 recordingId 与本地 outputPath
// 5. DB 落 recording 记录(写入 RemoteID + OutputPath via FailureReason 字段暂存)
// 6. 异步广播 meeting.recording.started 给所有活跃成员
//
// 返回:刚创建的 MeetingRecording 记录(业务主键 ID 即"录制 ID"
func (s *MeetingRecordingService) StartRecording(ctx context.Context, userID int64, code string) (*model.MeetingRecording, error) {
funcName := "service.meeting_recording_service.StartRecording"
room, err := s.roomDAO.GetByCode(ctx, code)
if err != nil {
return nil, err
}
if room == nil {
return nil, ErrMeetingNotFound
}
if room.Status == constants.MeetingStatusEnded {
return nil, ErrMeetingEnded
}
if room.HostID != userID {
return nil, ErrNotMeetingHost
}
// 防重:同一会议同时仅允许 1 条活跃录制
if active, err := s.recordingDAO.GetActiveByRoom(ctx, room.ID); err != nil {
return nil, err
} else if active != nil {
return nil, ErrRecordingAlreadyActive
}
// 收集 host 当前活跃的 Producer ID
producerIDs, err := s.listHostProducerIDs(ctx, code, userID)
if err != nil {
return nil, err
}
if len(producerIDs) == 0 {
return nil, ErrRecordingNoProducers
}
// 反查 Router若房间初始化阶段 Router 已被释放(极端场景),转为媒体服务不可用
routerID, ok := s.mediaOrchestrator.ResolveRouterID(code)
if !ok || routerID == "" {
logs.Warn(ctx, funcName, "未能解析 Router ID会议媒体未就绪",
zap.String("room_code", code))
return nil, ErrMediaServiceUnavailable
}
// 调 media-server
resp, err := s.mediaOrchestrator.StartRecording(ctx, &StartRecordingReq{
RouterID: routerID,
RoomCode: code,
ProducerIDs: producerIDs,
})
if err != nil {
logs.Error(ctx, funcName, "media-server 启动录制失败",
zap.String("room_code", code), zap.Int64("user_id", userID), zap.Error(err))
return nil, err
}
// 写 DB使用 FailureReason 字段暂存 outputPathstop 阶段读出来上传 MinIO 后清空)
// 这样不需要为 outputPath 加额外列Phase B1 取舍stop 时 MarkUploading 会覆盖该字段
now := time.Now()
rec := &model.MeetingRecording{
RoomID: room.ID,
RoomCode: code,
StartedBy: userID,
Status: model.MeetingRecordingStatusRecording,
RemoteID: resp.RecordingID,
FailureReason: resp.OutputPath, // 临时占位stop 时清空
StartedAt: now,
}
if err := s.recordingDAO.Create(ctx, rec); err != nil {
// DB 写失败必须回滚 media-server 上的录制,否则 ffmpeg 会一直写本地磁盘
logs.Error(ctx, funcName, "录制 DB 写入失败,回滚 media-server",
zap.String("remote_id", resp.RecordingID), zap.Error(err))
if _, stopErr := s.mediaOrchestrator.StopRecording(ctx, resp.RecordingID); stopErr != nil {
logs.Warn(ctx, funcName, "回滚停止录制失败(孤儿 ffmpeg 进程将由 media-server 自身 TTL 兜底)",
zap.String("remote_id", resp.RecordingID), zap.Error(stopErr))
}
return nil, err
}
// WS 广播detach context 避免 HTTP 取消导致广播半成品)
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventRecordingStarted, map[string]interface{}{
"recording_id": rec.ID,
"started_by": userID,
"started_at": rec.StartedAt.Unix(),
})
logs.Info(ctx, funcName, "录制已启动",
zap.Int64("recording_id", rec.ID),
zap.String("room_code", code),
zap.Int64("user_id", userID),
zap.Int("producer_count", len(producerIDs)))
return rec, nil
}
// StopRecording 停止录制host
//
// 流程:
// 1. 鉴权 + 状态校验recording 必须存在、status=recording、调用方是 host容忍 EndRoom 路径以系统 userID 调用)
// 2. media-server stop拿到本地 outputPath / size / duration / exitCode
// 3. DB MarkUploadingstatus: recording → uploading同时清空临时 FailureReason
// 4. 上传本地 mp4 到 MinIOrecordings/{roomCode}/{id}-{remoteId}.mp4
// 5. DB MarkReady + 广播 meeting.recording.stopped(status=ready, file_url)
// 6. 任一步失败MarkFailed + 广播 stopped(status=failed, reason)
func (s *MeetingRecordingService) StopRecording(ctx context.Context, userID int64, code string, recordingID int64) (*model.MeetingRecording, error) {
funcName := "service.meeting_recording_service.StopRecording"
room, err := s.roomDAO.GetByCode(ctx, code)
if err != nil {
return nil, err
}
if room == nil {
return nil, ErrMeetingNotFound
}
if room.HostID != userID {
return nil, ErrNotMeetingHost
}
rec, err := s.recordingDAO.GetByID(ctx, recordingID)
if err != nil {
return nil, err
}
if rec == nil || rec.RoomID != room.ID {
return nil, ErrRecordingNotFound
}
if rec.Status != model.MeetingRecordingStatusRecording {
return nil, ErrRecordingNotStoppable
}
// 调 media-server stop404 视为已停止 → 用本地 DB 字段降级处理)
stopResp, stopErr := s.mediaOrchestrator.StopRecording(ctx, rec.RemoteID)
if stopErr != nil {
if !errors.Is(stopErr, ErrMediaResourceNotFound) {
logs.Error(ctx, funcName, "media-server 停止录制失败",
zap.Int64("recording_id", recordingID), zap.String("remote_id", rec.RemoteID), zap.Error(stopErr))
s.markFailedAndBroadcast(ctx, room.ID, rec, "media-server 停止录制失败")
return nil, stopErr
}
logs.Warn(ctx, funcName, "media-server 录制已不存在,按本地状态降级处理",
zap.Int64("recording_id", recordingID), zap.String("remote_id", rec.RemoteID))
stopResp = &StopRecordingResp{RecordingID: rec.RemoteID, OutputPath: rec.FailureReason}
}
// 转入 uploading 状态
durationSec := int(stopResp.DurationMS / 1000)
if _, err := s.recordingDAO.MarkUploading(ctx, rec.ID, time.Now(), durationSec); err != nil {
logs.Error(ctx, funcName, "标记 uploading 失败", zap.Int64("recording_id", rec.ID), zap.Error(err))
s.markFailedAndBroadcast(ctx, room.ID, rec, "数据库状态切换失败")
return nil, err
}
outputPath := stopResp.OutputPath
if outputPath == "" {
// 容忍 media-server 旧版本 / 404 降级:用 DB 暂存的本地路径
outputPath = rec.FailureReason
}
// 上传到 MinIO
objectKey := fmt.Sprintf("recordings/%s/%d-%s.mp4", strings.ToLower(code), rec.ID, rec.RemoteID)
fileURL, sizeBytes, uploadErr := s.uploadRecording(ctx, outputPath, objectKey)
if uploadErr != nil {
logs.Error(ctx, funcName, "录制文件上传 MinIO 失败",
zap.Int64("recording_id", rec.ID), zap.String("local_path", outputPath), zap.Error(uploadErr))
s.markFailedAndBroadcast(ctx, room.ID, rec, "录制文件上传失败")
return nil, uploadErr
}
if _, err := s.recordingDAO.MarkReady(ctx, rec.ID, fileURL, objectKey, sizeBytes); err != nil {
logs.Error(ctx, funcName, "标记 ready 失败", zap.Int64("recording_id", rec.ID), zap.Error(err))
s.markFailedAndBroadcast(ctx, room.ID, rec, "数据库状态切换失败ready")
return nil, err
}
// 重新拉一份最终态
final, _ := s.recordingDAO.GetByID(ctx, rec.ID)
if final == nil {
final = rec
final.Status = model.MeetingRecordingStatusReady
final.FileURL = fileURL
final.FileObject = objectKey
final.SizeBytes = sizeBytes
final.DurationSec = durationSec
}
// 上传成功后再尝试删除本地文件(失败仅 Warn不影响业务
if outputPath != "" {
if err := os.Remove(outputPath); err != nil && !errors.Is(err, os.ErrNotExist) {
logs.Warn(ctx, funcName, "删除 media-server 本地录制文件失败(忽略)",
zap.String("local_path", outputPath), zap.Error(err))
}
}
// 广播 stopped(ready)
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventRecordingStopped, map[string]interface{}{
"recording_id": final.ID,
"status": model.MeetingRecordingStatusReady,
"file_url": fileURL,
"size_bytes": sizeBytes,
"duration_sec": durationSec,
})
logs.Info(ctx, funcName, "录制已结束并完成上传",
zap.Int64("recording_id", final.ID),
zap.String("room_code", code),
zap.Int64("size_bytes", sizeBytes),
zap.Int("duration_sec", durationSec))
return final, nil
}
// ListRecordings 拉取某会议的全部录制记录(倒序)
// 仅会议参与者可见;前端用于"录制历史"入口
func (s *MeetingRecordingService) ListRecordings(ctx context.Context, userID int64, code string) ([]model.MeetingRecording, error) {
room, err := s.roomDAO.GetByCode(ctx, code)
if err != nil {
return nil, err
}
if room == nil {
return nil, ErrMeetingNotFound
}
// host 总是可见;其他人需在 participants 表里出现过即可(包括已离会)
if room.HostID != userID {
p, err := s.participantDAO.GetByRoomAndUser(ctx, room.ID, userID)
if err != nil {
return nil, err
}
if p == nil {
return nil, ErrNotInMeeting
}
}
return s.recordingDAO.ListByRoom(ctx, room.ID)
}
// StopActiveForEndRoom 在 EndRoom / EmptyTTL 路径上由系统强制停止当前活跃录制
//
// 与对外 StopRecording 不同:
// - 跳过 host 校验(调用方是系统)
// - 失败仅 Warn 不阻断会议结束流程
// - 仍然走完整的 media-server stop + MinIO 上传 + DB 状态机
func (s *MeetingRecordingService) StopActiveForEndRoom(ctx context.Context, room *model.MeetingRoom) {
funcName := "service.meeting_recording_service.StopActiveForEndRoom"
if room == nil {
return
}
active, err := s.recordingDAO.GetActiveByRoom(ctx, room.ID)
if err != nil {
logs.Warn(ctx, funcName, "查询活跃录制失败(忽略)",
zap.Int64("room_id", room.ID), zap.Error(err))
return
}
if active == nil || active.Status != model.MeetingRecordingStatusRecording {
return
}
if _, err := s.StopRecording(ctx, room.HostID, room.RoomCode, active.ID); err != nil {
logs.Warn(ctx, funcName, "EndRoom 路径停止录制失败(已记录,定时任务会兜底)",
zap.Int64("recording_id", active.ID), zap.Error(err))
}
}
// ====== 内部工具 ======
// listHostProducerIDs 从 host 的资源追踪 set 中取出全部 "producer:<id>" 的 id 部分
//
// 选取规则Phase B1
// - 录 host 自己持有的全部 Producer含麦克风、摄像头、屏幕共享
// - 不包含 transport / consumer
// - 数组顺序无关media-server 内部会自动按 kind 处理
func (s *MeetingRecordingService) listHostProducerIDs(ctx context.Context, roomCode string, hostID int64) ([]string, error) {
key := resourceTrackKey(roomCode, hostID)
members, err := s.redis.SMembers(ctx, key).Result()
if err != nil {
return nil, err
}
out := make([]string, 0, len(members))
for _, m := range members {
if id, ok := strings.CutPrefix(m, "producer:"); ok && id != "" {
out = append(out, id)
}
}
return out, nil
}
// uploadRecording 把本地 mp4 上传到 MinIO返回可访问 URL 与文件大小
func (s *MeetingRecordingService) uploadRecording(ctx context.Context, localPath, objectKey string) (string, int64, error) {
if localPath == "" {
return "", 0, fmt.Errorf("本地录制文件路径为空")
}
stat, err := os.Stat(localPath)
if err != nil {
return "", 0, fmt.Errorf("stat 本地录制文件失败: %w", err)
}
if stat.Size() == 0 {
return "", 0, fmt.Errorf("本地录制文件为空0 byte疑似 ffmpeg 未写入任何数据")
}
f, err := os.Open(filepath.Clean(localPath))
if err != nil {
return "", 0, fmt.Errorf("打开本地录制文件失败: %w", err)
}
defer f.Close()
if _, err := s.minioClient.PutObject(ctx, s.minioCfg.Bucket, objectKey, f, stat.Size(), minio.PutObjectOptions{
ContentType: "video/mp4",
}); err != nil {
return "", 0, fmt.Errorf("PutObject 失败: %w", err)
}
scheme := "http"
if s.minioCfg.UseSSL {
scheme = "https"
}
url := fmt.Sprintf("%s://%s/%s/%s", scheme, s.minioCfg.Endpoint, s.minioCfg.Bucket, objectKey)
return url, stat.Size(), nil
}
// GetByRemoteID 反查本地录制记录(按 media-server 的 remoteId 字符串)
// webhook sink 入口前置查询,定位 Go 端主键
func (s *MeetingRecordingService) GetByRemoteID(ctx context.Context, remoteID string) (*model.MeetingRecording, error) {
return s.recordingDAO.GetByRemoteID(ctx, remoteID)
}
// HandleMediaServerFailure media-server 异常退出回调入口webhook sink
//
// 触发条件media-server 端 ffmpeg 在 Go 主动 Stop 之前自行退出,
// media-server 通过 RECORDING_FAILURE_WEBHOOK_URL 上报到本接口
//
// 行为:
// 1. 按 recordingID 查记录若已是终态ready / failed幂等返回 nil
// 2. 调 markFailedAndBroadcast 落 DB failed + 广播 meeting.recording.stopped(failed)
//
// 上层webhook controller需自行鉴权共享密钥 / 仅本地访问等)
func (s *MeetingRecordingService) HandleMediaServerFailure(ctx context.Context, recordingID int64, reason string) error {
funcName := "service.meeting_recording_service.HandleMediaServerFailure"
if recordingID <= 0 {
return fmt.Errorf("invalid recording id")
}
rec, err := s.recordingDAO.GetByID(ctx, recordingID)
if err != nil {
logs.Warn(ctx, funcName, "查询录制记录失败",
zap.Int64("recording_id", recordingID), zap.Error(err))
return err
}
if rec == nil {
logs.Warn(ctx, funcName, "录制记录不存在,忽略 webhook",
zap.Int64("recording_id", recordingID))
return nil
}
// 终态幂等
if rec.Status == model.MeetingRecordingStatusReady || rec.Status == model.MeetingRecordingStatusFailed {
logs.Info(ctx, funcName, "录制已处终态,忽略 webhook",
zap.Int64("recording_id", recordingID), zap.String("status", string(rec.Status)))
return nil
}
if reason == "" {
reason = "media-server 上报ffmpeg 异常退出"
}
s.markFailedAndBroadcast(ctx, rec.RoomID, rec, reason)
logs.Info(ctx, funcName, "已处理 media-server 失败 webhook",
zap.Int64("recording_id", recordingID), zap.String("reason", reason))
return nil
}
// markFailedAndBroadcast 标记失败并广播 stopped(failed)
// 内部容错DB / 广播任一失败仅 Warn不二次抛出
func (s *MeetingRecordingService) markFailedAndBroadcast(ctx context.Context, roomID int64, rec *model.MeetingRecording, reason string) {
if _, err := s.recordingDAO.MarkFailed(ctx, rec.ID, reason); err != nil {
logs.Warn(ctx, "service.meeting_recording_service.markFailedAndBroadcast",
"标记 failed 失败(忽略)",
zap.Int64("recording_id", rec.ID), zap.Error(err))
}
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), roomID, constants.MeetingWSEventRecordingStopped, map[string]interface{}{
"recording_id": rec.ID,
"status": model.MeetingRecordingStatusFailed,
"failure_reason": reason,
})
}

View File

@@ -80,6 +80,22 @@ type MeetingService struct {
onlineChecker OnlineChecker
mediaOrchestrator MediaOrchestrator
lifecycleSvc *MeetingLifecycleService
// recordingStopper 由 NewApp 在 Wire 注入完成后通过 SetRecordingStopper 注入
// 用 setter 而非构造函数参数:避免与 MeetingRecordingService 形成构造期循环依赖
// nil 时 EndRoom 跳过录制清理(向后兼容,生产路径始终非 nil
recordingStopper RecordingStopper
}
// RecordingStopper EndRoom / 空房 TTL 等系统路径强制停止活跃录制的钩子
// 由 *MeetingRecordingService 实现,通过 setter 注入,避免循环依赖
type RecordingStopper interface {
StopActiveForEndRoom(ctx context.Context, room *model.MeetingRoom)
}
// SetRecordingStopper 注入录制停止钩子(在 NewApp 完成 Wire 后调用)
func (s *MeetingService) SetRecordingStopper(stopper RecordingStopper) {
s.recordingStopper = stopper
}
// NewMeetingService 创建 MeetingService 实例
@@ -591,6 +607,12 @@ func (s *MeetingService) EndRoom(ctx context.Context, userID int64, code string)
_ = s.broadcaster.PublishToUser(ctx, p.UserID, constants.MeetingWSEventRoomEnded, payload)
}
// Phase B先停掉活跃录制再关 Router否则 CloseRouter 会触发录制 Consumer 异常退出,
// 录制文件可能丢失最后几秒。stopper 同步执行(含 ffmpeg 优雅退出 + MinIO 上传),失败仅 Warn 不阻断
if s.recordingStopper != nil {
s.recordingStopper.StopActiveForEndRoom(ctx, room)
}
if err := s.mediaOrchestrator.CloseRouter(ctx, code); err != nil {
logs.Warn(ctx, funcName, "关闭 mediasoup Router 失败", zap.Error(err))
}

View File

@@ -69,6 +69,39 @@ func memberStateKey(roomCode string, userID int64) string {
return fmt.Sprintf("echo:meeting:member_state:%s:%d", roomCode, userID)
}
// screenOwnerKey Redis String记录当前会议正在共享屏幕的人 + Producer ID
// 值格式:"{userId}:{producerId}";为空 / Nil 表示当前无人共享
// 单字符串而非 Hash 是因为"会议同时只允许 1 份屏幕共享",无需多字段
//
// 生命周期:
// - OnProduceStart 收到 appData.screen=true 时 SET同时广播 meeting.screen.started
// - OnProducerClose / cleanupUserResources 命中该 Producer 时 DEL同时广播 meeting.screen.stopped
// - 房间销毁路径cleanupRoomRedisResidual一并 Del
func screenOwnerKey(roomCode string) string {
return fmt.Sprintf("echo:meeting:screen_owner:%s", roomCode)
}
// formatScreenOwnerValue / parseScreenOwnerValue 统一 screenOwnerKey 的取值格式
// 复用于 set / parse 双向,避免格式漂移
func formatScreenOwnerValue(userID int64, producerID string) string {
return fmt.Sprintf("%d:%s", userID, producerID)
}
func parseScreenOwnerValue(raw string) (int64, string, bool) {
if raw == "" {
return 0, "", false
}
parts := strings.SplitN(raw, ":", 2)
if len(parts) != 2 || parts[1] == "" {
return 0, "", false
}
var uid int64
if _, err := fmt.Sscanf(parts[0], "%d", &uid); err != nil || uid <= 0 {
return 0, "", false
}
return uid, parts[1], true
}
// resourceTTL 单个用户资源追踪集合 TTL
// 设计:会议期间维持可达即可;若用户长期不活跃由断线清理接管
// Task 16 Nit常量已迁出至 constants.MeetingResourceTrackTTLSeconds此处保留计算式 wrapper 方便调用侧零改动
@@ -237,6 +270,43 @@ func (s *MeetingSignalService) OnRoomJoin(ctx context.Context, userID int64, roo
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)
s.pushExistingScreenShare(ctx, roomCode, userID)
}
// pushExistingScreenShare 后入者定向补推当前屏幕共享状态Phase 3
// 复用 meeting.screen.started 事件语义,前端共用同一 handler无需新事件类型
// 没有人在共享 / Redis 异常 / 解析失败 都静默 no-op
func (s *MeetingSignalService) pushExistingScreenShare(ctx context.Context, roomCode string, userID int64) {
funcName := "service.meeting_signal_service.pushExistingScreenShare"
raw, err := s.redis.Get(ctx, screenOwnerKey(roomCode)).Result()
if err != nil {
if err != redis.Nil {
logs.Warn(ctx, funcName, "读取屏幕共享拥有者失败(忽略)",
zap.String("room_code", roomCode), zap.Error(err))
}
return
}
ownerUID, producerID, ok := parseScreenOwnerValue(raw)
if !ok {
return
}
if ownerUID == userID {
// 自己就是共享者(理论上不会发生:刚 join 不应已是 owner跳过避免回响
return
}
if err := s.broadcaster.PublishToUser(ctx, userID, constants.MeetingWSEventScreenStarted, map[string]interface{}{
"room_code": roomCode,
"user_id": ownerUID,
"producer_id": producerID,
"existing": true,
}); err != nil {
logs.Warn(ctx, funcName, "定向推送 existing screen.started 失败",
zap.Int64("to_user", userID),
zap.Int64("owner_user", ownerUID),
zap.String("producer_id", producerID),
zap.Error(err))
}
}
// pushExistingProducers 向刚加入者定向推送房间里其他用户已产生的 producer 列表
@@ -384,7 +454,7 @@ func (s *MeetingSignalService) OnWSDisconnect(ctx context.Context, userID int64)
return
}
s.cleanupUserResources(ctx, room.RoomCode, userID)
s.cleanupUserResources(ctx, room.ID, room.RoomCode, userID)
if s.lifecycleSvc != nil && room.HostID == userID {
s.lifecycleSvc.OnHostDisconnect(ctx, room.RoomCode, userID)
@@ -405,7 +475,7 @@ func (s *MeetingSignalService) OnRoomLeave(ctx context.Context, userID int64, ro
if err != nil {
return err
}
s.cleanupUserResources(ctx, roomCode, userID)
s.cleanupUserResources(ctx, room.ID, roomCode, userID)
// P2-8 修复:使用常量 MeetingLeftReasonDisconnect避免 "ws_disconnect" 等硬编码
// 与前端 MEETING_LEFT_REASON_LABEL 字面值不一致
@@ -524,6 +594,36 @@ type ProduceStartPayload struct {
TransportID string `json:"transport_id"`
Kind string `json:"kind"` // "audio" | "video"
RtpParameters json.RawMessage `json:"rtp_parameters"`
// AppData 客户端透传给 mediasoup Producer 的自定义元信息
// Phase 3 屏幕共享引入:约定 {"screen": true} 表示该 video Producer 是屏幕分享流;
// 服务端识别后会写入 screenOwnerKey + 广播 meeting.screen.started前端将该流提升大画面。
// 服务端会强制把 userId / roomCode 注入回 appData覆盖客户端伪造尝试见 HTTPMediaOrchestrator.CreateProducer
AppData json.RawMessage `json:"app_data,omitempty"`
}
// isScreenAppData 判断客户端 app_data 是否声明了屏幕共享标识
// 容错raw 为空 / 非对象 / 字段缺失均返回 false
func isScreenAppData(raw json.RawMessage) bool {
if len(raw) == 0 {
return false
}
var m map[string]any
if err := json.Unmarshal(raw, &m); err != nil {
return false
}
v, ok := m["screen"]
if !ok {
return false
}
switch x := v.(type) {
case bool:
return x
case string:
return x == "true" || x == "1"
case float64:
return x != 0
}
return false
}
// ProduceStartResult 返回给客户端的 producerID
@@ -544,25 +644,61 @@ func (s *MeetingSignalService) OnProduceStart(ctx context.Context, userID int64,
if err := s.assertOwnsResource(ctx, payload.RoomCode, userID, "transport", payload.TransportID); err != nil {
return nil, err
}
// Phase 3识别屏幕共享意图强制 video kind 才允许(音频流不参与屏幕共享语义)
isScreen := payload.Kind == "video" && isScreenAppData(payload.AppData)
// Phase 3单会议同时仅允许 1 份屏幕共享。若已有他人在共享,直接拒绝;
// 同一用户重复请求(例如客户端重发)则容忍,由后续 SET 覆盖
if isScreen {
if existingRaw, err := s.redis.Get(ctx, screenOwnerKey(payload.RoomCode)).Result(); err == nil {
if ownerUID, _, ok := parseScreenOwnerValue(existingRaw); ok && ownerUID != userID {
return nil, fmt.Errorf("当前会议已有成员正在共享屏幕")
}
}
}
producerID, err := s.mediaOrchestrator.CreateProducer(ctx, &CreateProducerReq{
RoomCode: payload.RoomCode,
UserID: userID,
TransportID: payload.TransportID,
Kind: payload.Kind,
RtpParameters: payload.RtpParameters,
AppData: payload.AppData,
})
if err != nil {
return nil, err
}
s.trackResource(ctx, payload.RoomCode, userID, "producer", producerID)
// Phase 3屏幕共享额外维护 owner key + 触发 meeting.screen.started 广播
// 注意:先广播 producer.new保证消费者侧 Consumer 创建),再广播 screen.started标记大画面提升
if isScreen {
if err := s.redis.Set(ctx, screenOwnerKey(payload.RoomCode), formatScreenOwnerValue(userID, producerID), resourceTTL).Err(); err != nil {
logs.Warn(ctx, "service.meeting_signal_service.OnProduceStart", "写入屏幕共享拥有者失败(不阻断业务)",
zap.String("room_code", payload.RoomCode),
zap.Int64("user_id", userID),
zap.String("producer_id", producerID),
zap.Error(err))
}
}
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,
"screen": isScreen,
}, userID)
if isScreen {
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventScreenStarted, map[string]interface{}{
"room_code": payload.RoomCode,
"user_id": userID,
"producer_id": producerID,
}) // 不排除任何人:共享者本人也接收,便于前端统一更新大画面 UI
}
return &ProduceStartResult{ProducerID: producerID}, nil
}
@@ -653,6 +789,9 @@ func (s *MeetingSignalService) OnProducerClose(ctx context.Context, userID int64
}
s.untrackResource(ctx, payload.RoomCode, userID, "producer", payload.ProducerID)
// Phase 3若关闭的恰好是当前屏幕共享 Producer则一并清 screenOwnerKey + 广播 screen.stopped
s.maybeReleaseScreenOwner(ctx, room.ID, payload.RoomCode, userID, payload.ProducerID)
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberProducerNew, map[string]interface{}{
"room_code": payload.RoomCode,
"user_id": userID,
@@ -662,11 +801,57 @@ func (s *MeetingSignalService) OnProducerClose(ctx context.Context, userID int64
return nil
}
// maybeReleaseScreenOwner 检查并释放屏幕共享拥有者状态
// 当 producerID 恰好等于 screenOwnerKey 当前值的 producer 段时,认为这次关闭是屏幕分享停止:
// - DEL screenOwnerKey
// - 广播 meeting.screen.stopped 给所有活跃成员(含共享者本人)
//
// 非屏幕 Producer 或 key 已被他人覆盖时静默 no-op保证幂等
func (s *MeetingSignalService) maybeReleaseScreenOwner(ctx context.Context, roomID int64, roomCode string, userID int64, producerID string) {
funcName := "service.meeting_signal_service.maybeReleaseScreenOwner"
key := screenOwnerKey(roomCode)
raw, err := s.redis.Get(ctx, key).Result()
if err != nil {
// redis.Nil无人共享属正常路径其他错误仅 Warn 不阻塞
if err != redis.Nil {
logs.Warn(ctx, funcName, "读取屏幕共享拥有者失败(忽略)",
zap.String("room_code", roomCode), zap.Error(err))
}
return
}
ownerUID, ownerProducerID, ok := parseScreenOwnerValue(raw)
if !ok || ownerProducerID != producerID {
return
}
// 防御性日志:理论上 ownerUID 必然等于 userIDtrackResource 已校验归属)
if ownerUID != userID {
logs.Warn(ctx, funcName, "屏幕共享拥有者与关闭者不一致(仍按停止处理)",
zap.Int64("owner_uid", ownerUID),
zap.Int64("closer_uid", userID),
zap.String("producer_id", producerID))
}
if err := s.redis.Del(ctx, key).Err(); err != nil {
logs.Warn(ctx, funcName, "删除屏幕共享拥有者失败(忽略)",
zap.String("room_code", roomCode), zap.Error(err))
}
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), roomID, constants.MeetingWSEventScreenStopped, map[string]interface{}{
"room_code": roomCode,
"user_id": ownerUID,
"producer_id": producerID,
})
}
// ========== 资源清理 ==========
// cleanupUserResources 批量关闭指定用户在某会议的所有媒体资源
// WS 断开、主动离会、被踢时使用;依赖 Redis 集合中追踪的资源 ID
func (s *MeetingSignalService) cleanupUserResources(ctx context.Context, roomCode string, userID int64) {
//
// Phase 3 屏幕共享:调用方需提供 roomID便于在该用户恰好是屏幕共享者时触发
// meeting.screen.stopped 广播。roomID==0 时跳过屏幕广播(仅清 Redis保持向后兼容。
func (s *MeetingSignalService) cleanupUserResources(ctx context.Context, roomID int64, roomCode string, userID int64) {
funcName := "service.meeting_signal_service.cleanupUserResources"
key := resourceTrackKey(roomCode, userID)
@@ -687,6 +872,11 @@ func (s *MeetingSignalService) cleanupUserResources(ctx context.Context, roomCod
switch kind {
case "producer":
_ = s.mediaOrchestrator.CloseProducer(ctx, id)
// Phase 3若该 Producer 是当前屏幕共享流,释放 owner key 并广播 screen.stopped
// 复用 OnProducerClose 同款幂等逻辑roomID==0 跳过避免无效广播
if roomID > 0 {
s.maybeReleaseScreenOwner(ctx, roomID, roomCode, userID, id)
}
case "consumer":
_ = s.mediaOrchestrator.CloseConsumer(ctx, id)
case "transport":
@@ -714,7 +904,14 @@ func (s *MeetingSignalService) cleanupUserResources(ctx context.Context, roomCod
// - 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 {
if rdb == nil {
return
}
// Phase 3会议销毁路径必删 screenOwnerKey即使无活跃用户也清理
// 单独 Del 不依赖 pipeline因为会议销毁是低频路径多一次 RTT 可接受
_ = rdb.Del(ctx, screenOwnerKey(roomCode)).Err()
if len(userIDs) == 0 {
return
}
pipe := rdb.Pipeline()