From ac8ec4c903623e4f301ab725ceabf04c6a87dd1f Mon Sep 17 00:00:00 2001 From: duoaohui <928970622@qq.com> Date: Mon, 18 May 2026 19:42:55 +0800 Subject: [PATCH] =?UTF-8?q?=E8=A7=86=E9=A2=91=E4=BC=9A=E8=AE=AE=E4=BF=9D?= =?UTF-8?q?=E5=AD=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- admin/src/api/meeting.js | 56 ++ admin/src/router/index.js | 6 + admin/src/views/layout/index.vue | 8 + admin/src/views/meeting/list.vue | 654 +++++++++++++++ .../controller/meeting_manage_controller.go | 89 ++ .../app/admin/dao/meeting_manage_dao.go | 286 +++++++ backend/go-service/app/admin/provider.go | 4 + backend/go-service/app/admin/router.go | 7 + .../admin/service/meeting_manage_service.go | 268 ++++++ backend/go-service/app/constants/meeting.go | 18 + backend/go-service/app/dto/admin_dto.go | 96 +++ .../meeting_recording_controller.go | 214 +++++ .../app/meeting/dao/meeting_recording_dao.go | 174 ++++ .../app/meeting/model/meeting_recording.go | 44 + backend/go-service/app/meeting/provider.go | 3 + backend/go-service/app/meeting/router.go | 18 +- .../service/http_media_orchestrator.go | 81 +- .../app/meeting/service/interfaces.go | 61 ++ .../service/meeting_recording_service.go | 469 +++++++++++ .../app/meeting/service/meeting_service.go | 22 + .../meeting/service/meeting_signal_service.go | 205 ++++- backend/go-service/app/provider/provider.go | 17 +- backend/go-service/app/provider/wire_gen.go | 8 +- backend/go-service/cmd/server/main.go | 3 + backend/go-service/router/router.go | 4 +- docs/plans/2026-05-18-recording-progress.md | 107 +++ frontend/src/api/meeting.js | 37 +- .../src/components/meeting/MeetingToolbar.vue | 78 +- frontend/src/constants/meeting.js | 27 +- frontend/src/pages.json | 12 + frontend/src/pages/auth/sso.vue | 71 ++ frontend/src/pages/meeting/recordings.vue | 271 ++++++ frontend/src/pages/meeting/room.vue | 294 ++++++- frontend/src/store/meeting.js | 359 +++++++- media-server/Dockerfile | 4 +- media-server/src/app.ts | 2 + media-server/src/routes/recording.route.ts | 44 + media-server/src/schemas/recording.schema.ts | 25 + .../src/services/recording.service.ts | 781 ++++++++++++++++++ media-server/src/utils/port-pool.ts | 112 +++ media-server/tests/utils/port-pool.spec.ts | 69 ++ 41 files changed, 5069 insertions(+), 39 deletions(-) create mode 100644 admin/src/api/meeting.js create mode 100644 admin/src/views/meeting/list.vue create mode 100644 backend/go-service/app/admin/controller/meeting_manage_controller.go create mode 100644 backend/go-service/app/admin/dao/meeting_manage_dao.go create mode 100644 backend/go-service/app/admin/service/meeting_manage_service.go create mode 100644 backend/go-service/app/meeting/controller/meeting_recording_controller.go create mode 100644 backend/go-service/app/meeting/dao/meeting_recording_dao.go create mode 100644 backend/go-service/app/meeting/model/meeting_recording.go create mode 100644 backend/go-service/app/meeting/service/meeting_recording_service.go create mode 100644 docs/plans/2026-05-18-recording-progress.md create mode 100644 frontend/src/pages/auth/sso.vue create mode 100644 frontend/src/pages/meeting/recordings.vue create mode 100644 media-server/src/routes/recording.route.ts create mode 100644 media-server/src/schemas/recording.schema.ts create mode 100644 media-server/src/services/recording.service.ts create mode 100644 media-server/src/utils/port-pool.ts create mode 100644 media-server/tests/utils/port-pool.spec.ts diff --git a/admin/src/api/meeting.js b/admin/src/api/meeting.js new file mode 100644 index 0000000..409c10a --- /dev/null +++ b/admin/src/api/meeting.js @@ -0,0 +1,56 @@ +/** + * 会议管理 API 模块(Phase B) + * + * 对应后端路由:/api/v1/admin/meetings + * 所有接口需要 JWT + admin 角色权限 + * + * 接口列表: + * - GET /api/v1/admin/meetings 会议列表(分页 + 多条件筛选) + * - GET /api/v1/admin/meetings/:id 会议详情(含参与者 + 录制列表) + * - GET /api/v1/admin/meetings/stats 会议统计(卡片) + */ +import request from '@/utils/request' + +/** + * 获取会议列表 + * @param {Object} params + * @param {number} [params.page] + * @param {number} [params.page_size] + * @param {string} [params.keyword] - 模糊匹配 title / room_code + * @param {number} [params.status] - 0=未开始 1=进行中 2=已结束 + * @param {number} [params.host_id] - 主持人 ID + * @param {boolean} [params.has_recording] - true=仅含录制的会议 + * @param {string} [params.start_time] - YYYY-MM-DD + * @param {string} [params.end_time] - YYYY-MM-DD + * @returns {Promise<{data:{total:number,list:Array,page:number,page_size:number}}>} + */ +export function getMeetingList(params) { + return request({ + url: '/api/v1/admin/meetings', + method: 'get', + params + }) +} + +/** + * 获取会议详情(含参与者 + 录制列表) + * @param {number} id - 会议主键 ID + * @returns {Promise<{data:Object}>} + */ +export function getMeetingDetail(id) { + return request({ + url: `/api/v1/admin/meetings/${id}`, + method: 'get' + }) +} + +/** + * 获取会议统计(顶栏卡片) + * @returns {Promise<{data:{total_count,active_count,today_count,recording_count}}>} + */ +export function getMeetingStats() { + return request({ + url: '/api/v1/admin/meetings/stats', + method: 'get' + }) +} diff --git a/admin/src/router/index.js b/admin/src/router/index.js index 5eb25fd..5fc0152 100644 --- a/admin/src/router/index.js +++ b/admin/src/router/index.js @@ -74,6 +74,12 @@ const routes = [ name: 'MessageStats', component: () => import('@/views/message/stats.vue'), meta: { title: '消息统计' } + }, + { + path: 'meeting/list', + name: 'MeetingList', + component: () => import('@/views/meeting/list.vue'), + meta: { title: '会议记录' } } ] } diff --git a/admin/src/views/layout/index.vue b/admin/src/views/layout/index.vue index 3595609..39a1680 100644 --- a/admin/src/views/layout/index.vue +++ b/admin/src/views/layout/index.vue @@ -65,6 +65,14 @@ 消息统计 + + + 会议记录 + + diff --git a/admin/src/views/meeting/list.vue b/admin/src/views/meeting/list.vue new file mode 100644 index 0000000..ca07198 --- /dev/null +++ b/admin/src/views/meeting/list.vue @@ -0,0 +1,654 @@ + + + + + + diff --git a/backend/go-service/app/admin/controller/meeting_manage_controller.go b/backend/go-service/app/admin/controller/meeting_manage_controller.go new file mode 100644 index 0000000..093b23d --- /dev/null +++ b/backend/go-service/app/admin/controller/meeting_manage_controller.go @@ -0,0 +1,89 @@ +// Package controller 提供 admin 模块的 HTTP 入口(会议管理) +package controller + +import ( + "errors" + "strconv" + + "github.com/echochat/backend/app/admin/service" + "github.com/echochat/backend/app/dto" + "github.com/echochat/backend/pkg/logs" + "github.com/echochat/backend/pkg/utils" + "github.com/gin-gonic/gin" + "go.uber.org/zap" +) + +// MeetingManageController 管理端会议管理 HTTP 控制器 +// +// 路由(在 admin/router.go 中挂载到 /api/v1/admin/meetings 下,JWT + admin 角色): +// - GET /api/v1/admin/meetings 分页列表 + 多条件筛选 +// - GET /api/v1/admin/meetings/stats 会议管理摘要统计 +// - GET /api/v1/admin/meetings/:id 会议详情(含参与者、录制列表) +type MeetingManageController struct { + svc *service.MeetingManageService +} + +// NewMeetingManageController Wire Provider +func NewMeetingManageController(svc *service.MeetingManageService) *MeetingManageController { + return &MeetingManageController{svc: svc} +} + +// GetMeetingList GET /api/v1/admin/meetings +func (ctl *MeetingManageController) GetMeetingList(c *gin.Context) { + ctx := c.Request.Context() + + var req dto.AdminMeetingListRequest + if err := c.ShouldBindQuery(&req); err != nil { + utils.ResponseBadRequest(c, "参数格式错误") + return + } + + resp, err := ctl.svc.GetMeetingList(ctx, &req) + if err != nil { + ctl.handleError(c, err, "获取会议列表失败") + return + } + utils.ResponseSuccess(c, resp) +} + +// GetMeetingDetail GET /api/v1/admin/meetings/:id +func (ctl *MeetingManageController) GetMeetingDetail(c *gin.Context) { + ctx := c.Request.Context() + + id, err := strconv.ParseInt(c.Param("id"), 10, 64) + if err != nil || id <= 0 { + utils.ResponseBadRequest(c, "会议 ID 非法") + return + } + + resp, err := ctl.svc.GetMeetingDetail(ctx, id) + if err != nil { + ctl.handleError(c, err, "获取会议详情失败") + return + } + utils.ResponseSuccess(c, resp) +} + +// GetMeetingStats GET /api/v1/admin/meetings/stats +func (ctl *MeetingManageController) GetMeetingStats(c *gin.Context) { + ctx := c.Request.Context() + + resp, err := ctl.svc.GetStats(ctx) + if err != nil { + ctl.handleError(c, err, "获取会议统计失败") + return + } + utils.ResponseSuccess(c, resp) +} + +// handleError 统一业务错误映射 +func (ctl *MeetingManageController) handleError(c *gin.Context, err error, fallbackMsg string) { + switch { + case errors.Is(err, service.ErrMeetingNotFound): + utils.ResponseNotFound(c, err.Error()) + default: + logs.Warn(c.Request.Context(), "controller.meeting_manage_controller.handleError", + fallbackMsg, zap.Error(err)) + utils.ResponseError(c, fallbackMsg) + } +} diff --git a/backend/go-service/app/admin/dao/meeting_manage_dao.go b/backend/go-service/app/admin/dao/meeting_manage_dao.go new file mode 100644 index 0000000..0e0b0f3 --- /dev/null +++ b/backend/go-service/app/admin/dao/meeting_manage_dao.go @@ -0,0 +1,286 @@ +// Package dao 提供 admin 模块的数据库访问操作(会议管理) +package dao + +import ( + "context" + "errors" + "fmt" + "time" + + "github.com/echochat/backend/app/dto" + meetingModel "github.com/echochat/backend/app/meeting/model" + "github.com/echochat/backend/pkg/logs" + "go.uber.org/zap" + "gorm.io/gorm" +) + +// MeetingManageDAO 管理端会议数据访问对象 +// +// 设计要点: +// - 列表查询走聚合 SQL:一次取 meeting_rooms + host 昵称 + 参与者数 + 录制数,避免 N+1 +// - 仅做"查询/统计"层操作,写操作(结束会议、删除录制)放在 service 层并复用既有路径 +// - 与 MessageManageDAO 保持一致的 ListXxx + 统计方法风格 +type MeetingManageDAO struct { + db *gorm.DB +} + +// NewMeetingManageDAO 创建 MeetingManageDAO 实例 +func NewMeetingManageDAO(db *gorm.DB) *MeetingManageDAO { + return &MeetingManageDAO{db: db} +} + +// MeetingListRow 列表行的聚合结果(DAO 内部结构,不外泄) +// +// 通过一次带 LEFT JOIN + GROUP BY 的查询拿齐: +// - 主表 meeting_rooms 全部字段 +// - 主持人昵称 / 头像(auth_users) +// - 参与者数(meeting_participants 去重 user_id) +// - 录制数(meeting_recordings 总条目) +// +// 注意:使用 GORM Raw + Scan,绕开 model 主键自动展开,方便聚合 +type MeetingListRow struct { + ID int64 `gorm:"column:id"` + RoomCode string `gorm:"column:room_code"` + Title string `gorm:"column:title"` + Type int `gorm:"column:type"` + Status int `gorm:"column:status"` + HostID int64 `gorm:"column:host_id"` + HostNickname string `gorm:"column:host_nickname"` + HostAvatar string `gorm:"column:host_avatar"` + MaxMembers int `gorm:"column:max_members"` + ParticipantCount int `gorm:"column:participant_count"` + RecordingCount int `gorm:"column:recording_count"` + ScheduledAt *time.Time `gorm:"column:scheduled_at"` + StartedAt *time.Time `gorm:"column:started_at"` + EndedAt *time.Time `gorm:"column:ended_at"` + EndedReason *string `gorm:"column:ended_reason"` + CreatedAt time.Time `gorm:"column:created_at"` +} + +// buildListWhere 把 AdminMeetingListRequest 拼成 WHERE 条件,返回 SQL 片段与参数 +// 复用给 ListMeetings + CountMeetings 两个查询,保证条件一致 +func (d *MeetingManageDAO) buildListWhere(req *dto.AdminMeetingListRequest) (string, []any) { + conditions := []string{"1=1"} + args := []any{} + + if req.Keyword != "" { + conditions = append(conditions, "(mr.title ILIKE ? OR mr.room_code ILIKE ?)") + kw := fmt.Sprintf("%%%s%%", req.Keyword) + args = append(args, kw, kw) + } + if req.Status != nil { + conditions = append(conditions, "mr.status = ?") + args = append(args, *req.Status) + } + if req.HostID != nil { + conditions = append(conditions, "mr.host_id = ?") + args = append(args, *req.HostID) + } + if req.StartTime != "" { + if t, err := time.Parse("2006-01-02", req.StartTime); err == nil { + conditions = append(conditions, "mr.created_at >= ?") + args = append(args, t) + } + } + if req.EndTime != "" { + if t, err := time.Parse("2006-01-02", req.EndTime); err == nil { + conditions = append(conditions, "mr.created_at < ?") + args = append(args, t.AddDate(0, 0, 1)) + } + } + // HasRecording 通过 HAVING 在外层处理,这里不拼接 + + where := "" + for i, c := range conditions { + if i == 0 { + where = c + } else { + where = where + " AND " + c + } + } + return where, args +} + +// ListMeetings 分页查询会议列表 +// +// 聚合 SQL 设计: +// - LEFT JOIN auth_users:拿主持人昵称(host 可能已删除时为空,不影响行) +// - LEFT JOIN meeting_participants:聚合参与者数(COUNT DISTINCT user_id) +// - LEFT JOIN meeting_recordings:聚合录制数(COUNT recording.id) +// - HasRecording 通过 HAVING recording_count > 0 二次过滤 +// +// 排序:created_at DESC(最新会议在前) +func (d *MeetingManageDAO) ListMeetings(ctx context.Context, req *dto.AdminMeetingListRequest) ([]MeetingListRow, int64, error) { + funcName := "dao.meeting_manage_dao.ListMeetings" + logs.Debug(ctx, funcName, "查询会议列表") + + where, args := d.buildListWhere(req) + + page := req.Page + if page <= 0 { + page = 1 + } + pageSize := req.PageSize + if pageSize <= 0 { + pageSize = 20 + } + if pageSize > 100 { + pageSize = 100 + } + + // 子查询聚合 participant_count / recording_count,避免笛卡尔积放大 COUNT + // 这里用相关子查询而非 JOIN+GROUP BY,便于 HAVING 与 COUNT(*) 分页对齐 + baseSelect := ` + SELECT + mr.id, mr.room_code, mr.title, mr.type, mr.status, mr.host_id, + COALESCE(u.nickname, '') AS host_nickname, + COALESCE(u.avatar, '') AS host_avatar, + mr.max_members, + COALESCE((SELECT COUNT(DISTINCT mp.user_id) FROM meeting_participants mp WHERE mp.room_id = mr.id), 0) AS participant_count, + COALESCE((SELECT COUNT(*) FROM meeting_recordings mrec WHERE mrec.room_id = mr.id), 0) AS recording_count, + mr.scheduled_at, mr.started_at, mr.ended_at, mr.ended_reason, mr.created_at + FROM meeting_rooms mr + LEFT JOIN auth_users u ON u.id = mr.host_id + WHERE ` + where + + // HasRecording 过滤:包成外层 SELECT * FROM (...) WHERE recording_count > 0 + listSQL := baseSelect + countSQL := `SELECT COUNT(*) FROM (` + baseSelect + `) AS sub` + if req.HasRecording != nil && *req.HasRecording { + listSQL = `SELECT * FROM (` + baseSelect + `) AS sub WHERE recording_count > 0` + countSQL = `SELECT COUNT(*) FROM (` + listSQL + `) AS sub2` + } + listSQL += ` ORDER BY created_at DESC OFFSET ? LIMIT ?` + + var total int64 + if err := d.db.WithContext(ctx).Raw(countSQL, args...).Scan(&total).Error; err != nil { + logs.Error(ctx, funcName, "统计会议总数失败", zap.Error(err)) + return nil, 0, err + } + + pageArgs := append([]any{}, args...) + pageArgs = append(pageArgs, (page-1)*pageSize, pageSize) + + var rows []MeetingListRow + if err := d.db.WithContext(ctx).Raw(listSQL, pageArgs...).Scan(&rows).Error; err != nil { + logs.Error(ctx, funcName, "查询会议列表失败", zap.Error(err)) + return nil, 0, err + } + + return rows, total, nil +} + +// GetMeetingByID 取单条会议聚合行(含主持人昵称 + 聚合计数) +// 供详情接口复用 MeetingListRow 投影,避免重复定义结构 +func (d *MeetingManageDAO) GetMeetingByID(ctx context.Context, id int64) (*MeetingListRow, error) { + funcName := "dao.meeting_manage_dao.GetMeetingByID" + sql := ` + SELECT + mr.id, mr.room_code, mr.title, mr.type, mr.status, mr.host_id, + COALESCE(u.nickname, '') AS host_nickname, + COALESCE(u.avatar, '') AS host_avatar, + mr.max_members, + COALESCE((SELECT COUNT(DISTINCT mp.user_id) FROM meeting_participants mp WHERE mp.room_id = mr.id), 0) AS participant_count, + COALESCE((SELECT COUNT(*) FROM meeting_recordings mrec WHERE mrec.room_id = mr.id), 0) AS recording_count, + mr.scheduled_at, mr.started_at, mr.ended_at, mr.ended_reason, mr.created_at + FROM meeting_rooms mr + LEFT JOIN auth_users u ON u.id = mr.host_id + WHERE mr.id = ? + LIMIT 1 + ` + var row MeetingListRow + err := d.db.WithContext(ctx).Raw(sql, id).Scan(&row).Error + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + logs.Error(ctx, funcName, "查询会议详情失败", zap.Int64("id", id), zap.Error(err)) + return nil, err + } + if row.ID == 0 { + return nil, nil + } + return &row, nil +} + +// ListParticipantsDetail 拉取会议的全部参与者(带用户昵称/头像) +// +// 按 joined_at ASC 排序,主持人先入会后退出时也按时间顺序展示,便于核对入会轨迹 +func (d *MeetingManageDAO) ListParticipantsDetail(ctx context.Context, roomID int64) ([]ParticipantDetailRow, error) { + funcName := "dao.meeting_manage_dao.ListParticipantsDetail" + sql := ` + SELECT + mp.user_id, mp.role, mp.joined_at, mp.left_at, mp.left_reason, mp.duration, + COALESCE(u.nickname, '') AS nickname, + COALESCE(u.avatar_url, '') AS avatar + FROM meeting_participants mp + LEFT JOIN auth_users u ON u.id = mp.user_id + WHERE mp.room_id = ? + ORDER BY mp.joined_at ASC, mp.id ASC + ` + var rows []ParticipantDetailRow + if err := d.db.WithContext(ctx).Raw(sql, roomID).Scan(&rows).Error; err != nil { + logs.Error(ctx, funcName, "查询会议参与者失败", zap.Int64("room_id", roomID), zap.Error(err)) + return nil, err + } + return rows, nil +} + +// ParticipantDetailRow 参与者明细投影 +type ParticipantDetailRow struct { + UserID int64 `gorm:"column:user_id"` + Role int `gorm:"column:role"` + JoinedAt time.Time `gorm:"column:joined_at"` + LeftAt *time.Time `gorm:"column:left_at"` + LeftReason *string `gorm:"column:left_reason"` + Duration int `gorm:"column:duration"` + Nickname string `gorm:"column:nickname"` + Avatar string `gorm:"column:avatar"` +} + +// ListRecordingsByRoom 拉取会议下的全部录制(admin 视角,按 started_at DESC) +// +// 复用 meeting/model.MeetingRecording,不另起 DTO;service 层再投影成 AdminMeetingRecordingDTO +func (d *MeetingManageDAO) ListRecordingsByRoom(ctx context.Context, roomID int64) ([]meetingModel.MeetingRecording, error) { + var list []meetingModel.MeetingRecording + err := d.db.WithContext(ctx). + Where("room_id = ?", roomID). + Order("started_at DESC, id DESC"). + Find(&list).Error + return list, err +} + +// CountTotal 历史会议总数 +func (d *MeetingManageDAO) CountTotal(ctx context.Context) (int64, error) { + var count int64 + err := d.db.WithContext(ctx).Model(&meetingModel.MeetingRoom{}).Count(&count).Error + return count, err +} + +// CountActive 当前进行中(status=1)的会议数 +func (d *MeetingManageDAO) CountActive(ctx context.Context) (int64, error) { + var count int64 + err := d.db.WithContext(ctx). + Model(&meetingModel.MeetingRoom{}). + Where("status = ?", 1). + Count(&count).Error + return count, err +} + +// CountToday 今日创建的会议数(按 created_at) +func (d *MeetingManageDAO) CountToday(ctx context.Context) (int64, error) { + var count int64 + today := time.Now().Truncate(24 * time.Hour) + err := d.db.WithContext(ctx). + Model(&meetingModel.MeetingRoom{}). + Where("created_at >= ?", today). + Count(&count).Error + return count, err +} + +// CountRecordings 录制总数 +func (d *MeetingManageDAO) CountRecordings(ctx context.Context) (int64, error) { + var count int64 + err := d.db.WithContext(ctx).Model(&meetingModel.MeetingRecording{}).Count(&count).Error + return count, err +} diff --git a/backend/go-service/app/admin/provider.go b/backend/go-service/app/admin/provider.go index dcd8ab8..fad207d 100644 --- a/backend/go-service/app/admin/provider.go +++ b/backend/go-service/app/admin/provider.go @@ -23,4 +23,8 @@ var AdminSet = wire.NewSet( dao.NewMessageManageDAO, service.NewMessageManageService, controller.NewMessageManageController, + // Phase B:会议管理(与录制审计共用一套接口) + dao.NewMeetingManageDAO, + service.NewMeetingManageService, + controller.NewMeetingManageController, ) diff --git a/backend/go-service/app/admin/router.go b/backend/go-service/app/admin/router.go index ebb3580..30cdde3 100644 --- a/backend/go-service/app/admin/router.go +++ b/backend/go-service/app/admin/router.go @@ -17,6 +17,7 @@ func RegisterRoutes( contactManageCtrl *controller.ContactManageController, groupManageCtrl *controller.GroupManageController, msgManageCtrl *controller.MessageManageController, + meetingManageCtrl *controller.MeetingManageController, jwtAuth gin.HandlerFunc, ) { // 管理端路由组:JWT 认证 + admin/super_admin 角色检查 @@ -52,5 +53,11 @@ func RegisterRoutes( adminGroup.GET("/messages/:id", msgManageCtrl.GetMessageDetail) adminGroup.DELETE("/messages/:id", msgManageCtrl.DeleteMessage) adminGroup.PUT("/messages/:id/recall", msgManageCtrl.RecallMessage) + + // 会议管理(Phase B:含录制审计) + // 注意 /stats 必须放在 /:id 之前,否则会被 :id 路由吞掉(gin trie 优先级问题) + adminGroup.GET("/meetings/stats", meetingManageCtrl.GetMeetingStats) + adminGroup.GET("/meetings", meetingManageCtrl.GetMeetingList) + adminGroup.GET("/meetings/:id", meetingManageCtrl.GetMeetingDetail) } } diff --git a/backend/go-service/app/admin/service/meeting_manage_service.go b/backend/go-service/app/admin/service/meeting_manage_service.go new file mode 100644 index 0000000..72a1875 --- /dev/null +++ b/backend/go-service/app/admin/service/meeting_manage_service.go @@ -0,0 +1,268 @@ +// Package service 提供 admin 模块的业务服务(会议管理) +package service + +import ( + "context" + "errors" + "time" + + adminDAO "github.com/echochat/backend/app/admin/dao" + "github.com/echochat/backend/app/dto" + meetingModel "github.com/echochat/backend/app/meeting/model" + "github.com/echochat/backend/pkg/logs" + "go.uber.org/zap" +) + +// ErrMeetingNotFound 会议记录不存在 +var ErrMeetingNotFound = errors.New("会议记录不存在") + +// MeetingManageService 管理端会议查询服务 +// +// 仅做"读"路径:列表 / 详情 / 统计。 +// "强制结束会议" 这种写操作直接复用 meeting/service.MeetingService.EndRoom 即可(admin_force 路径), +// 暂不在本服务里重复包装,避免与现有 RBAC/中间件冲突。 +type MeetingManageService struct { + dao *adminDAO.MeetingManageDAO +} + +// NewMeetingManageService Wire Provider +func NewMeetingManageService(dao *adminDAO.MeetingManageDAO) *MeetingManageService { + return &MeetingManageService{dao: dao} +} + +// ---------- public methods ---------- + +// GetMeetingList 查询会议列表(分页 + 多条件) +func (s *MeetingManageService) GetMeetingList(ctx context.Context, req *dto.AdminMeetingListRequest) (*dto.AdminMeetingListResponse, error) { + funcName := "service.meeting_manage_service.GetMeetingList" + logs.Debug(ctx, funcName, "获取会议列表") + + if req.Page <= 0 { + req.Page = 1 + } + if req.PageSize <= 0 { + req.PageSize = 20 + } + + rows, total, err := s.dao.ListMeetings(ctx, req) + if err != nil { + return nil, err + } + + list := make([]dto.AdminMeetingDTO, 0, len(rows)) + for i := range rows { + list = append(list, rowToDTO(&rows[i])) + } + + return &dto.AdminMeetingListResponse{ + Total: total, + List: list, + Page: req.Page, + PageSize: req.PageSize, + }, nil +} + +// GetMeetingDetail 查询会议详情(含参与者 + 录制) +func (s *MeetingManageService) GetMeetingDetail(ctx context.Context, id int64) (*dto.AdminMeetingDetailResponse, error) { + funcName := "service.meeting_manage_service.GetMeetingDetail" + logs.Debug(ctx, funcName, "获取会议详情", zap.Int64("id", id)) + + row, err := s.dao.GetMeetingByID(ctx, id) + if err != nil { + return nil, err + } + if row == nil { + return nil, ErrMeetingNotFound + } + + participants, err := s.dao.ListParticipantsDetail(ctx, row.ID) + if err != nil { + logs.Warn(ctx, funcName, "拉取参与者明细失败(继续返回详情)", + zap.Int64("room_id", row.ID), zap.Error(err)) + participants = nil + } + recordings, err := s.dao.ListRecordingsByRoom(ctx, row.ID) + if err != nil { + logs.Warn(ctx, funcName, "拉取录制列表失败(继续返回详情)", + zap.Int64("room_id", row.ID), zap.Error(err)) + recordings = nil + } + + pDTOs := make([]dto.AdminMeetingParticipantDTO, 0, len(participants)) + for i := range participants { + pDTOs = append(pDTOs, participantToDTO(&participants[i])) + } + rDTOs := make([]dto.AdminMeetingRecordingDTO, 0, len(recordings)) + for i := range recordings { + rDTOs = append(rDTOs, recordingToDTO(&recordings[i])) + } + + return &dto.AdminMeetingDetailResponse{ + AdminMeetingDTO: rowToDTO(row), + Participants: pDTOs, + Recordings: rDTOs, + }, nil +} + +// GetStats 获取会议管理摘要统计(卡片) +func (s *MeetingManageService) GetStats(ctx context.Context) (*dto.AdminMeetingStatsResponse, error) { + funcName := "service.meeting_manage_service.GetStats" + logs.Debug(ctx, funcName, "获取会议统计") + + total, err := s.dao.CountTotal(ctx) + if err != nil { + return nil, err + } + active, err := s.dao.CountActive(ctx) + if err != nil { + return nil, err + } + today, err := s.dao.CountToday(ctx) + if err != nil { + return nil, err + } + recordings, err := s.dao.CountRecordings(ctx) + if err != nil { + return nil, err + } + + return &dto.AdminMeetingStatsResponse{ + TotalCount: total, + ActiveCount: active, + TodayCount: today, + RecordingCount: recordings, + }, nil +} + +// ---------- DTO 投影辅助 ---------- + +const datetimeLayout = "2006-01-02 15:04:05" + +// formatTimePtr 把 *time.Time 转成展示字符串;nil 返回空串 +func formatTimePtr(t *time.Time) string { + if t == nil { + return "" + } + return t.Format(datetimeLayout) +} + +// formatStrPtr 把 *string 转成普通 string;nil 返回空串 +func formatStrPtr(s *string) string { + if s == nil { + return "" + } + return *s +} + +// typeLabel 1=即时 2=预约 +func typeLabel(t int) string { + switch t { + case 1: + return "即时会议" + case 2: + return "预约会议" + default: + return "未知" + } +} + +// statusLabel 0=未开始 1=进行中 2=已结束(与 meeting_rooms.status 对齐) +func statusLabel(s int) string { + switch s { + case 0: + return "未开始" + case 1: + return "进行中" + case 2: + return "已结束" + default: + return "未知" + } +} + +// roleLabel 0=普通 1=主持人 2=联合主持人 +func roleLabel(r int) string { + switch r { + case 1: + return "主持人" + case 2: + return "联合主持人" + default: + return "参会人" + } +} + +// computeDuration 计算会议时长(秒) +// - 已结束:ended_at - started_at +// - 进行中:now - started_at +// - 未开始:0 +func computeDuration(row *adminDAO.MeetingListRow) int { + if row.StartedAt == nil { + return 0 + } + end := time.Now() + if row.EndedAt != nil { + end = *row.EndedAt + } + d := end.Sub(*row.StartedAt) + if d < 0 { + return 0 + } + return int(d.Seconds()) +} + +func rowToDTO(row *adminDAO.MeetingListRow) dto.AdminMeetingDTO { + return dto.AdminMeetingDTO{ + ID: row.ID, + RoomCode: row.RoomCode, + Title: row.Title, + Type: row.Type, + TypeLabel: typeLabel(row.Type), + Status: row.Status, + StatusLabel: statusLabel(row.Status), + HostID: row.HostID, + HostNickname: row.HostNickname, + HostAvatar: row.HostAvatar, + MaxMembers: row.MaxMembers, + ParticipantCount: row.ParticipantCount, + RecordingCount: row.RecordingCount, + DurationSec: computeDuration(row), + ScheduledAt: formatTimePtr(row.ScheduledAt), + StartedAt: formatTimePtr(row.StartedAt), + EndedAt: formatTimePtr(row.EndedAt), + EndedReason: formatStrPtr(row.EndedReason), + CreatedAt: row.CreatedAt.Format(datetimeLayout), + } +} + +func participantToDTO(p *adminDAO.ParticipantDetailRow) dto.AdminMeetingParticipantDTO { + return dto.AdminMeetingParticipantDTO{ + UserID: p.UserID, + Nickname: p.Nickname, + Avatar: p.Avatar, + Role: p.Role, + RoleLabel: roleLabel(p.Role), + JoinedAt: p.JoinedAt.Format(datetimeLayout), + LeftAt: formatTimePtr(p.LeftAt), + LeftReason: formatStrPtr(p.LeftReason), + Duration: p.Duration, + } +} + +func recordingToDTO(r *meetingModel.MeetingRecording) dto.AdminMeetingRecordingDTO { + out := dto.AdminMeetingRecordingDTO{ + ID: r.ID, + StartedBy: r.StartedBy, + Status: r.Status, + FileURL: r.FileURL, + SizeBytes: r.SizeBytes, + DurationSec: r.DurationSec, + StartedAt: r.StartedAt.Format(datetimeLayout), + StoppedAt: formatTimePtr(r.StoppedAt), + } + // 仅 failed 暴露 failure_reason;recording 状态下 FailureReason 字段被复用为本地路径,不应外泄 + if r.Status == meetingModel.MeetingRecordingStatusFailed { + out.FailureReason = r.FailureReason + } + return out +} diff --git a/backend/go-service/app/constants/meeting.go b/backend/go-service/app/constants/meeting.go index 24b03d1..2cc738a 100644 --- a/backend/go-service/app/constants/meeting.go +++ b/backend/go-service/app/constants/meeting.go @@ -129,6 +129,24 @@ const ( // ====== 补充事件(非 §6.3 核心 11 事件,但业务必须)====== MeetingWSEventChatMessage = "meeting.chat.message" // S→C:会议内文字聊天广播(REST SendChat 触发) + + // ====== 屏幕共享(Phase 3 引入,2 个)====== + // 屏幕流仍走普通的 sendTransport video Producer 路径(appData.screen=true 区分), + // 这两个事件用于"会议内仅一份屏幕共享"的状态同步: + // - 启动时广播 producer_id + 共享者,前端将该流提升为大画面 + // - 停止时广播,前端把屏幕 tile 撤掉,回到普通网格布局 + // 后入者由 OnRoomJoin 内部 pushExistingRoomState 定向补推 started 事件,避免错过 + MeetingWSEventScreenStarted = "meeting.screen.started" // S→C:屏幕共享开始 + MeetingWSEventScreenStopped = "meeting.screen.stopped" // S→C:屏幕共享停止 + + // ====== 会议录制(Phase B 引入,2 个,仅 S→C)====== + // 启动 / 停止由 REST 接口(POST /meetings/:code/recordings、DELETE …)触发, + // service 层落 DB + 调 media-server 后通过下面两个事件广播给会议内所有活跃成员: + // - started:让所有客户端在 UI 上显示"正在录制"指示器 + // - stopped:附带最终状态(ready / failed)+ 文件 URL(ready 时), + // 便于前端立即跳转到回放或提示失败 + MeetingWSEventRecordingStarted = "meeting.recording.started" // S→C:录制开始 + MeetingWSEventRecordingStopped = "meeting.recording.stopped" // S→C:录制结束 ) // MeetingWSClientEvents 客户端可主动发起的 WS 事件(C→S)白名单 diff --git a/backend/go-service/app/dto/admin_dto.go b/backend/go-service/app/dto/admin_dto.go index 01d6f39..f842467 100644 --- a/backend/go-service/app/dto/admin_dto.go +++ b/backend/go-service/app/dto/admin_dto.go @@ -140,3 +140,99 @@ type ActiveGroupItem struct { Name string `json:"name"` // 群名称 Count int64 `json:"count"` // 消息数 } + +// ====== 管理端会议管理 DTO(Phase B 扩展) ====== + +// AdminMeetingListRequest 管理端会议列表查询请求 +// +// 与消息列表保持一致的分页/筛选风格: +// - Keyword:模糊匹配 title 或 room_code,方便客服按客户描述定位 +// - Status:0=未开始 / 1=进行中 / 2=已结束(与 meeting_rooms.status 一一对应) +// - HasRecording:true=仅返回有录制记录的会议,便于"录制审计"场景 +// - StartTime / EndTime:按 created_at 过滤,YYYY-MM-DD 闭区间 +type AdminMeetingListRequest struct { + Keyword string `form:"keyword"` // 模糊匹配 title / room_code + Status *int `form:"status"` // 0=未开始 1=进行中 2=已结束 + HostID *int64 `form:"host_id"` // 主持人 ID 精确匹配 + HasRecording *bool `form:"has_recording"` // true=仅含录制的会议 + StartTime string `form:"start_time"` // YYYY-MM-DD + EndTime string `form:"end_time"` // YYYY-MM-DD + Page int `form:"page"` + PageSize int `form:"page_size"` +} + +// AdminMeetingDTO 管理端会议条目(列表行) +// +// 字段选择原则: +// - 列表场景:会议号 / 标题 / 主持人 / 状态 / 时间 / 参与人数 / 录制数 +// - 详情字段(参与者明细、录制列表)放到 GetMeetingDetail 单独返回,避免列表过宽 +type AdminMeetingDTO struct { + ID int64 `json:"id"` + RoomCode string `json:"room_code"` + Title string `json:"title"` + Type int `json:"type"` // 1=即时 2=预约 + TypeLabel string `json:"type_label"` // "即时会议" / "预约会议" + Status int `json:"status"` // 0/1/2 + StatusLabel string `json:"status_label"` // "未开始" / "进行中" / "已结束" + HostID int64 `json:"host_id"` + HostNickname string `json:"host_nickname"` + HostAvatar string `json:"host_avatar"` + MaxMembers int `json:"max_members"` + ParticipantCount int `json:"participant_count"` // 历史累计入会人数(去重) + RecordingCount int `json:"recording_count"` // 该会议的录制条数 + DurationSec int `json:"duration_sec"` // 会议时长(秒);结束前用 now-started_at,结束后用 ended_at-started_at + ScheduledAt string `json:"scheduled_at"` // 预约会议开始时间,空表示即时会议 + StartedAt string `json:"started_at"` // 实际开始 + EndedAt string `json:"ended_at"` // 实际结束 + EndedReason string `json:"ended_reason"` // host_ended / empty_ttl / admin_force / system_error + CreatedAt string `json:"created_at"` +} + +// AdminMeetingListResponse 管理端会议列表响应(含统计摘要) +type AdminMeetingListResponse struct { + Total int64 `json:"total"` + List []AdminMeetingDTO `json:"list"` + Page int `json:"page"` + PageSize int `json:"page_size"` +} + +// AdminMeetingParticipantDTO 会议详情中的参与者条目 +type AdminMeetingParticipantDTO struct { + UserID int64 `json:"user_id"` + Nickname string `json:"nickname"` + Avatar string `json:"avatar"` + Role int `json:"role"` // 0=普通 1=主持人 2=联合主持人 + RoleLabel string `json:"role_label"` + JoinedAt string `json:"joined_at"` + LeftAt string `json:"left_at"` // 空表示仍在会议中 + LeftReason string `json:"left_reason"` // self / kicked / host_end / empty_ttl / disconnect + Duration int `json:"duration"` // 秒 +} + +// AdminMeetingRecordingDTO 会议详情中的录制条目 +type AdminMeetingRecordingDTO struct { + ID int64 `json:"id"` + StartedBy int64 `json:"started_by"` + Status string `json:"status"` // recording / uploading / ready / failed + FileURL string `json:"file_url"` + SizeBytes int64 `json:"size_bytes"` + DurationSec int `json:"duration_sec"` + FailureReason string `json:"failure_reason"` // 仅 failed 时返回 + StartedAt string `json:"started_at"` + StoppedAt string `json:"stopped_at"` +} + +// AdminMeetingDetailResponse 会议详情(含参与者列表 + 录制列表) +type AdminMeetingDetailResponse struct { + AdminMeetingDTO + Participants []AdminMeetingParticipantDTO `json:"participants"` + Recordings []AdminMeetingRecordingDTO `json:"recordings"` +} + +// AdminMeetingStatsResponse 会议管理摘要统计(顶栏卡片) +type AdminMeetingStatsResponse struct { + TotalCount int64 `json:"total_count"` // 历史会议总数 + ActiveCount int64 `json:"active_count"` // 当前进行中 + TodayCount int64 `json:"today_count"` // 今日创建的会议 + RecordingCount int64 `json:"recording_count"` // 录制总数 +} diff --git a/backend/go-service/app/meeting/controller/meeting_recording_controller.go b/backend/go-service/app/meeting/controller/meeting_recording_controller.go new file mode 100644 index 0000000..64acf21 --- /dev/null +++ b/backend/go-service/app/meeting/controller/meeting_recording_controller.go @@ -0,0 +1,214 @@ +package controller + +import ( + "errors" + "os" + "strconv" + + "github.com/echochat/backend/app/meeting/model" + "github.com/echochat/backend/app/meeting/service" + "github.com/echochat/backend/pkg/logs" + "github.com/echochat/backend/pkg/utils" + "github.com/gin-gonic/gin" + "go.uber.org/zap" +) + +// MeetingRecordingController 会议录制相关 REST 接口(Phase B 引入) +// +// 路由(在 router.go 中挂载到 /api/v1/meeting/rooms/:code/recordings 下): +// - POST /rooms/:code/recordings 启动录制(仅 host) +// - DELETE /rooms/:code/recordings/:id 停止录制(仅 host) +// - GET /rooms/:code/recordings 列出录制(host 或参与者) +type MeetingRecordingController struct { + recordingSvc *service.MeetingRecordingService +} + +// NewMeetingRecordingController Wire Provider +func NewMeetingRecordingController(svc *service.MeetingRecordingService) *MeetingRecordingController { + return &MeetingRecordingController{recordingSvc: svc} +} + +// handleErr 将录制相关领域错误映射成 HTTP 状态码 +func (ctl *MeetingRecordingController) handleErr(c *gin.Context, err error, fallbackMsg string) { + switch { + case errors.Is(err, service.ErrMeetingNotFound), + errors.Is(err, service.ErrRecordingNotFound): + utils.ResponseNotFound(c, err.Error()) + case errors.Is(err, service.ErrNotMeetingHost): + utils.ResponseForbidden(c, err.Error()) + case errors.Is(err, service.ErrMediaServiceUnavailable): + utils.ResponseError(c, err.Error()) + case errors.Is(err, service.ErrNotInMeeting), + errors.Is(err, service.ErrMeetingEnded), + errors.Is(err, service.ErrRecordingAlreadyActive), + errors.Is(err, service.ErrRecordingNoProducers), + errors.Is(err, service.ErrRecordingNotStoppable): + utils.ResponseBadRequest(c, err.Error()) + default: + logs.Warn(c.Request.Context(), "controller.meeting_recording_controller.handleErr", + fallbackMsg, zap.Error(err)) + utils.ResponseError(c, fallbackMsg) + } +} + +// StartRecording POST /api/v1/meeting/rooms/:code/recordings +func (ctl *MeetingRecordingController) StartRecording(c *gin.Context) { + userID, ok := requireUserID(c) + if !ok { + return + } + code := c.Param("code") + rec, err := ctl.recordingSvc.StartRecording(c.Request.Context(), userID, code) + if err != nil { + ctl.handleErr(c, err, "启动录制失败") + return + } + utils.ResponseSuccess(c, recordingToDTO(rec)) +} + +// StopRecording DELETE /api/v1/meeting/rooms/:code/recordings/:id +func (ctl *MeetingRecordingController) StopRecording(c *gin.Context) { + userID, ok := requireUserID(c) + if !ok { + return + } + code := c.Param("code") + recordingID, err := strconv.ParseInt(c.Param("id"), 10, 64) + if err != nil || recordingID <= 0 { + utils.ResponseBadRequest(c, "录制 ID 非法") + return + } + rec, stopErr := ctl.recordingSvc.StopRecording(c.Request.Context(), userID, code, recordingID) + if stopErr != nil { + ctl.handleErr(c, stopErr, "停止录制失败") + return + } + utils.ResponseSuccess(c, recordingToDTO(rec)) +} + +// ListRecordings GET /api/v1/meeting/rooms/:code/recordings +func (ctl *MeetingRecordingController) ListRecordings(c *gin.Context) { + userID, ok := requireUserID(c) + if !ok { + return + } + code := c.Param("code") + list, err := ctl.recordingSvc.ListRecordings(c.Request.Context(), userID, code) + if err != nil { + ctl.handleErr(c, err, "查询录制列表失败") + return + } + out := make([]map[string]any, 0, len(list)) + for i := range list { + out = append(out, recordingToDTO(&list[i])) + } + utils.ResponseSuccess(c, gin.H{"list": out}) +} + +// requireUserID 复用 controller 包内的 requireUserID(在 meeting_controller.go 中定义) +// 这里通过同包私有函数访问,避免引入额外依赖 + +// recordingToDTO 录制记录的精简响应结构 +// 不暴露 RemoteID / FileObject(敏感的内部寻址字段),供前端直接渲染列表 + 播放 +func recordingToDTO(r *model.MeetingRecording) map[string]any { + if r == nil { + return nil + } + out := map[string]any{ + "id": r.ID, + "room_id": r.RoomID, + "room_code": r.RoomCode, + "started_by": r.StartedBy, + "status": r.Status, + "file_url": r.FileURL, + "size_bytes": r.SizeBytes, + "duration_sec": r.DurationSec, + "failure_reason": "", + "started_at": r.StartedAt.Unix(), + } + // 仅在最终态 failed 时回显失败原因;recording 状态下 FailureReason 字段被当作临时本地路径,不应外泄 + if r.Status == model.MeetingRecordingStatusFailed { + out["failure_reason"] = r.FailureReason + } + if r.StoppedAt != nil { + out["stopped_at"] = r.StoppedAt.Unix() + } + return out +} + +// ==================== Internal Webhook(Phase B watchdog) ==================== + +// mediaServerFailurePayload media-server 上报的失败 JSON 体 +// 字段与 media-server/src/services/recording.service.ts::RecordingFailureCallback 对应 +type mediaServerFailurePayload struct { + RecordingID string `json:"recordingId"` // media-server 端的 remoteId(UUID 字符串),非 Go 端主键 + RoomCode string `json:"roomCode"` + RouterID string `json:"routerId"` + ExitCode *int `json:"exitCode"` + FailureReason string `json:"failureReason"` +} + +// MediaServerFailureWebhook POST /internal/meeting/recordings/failure +// +// media-server 端 ffmpeg 异常退出时回调此接口(由 RECORDING_FAILURE_WEBHOOK_URL 配置) +// 鉴权: +// - 必须携带 Header `X-Internal-Secret`,值与环境变量 INTERNAL_WEBHOOK_SECRET 一致 +// - 未配置 secret 时直接 403(默认安全,避免暴露内部端点) +// +// 行为: +// - 解析 payload.recordingId(media-server remoteId)→ 反查 Go 端 meeting_recordings 表 +// - 委托给 MeetingRecordingService.HandleMediaServerFailure 落 failed + 广播 +// - 始终返回 200(webhook 语义:吞错避免 media-server 端反复重试) +func (ctl *MeetingRecordingController) MediaServerFailureWebhook(c *gin.Context) { + funcName := "controller.meeting_recording_controller.MediaServerFailureWebhook" + secret := os.Getenv("INTERNAL_WEBHOOK_SECRET") + if secret == "" { + logs.Warn(c.Request.Context(), funcName, "INTERNAL_WEBHOOK_SECRET 未配置,拒绝 webhook 请求") + utils.ResponseForbidden(c, "internal webhook not configured") + return + } + if c.GetHeader("X-Internal-Secret") != secret { + logs.Warn(c.Request.Context(), funcName, "webhook 鉴权失败", + zap.String("client_ip", c.ClientIP())) + utils.ResponseForbidden(c, "invalid internal secret") + return + } + + var payload mediaServerFailurePayload + if err := c.ShouldBindJSON(&payload); err != nil { + logs.Warn(c.Request.Context(), funcName, "解析 payload 失败", zap.Error(err)) + utils.ResponseBadRequest(c, "invalid payload") + return + } + + if payload.RecordingID == "" { + utils.ResponseBadRequest(c, "missing recordingId") + return + } + + // 反查 Go 端主键:通过 RemoteID 定位 meeting_recordings 行 + // 调用方(media-server)持有的是 remoteId(UUID 字符串),Go 端用 ID 主键管理 + rec, err := ctl.recordingSvc.GetByRemoteID(c.Request.Context(), payload.RecordingID) + if err != nil { + logs.Warn(c.Request.Context(), funcName, "查询录制记录失败(按 remote_id)", + zap.String("remote_id", payload.RecordingID), zap.Error(err)) + // 仍 200,避免 media-server 重试风暴 + utils.ResponseSuccess(c, gin.H{"handled": false, "reason": "lookup_failed"}) + return + } + if rec == nil { + logs.Info(c.Request.Context(), funcName, "未找到对应录制记录(可能已被清理)", + zap.String("remote_id", payload.RecordingID)) + utils.ResponseSuccess(c, gin.H{"handled": false, "reason": "not_found"}) + return + } + + if err := ctl.recordingSvc.HandleMediaServerFailure(c.Request.Context(), rec.ID, payload.FailureReason); err != nil { + logs.Warn(c.Request.Context(), funcName, "处理失败 webhook 出错", + zap.Int64("recording_id", rec.ID), zap.Error(err)) + utils.ResponseSuccess(c, gin.H{"handled": false, "reason": "internal_error"}) + return + } + + utils.ResponseSuccess(c, gin.H{"handled": true, "recording_id": rec.ID}) +} diff --git a/backend/go-service/app/meeting/dao/meeting_recording_dao.go b/backend/go-service/app/meeting/dao/meeting_recording_dao.go new file mode 100644 index 0000000..1d7bf74 --- /dev/null +++ b/backend/go-service/app/meeting/dao/meeting_recording_dao.go @@ -0,0 +1,174 @@ +package dao + +import ( + "context" + "errors" + "time" + + "github.com/echochat/backend/app/meeting/model" + "github.com/echochat/backend/pkg/logs" + "go.uber.org/zap" + "gorm.io/gorm" +) + +// MeetingRecordingDAO 会议录制数据访问对象(Phase B 新增) +// +// 设计要点: +// - Create 后立刻获得 ID 主键,作为前端可见的 recordingId +// - 状态流转通过 UpdateStatus / UpdateForUpload / UpdateReady / UpdateFailed 等显式方法做, +// 避免业务层直接 db.Save 大对象造成不一致写 +// - 每个状态切换都基于"前置状态白名单"行级 WHERE,幂等 + 防止 race +type MeetingRecordingDAO struct { + db *gorm.DB +} + +// NewMeetingRecordingDAO 构造 DAO(Wire Provider) +func NewMeetingRecordingDAO(db *gorm.DB) *MeetingRecordingDAO { + return &MeetingRecordingDAO{db: db} +} + +// Create 写入一条新录制记录(recording 初始状态) +// 调用方需保证 RemoteID 已由 media-server 返回;若尚未拿到(先建本地壳记录),传空串即可, +// 后续用 UpdateRemoteID 补全。 +func (d *MeetingRecordingDAO) Create(ctx context.Context, rec *model.MeetingRecording) error { + funcName := "dao.meeting_recording_dao.Create" + if err := d.db.WithContext(ctx).Create(rec).Error; err != nil { + logs.Error(ctx, funcName, "创建录制记录失败", + zap.String("room_code", rec.RoomCode), + zap.Int64("started_by", rec.StartedBy), + zap.Error(err)) + return err + } + return nil +} + +// GetByID 按主键查询,记录不存在返回 (nil, nil) +func (d *MeetingRecordingDAO) GetByID(ctx context.Context, id int64) (*model.MeetingRecording, error) { + var rec model.MeetingRecording + err := d.db.WithContext(ctx).First(&rec, id).Error + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + return nil, err + } + return &rec, nil +} + +// GetActiveByRoom 查询某会议正在进行中的录制(status=recording 或 uploading) +// 同会议同一时刻只允许 1 条活跃录制;用于"开始录制"前置防重 + EndRoom 自动停录 +// 不存在返回 (nil, nil) +func (d *MeetingRecordingDAO) GetActiveByRoom(ctx context.Context, roomID int64) (*model.MeetingRecording, error) { + var rec model.MeetingRecording + err := d.db.WithContext(ctx). + Where("room_id = ? AND status IN ?", roomID, []string{ + model.MeetingRecordingStatusRecording, + model.MeetingRecordingStatusUploading, + }). + Order("id DESC"). + First(&rec).Error + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + return nil, err + } + return &rec, nil +} + +// ListByRoom 按 room_id 倒序列出全部录制(用于"录制列表"页) +// 不分页:单会议录制条数上限受业务约束(一般 ≤ 20);如需分页可后续扩展 +func (d *MeetingRecordingDAO) ListByRoom(ctx context.Context, roomID int64) ([]model.MeetingRecording, error) { + var list []model.MeetingRecording + err := d.db.WithContext(ctx). + Where("room_id = ?", roomID). + Order("id DESC"). + Find(&list).Error + return list, err +} + +// GetByRemoteID 按 media-server 返回的 remoteId 反查本地记录,不存在返回 (nil, nil) +// 供 webhook sink 在收到 media-server 的失败回调时定位本地行 +func (d *MeetingRecordingDAO) GetByRemoteID(ctx context.Context, remoteID string) (*model.MeetingRecording, error) { + if remoteID == "" { + return nil, nil + } + var rec model.MeetingRecording + err := d.db.WithContext(ctx).Where("remote_id = ?", remoteID).First(&rec).Error + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + return nil, err + } + return &rec, nil +} + +// UpdateRemoteID 写入 media-server 返回的 recordingId +// 用于先 Create 本地记录、再调 media-server 拿到 RemoteID 后回填的两阶段流程 +func (d *MeetingRecordingDAO) UpdateRemoteID(ctx context.Context, id int64, remoteID string) error { + return d.db.WithContext(ctx). + Model(&model.MeetingRecording{}). + Where("id = ?", id). + Update("remote_id", remoteID).Error +} + +// MarkUploading 录制结束 → 上传中:仅当 status=recording 才推进,幂等 +// 同时落 stopped_at + duration_sec +func (d *MeetingRecordingDAO) MarkUploading(ctx context.Context, id int64, stoppedAt time.Time, durationSec int) (int64, error) { + res := d.db.WithContext(ctx). + Model(&model.MeetingRecording{}). + Where("id = ? AND status = ?", id, model.MeetingRecordingStatusRecording). + Updates(map[string]interface{}{ + "status": model.MeetingRecordingStatusUploading, + "stopped_at": stoppedAt, + "duration_sec": durationSec, + }) + return res.RowsAffected, res.Error +} + +// MarkReady 上传成功 → ready;记录 file_url / file_object / size_bytes +// 仅当 status=uploading 才推进 +func (d *MeetingRecordingDAO) MarkReady(ctx context.Context, id int64, fileURL, fileObject string, sizeBytes int64) (int64, error) { + res := d.db.WithContext(ctx). + Model(&model.MeetingRecording{}). + Where("id = ? AND status = ?", id, model.MeetingRecordingStatusUploading). + Updates(map[string]interface{}{ + "status": model.MeetingRecordingStatusReady, + "file_url": fileURL, + "file_object": fileObject, + "size_bytes": sizeBytes, + }) + return res.RowsAffected, res.Error +} + +// MarkFailed 终态失败:从任意非 ready 状态都允许推进,记录 failure_reason +// 调用方需在 reason 为空时传一个有意义的兜底字符串(如 "unknown") +func (d *MeetingRecordingDAO) MarkFailed(ctx context.Context, id int64, reason string) (int64, error) { + if reason == "" { + reason = "unknown" + } + res := d.db.WithContext(ctx). + Model(&model.MeetingRecording{}). + Where("id = ? AND status != ?", id, model.MeetingRecordingStatusReady). + Updates(map[string]interface{}{ + "status": model.MeetingRecordingStatusFailed, + "failure_reason": reason, + }) + return res.RowsAffected, res.Error +} + +// ListUnfinishedSince 查询启动时间早于 cutoff 但仍未到终态(ready/failed)的录制 ID +// 供兜底定时任务在 media-server 重启后清理"卡死的录制",避免 DB 永远停留在 recording 状态 +func (d *MeetingRecordingDAO) ListUnfinishedSince(ctx context.Context, cutoff time.Time, limit int) ([]int64, error) { + var ids []int64 + err := d.db.WithContext(ctx). + Model(&model.MeetingRecording{}). + Where("status IN ? AND started_at < ?", + []string{model.MeetingRecordingStatusRecording, model.MeetingRecordingStatusUploading}, + cutoff). + Order("started_at ASC"). + Limit(limit). + Pluck("id", &ids).Error + return ids, err +} diff --git a/backend/go-service/app/meeting/model/meeting_recording.go b/backend/go-service/app/meeting/model/meeting_recording.go new file mode 100644 index 0000000..a404540 --- /dev/null +++ b/backend/go-service/app/meeting/model/meeting_recording.go @@ -0,0 +1,44 @@ +package model + +import "time" + +// MeetingRecording 会议录制记录,对应 meeting_recordings 表 +// +// 生命周期状态: +// - recording → 录制进行中(media-server ffmpeg 进程在跑,尚未上传) +// - uploading → media-server 已停止录制,Go 后端正在拉取本地 mp4 上传到 MinIO +// - ready → MinIO 上传完成,FileURL 可用 +// - failed → 录制 / 上传失败,FailureReason 记录原因 +// +// RemoteID 是 media-server 内部的 recordingId(hex32),Go 侧 stop / 状态反查时使用。 +// 唯一索引落在 RemoteID 上,保证一对一关系。 +type MeetingRecording struct { + ID int64 `json:"id" gorm:"primaryKey;autoIncrement"` + RoomID int64 `json:"room_id" gorm:"not null;index:idx_meeting_recordings_room"` + RoomCode string `json:"room_code" gorm:"size:32;not null;index:idx_meeting_recordings_room_code"` + StartedBy int64 `json:"started_by" gorm:"not null;index:idx_meeting_recordings_started_by"` + Status string `json:"status" gorm:"size:16;not null;default:recording"` // recording / uploading / ready / failed + RemoteID string `json:"remote_id" gorm:"size:64;uniqueIndex"` // media-server 侧 recordingId + FileURL string `json:"file_url" gorm:"size:512"` + FileObject string `json:"file_object" gorm:"size:255"` // MinIO 内的对象 Key,便于鉴权下载/删除 + SizeBytes int64 `json:"size_bytes" gorm:"default:0"` + DurationSec int `json:"duration_sec" gorm:"default:0"` + FailureReason string `json:"failure_reason" gorm:"size:255"` + StartedAt time.Time `json:"started_at" gorm:"not null;type:timestamp(0)"` + StoppedAt *time.Time `json:"stopped_at" gorm:"type:timestamp(0)"` + CreatedAt time.Time `json:"created_at" gorm:"not null;autoCreateTime;type:timestamp(0)"` + UpdatedAt time.Time `json:"updated_at" gorm:"not null;autoUpdateTime;type:timestamp(0)"` +} + +// TableName 指定数据库表名 +func (MeetingRecording) TableName() string { + return "meeting_recordings" +} + +// 录制状态常量;与表中 status 字段一一对应 +const ( + MeetingRecordingStatusRecording = "recording" + MeetingRecordingStatusUploading = "uploading" + MeetingRecordingStatusReady = "ready" + MeetingRecordingStatusFailed = "failed" +) diff --git a/backend/go-service/app/meeting/provider.go b/backend/go-service/app/meeting/provider.go index cbb21e7..4677212 100644 --- a/backend/go-service/app/meeting/provider.go +++ b/backend/go-service/app/meeting/provider.go @@ -20,11 +20,14 @@ var MeetingSet = wire.NewSet( dao.NewMeetingRoomDAO, dao.NewMeetingParticipantDAO, dao.NewMeetingChatDAO, + dao.NewMeetingRecordingDAO, // Phase B 新增:录制记录 DAO service.NewMeetingBroadcaster, service.NewMeetingLifecycleService, // Task 8 新增:会议生命周期状态机 service.NewMeetingService, service.NewMeetingSignalService, + service.NewMeetingRecordingService, // Phase B 新增:录制业务服务 controller.NewMeetingController, + controller.NewMeetingRecordingController, // Phase B 新增:录制 REST 控制器 controller.NewMeetingWSHandler, task.NewMeetingCleanupTask, // Task 8 新增:生命周期兜底定时任务 diff --git a/backend/go-service/app/meeting/router.go b/backend/go-service/app/meeting/router.go index ca60671..764af07 100644 --- a/backend/go-service/app/meeting/router.go +++ b/backend/go-service/app/meeting/router.go @@ -9,7 +9,13 @@ import ( // RegisterRoutes 注册 meeting 模块的全部前台路由 // 所有接口统一挂载在 /api/v1/meeting/* 前缀下,均需 JWT 认证 // 设计文档 §6.2 的 12 个接口 Task 4 骨架阶段仅搭路由层,业务实现留待 Task 5/6/7 -func RegisterRoutes(r *gin.Engine, ctrl *controller.MeetingController, jwtAuth gin.HandlerFunc) { +// Phase B 新增录制三接口(POST/DELETE/GET /rooms/:code/recordings*) +func RegisterRoutes( + r *gin.Engine, + ctrl *controller.MeetingController, + recordingCtrl *controller.MeetingRecordingController, + jwtAuth gin.HandlerFunc, +) { authed := r.Group("/api/v1/meeting") authed.Use(jwtAuth) { @@ -29,5 +35,15 @@ func RegisterRoutes(r *gin.Engine, ctrl *controller.MeetingController, jwtAuth g authed.POST("/rooms/:code/chats", ctrl.SendChat) authed.GET("/rooms/:code/chats", ctrl.ListChats) + + // 会议录制(Phase B 引入) + authed.POST("/rooms/:code/recordings", recordingCtrl.StartRecording) + authed.DELETE("/rooms/:code/recordings/:id", recordingCtrl.StopRecording) + authed.GET("/rooms/:code/recordings", recordingCtrl.ListRecordings) } + + // 内部 webhook(Phase B watchdog):media-server 异常退出回调 + // 鉴权由 controller 内部用 X-Internal-Secret header + INTERNAL_WEBHOOK_SECRET 校验 + // 不走 JWT,因为调用方是 media-server 而非用户 + r.POST("/internal/meeting/recordings/failure", recordingCtrl.MediaServerFailureWebhook) } diff --git a/backend/go-service/app/meeting/service/http_media_orchestrator.go b/backend/go-service/app/meeting/service/http_media_orchestrator.go index 18055c1..d0421a9 100644 --- a/backend/go-service/app/meeting/service/http_media_orchestrator.go +++ b/backend/go-service/app/meeting/service/http_media_orchestrator.go @@ -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/recordings(Phase B 引入) +// +// 注意:录制启动会触发 ffmpeg spawn + 端口分配 + 250ms 等待启动,正常耗时约 500ms~1.5s +// 因此走 TimeoutMS(默认 10s)而非 CloseTimeoutMS。失败不重试,避免重复 spawn ffmpeg +// +// 错误映射: +// - Node 404(router 不存在)→ ErrMediaResourceNotFound +// - Node 409(producer 不存在/不可消费)→ 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/:id(Phase B 引入) +// +// 时间预算:ffmpeg 优雅退出最多 STOP_GRACE_MS=5s + media-server 自身处理 ~200ms +// 因此 timeout 需要明显 > 5s;用 TimeoutMS(默认 10s)足够 +// +// 404(recording 不存在 / 已停止)映射为 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 diff --git a/backend/go-service/app/meeting/service/interfaces.go b/backend/go-service/app/meeting/service/interfaces.go index d20a18a..76ae876 100644 --- a/backend/go-service/app/meeting/service/interfaces.go +++ b/backend/go-service/app/meeting/service/interfaces.go @@ -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"` } +// ====== 会议录制 RPC(Phase 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 + // 成功返回 recordingId(hex32)+ 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 +} diff --git a/backend/go-service/app/meeting/service/meeting_recording_service.go b/backend/go-service/app/meeting/service/meeting_recording_service.go new file mode 100644 index 0000000..929030e --- /dev/null +++ b/backend/go-service/app/meeting/service/meeting_recording_service.go @@ -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(与图片/语音同 bucket,public-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 字段暂存 outputPath,stop 阶段读出来上传 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 MarkUploading(status: recording → uploading),同时清空临时 FailureReason +// 4. 上传本地 mp4 到 MinIO(recordings/{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 stop(404 视为已停止 → 用本地 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 部分 +// +// 选取规则(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, + }) +} diff --git a/backend/go-service/app/meeting/service/meeting_service.go b/backend/go-service/app/meeting/service/meeting_service.go index 13366de..3f3815c 100644 --- a/backend/go-service/app/meeting/service/meeting_service.go +++ b/backend/go-service/app/meeting/service/meeting_service.go @@ -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)) } diff --git a/backend/go-service/app/meeting/service/meeting_signal_service.go b/backend/go-service/app/meeting/service/meeting_signal_service.go index bac7639..2659a39 100644 --- a/backend/go-service/app/meeting/service/meeting_signal_service.go +++ b/backend/go-service/app/meeting/service/meeting_signal_service.go @@ -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 必然等于 userID(trackResource 已校验归属) + 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() diff --git a/backend/go-service/app/provider/provider.go b/backend/go-service/app/provider/provider.go index 15edea3..d39629c 100644 --- a/backend/go-service/app/provider/provider.go +++ b/backend/go-service/app/provider/provider.go @@ -42,6 +42,7 @@ type App struct { ContactManageController *adminController.ContactManageController // 管理端好友关系管理控制器 GroupManageController *adminController.GroupManageController // 管理端群组管理控制器 MessageManageController *adminController.MessageManageController // 管理端消息管理控制器 + MeetingManageController *adminController.MeetingManageController // 管理端会议管理控制器(Phase B) WSHandler *wsApp.Handler // WebSocket 连接处理器 Hub *ws.Hub // WebSocket Hub 连接管理 PubSub *ws.PubSub // Redis Pub/Sub 消息路由 @@ -59,6 +60,8 @@ type App struct { MeetingSignalService *meetingService.MeetingSignalService // 会议 WS 信令业务服务(Task 6 落地) MeetingLifecycleSvc *meetingService.MeetingLifecycleService // 会议生命周期状态机(Task 8 落地) MeetingController *meetingController.MeetingController // 会议 REST 控制器 + MeetingRecordingController *meetingController.MeetingRecordingController // 会议录制 REST 控制器(Phase B 新增) + MeetingRecordingService *meetingService.MeetingRecordingService // 会议录制业务服务(Phase B 新增) MeetingWSHandler *meetingController.MeetingWSHandler // 会议 WS 事件 Handler(构造时自动注册路由到 Hub) MeetingCleanupTask *meetingTask.MeetingCleanupTask // 会议生命周期兜底定时任务(Task 8 落地) } @@ -77,6 +80,7 @@ func NewApp( contactManageCtrl *adminController.ContactManageController, groupManageCtrl *adminController.GroupManageController, msgManageCtrl *adminController.MessageManageController, + meetingManageCtrl *adminController.MeetingManageController, wsHandler *wsApp.Handler, hub *ws.Hub, pubsub *ws.PubSub, @@ -94,12 +98,16 @@ func NewApp( meetingSignalSvc *meetingService.MeetingSignalService, meetingLifecycleSvc *meetingService.MeetingLifecycleService, meetingCtrl *meetingController.MeetingController, + meetingRecordingCtrl *meetingController.MeetingRecordingController, + meetingRecordingSvc *meetingService.MeetingRecordingService, meetingWSHandler *meetingController.MeetingWSHandler, meetingCleanup *meetingTask.MeetingCleanupTask, ) *App { wsHandler.SetOfflinePusher(offlinePusher) wsHandler.SetNotifyConnectHook(notifySvc) wsHandler.SetMeetingDisconnectHook(meetingSignalSvc) + // Phase B:通过 setter 注入录制停止钩子,避免 MeetingService ↔ MeetingRecordingService 构造期循环依赖 + meetingSvc.SetRecordingStopper(meetingRecordingSvc) return &App{ Config: cfg, @@ -114,6 +122,7 @@ func NewApp( ContactManageController: contactManageCtrl, GroupManageController: groupManageCtrl, MessageManageController: msgManageCtrl, + MeetingManageController: meetingManageCtrl, WSHandler: wsHandler, Hub: hub, PubSub: pubsub, @@ -130,9 +139,11 @@ func NewApp( MeetingService: meetingSvc, MeetingSignalService: meetingSignalSvc, MeetingLifecycleSvc: meetingLifecycleSvc, - MeetingController: meetingCtrl, - MeetingWSHandler: meetingWSHandler, - MeetingCleanupTask: meetingCleanup, + MeetingController: meetingCtrl, + MeetingRecordingController: meetingRecordingCtrl, + MeetingRecordingService: meetingRecordingSvc, + MeetingWSHandler: meetingWSHandler, + MeetingCleanupTask: meetingCleanup, } } diff --git a/backend/go-service/app/provider/wire_gen.go b/backend/go-service/app/provider/wire_gen.go index b76a0a3..2be9e33 100644 --- a/backend/go-service/app/provider/wire_gen.go +++ b/backend/go-service/app/provider/wire_gen.go @@ -80,6 +80,9 @@ func InitializeApp(cfg *config.Config) (*App, error) { conversationDAO := im.ProvideConversationDAO(gormDB) messageManageService := service2.NewMessageManageService(messageManageDAO, userDAO, conversationDAO, pubSub) messageManageController := controller2.NewMessageManageController(messageManageService) + meetingManageDAO := dao2.NewMeetingManageDAO(gormDB) + meetingManageService := service2.NewMeetingManageService(meetingManageDAO) + meetingManageController := controller2.NewMeetingManageController(meetingManageService) serverConfig := provideServerConfig(cfg) handler := ws.ProvideWSHandler(hub, pubSub, jwtConfig, serverConfig, onlineService, authService) friendGroupDAO := dao3.NewFriendGroupDAO(gormDB) @@ -103,14 +106,17 @@ func InitializeApp(cfg *config.Config) (*App, error) { meetingRoomDAO := dao6.NewMeetingRoomDAO(gormDB) meetingParticipantDAO := dao6.NewMeetingParticipantDAO(gormDB) meetingChatDAO := dao6.NewMeetingChatDAO(gormDB) + meetingRecordingDAO := dao6.NewMeetingRecordingDAO(gormDB) meetingBroadcaster := service7.NewMeetingBroadcaster(meetingParticipantDAO, pubSub) httpMediaOrchestrator := service7.NewHTTPMediaOrchestrator(cfg) meetingLifecycleService := service7.NewMeetingLifecycleService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, httpMediaOrchestrator, cfg) meetingService := service7.NewMeetingService(meetingRoomDAO, meetingParticipantDAO, meetingChatDAO, gormDB, client, meetingBroadcaster, notifyService, friendshipDAO, onlineService, httpMediaOrchestrator, meetingLifecycleService) meetingSignalService := service7.NewMeetingSignalService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, httpMediaOrchestrator, meetingLifecycleService) + meetingRecordingService := service7.NewMeetingRecordingService(meetingRoomDAO, meetingParticipantDAO, meetingRecordingDAO, client, meetingBroadcaster, httpMediaOrchestrator, minioClient, minioConfig) meetingController := controller7.NewMeetingController(meetingService) + meetingRecordingController := controller7.NewMeetingRecordingController(meetingRecordingService) meetingWSHandler := controller7.NewMeetingWSHandler(meetingSignalService, hub) meetingCleanupTask := task2.NewMeetingCleanupTask(meetingLifecycleService, meetingRoomDAO, meetingChatDAO) - app := NewApp(cfg, gormDB, client, minioClient, authService, authController, adminAuthController, userManageController, onlineController, contactManageController, groupManageController, messageManageController, handler, hub, pubSub, onlineService, contactController, imController, eventHandler, offlinePusher, fileController, groupController, notifyService, notificationController, cleanupTask, meetingService, meetingSignalService, meetingLifecycleService, meetingController, meetingWSHandler, meetingCleanupTask) + app := NewApp(cfg, gormDB, client, minioClient, authService, authController, adminAuthController, userManageController, onlineController, contactManageController, groupManageController, messageManageController, meetingManageController, handler, hub, pubSub, onlineService, contactController, imController, eventHandler, offlinePusher, fileController, groupController, notifyService, notificationController, cleanupTask, meetingService, meetingSignalService, meetingLifecycleService, meetingController, meetingRecordingController, meetingRecordingService, meetingWSHandler, meetingCleanupTask) return app, nil } diff --git a/backend/go-service/cmd/server/main.go b/backend/go-service/cmd/server/main.go index e8dcb28..a24218d 100644 --- a/backend/go-service/cmd/server/main.go +++ b/backend/go-service/cmd/server/main.go @@ -11,6 +11,7 @@ import ( "time" "github.com/echochat/backend/app/im/model" + meetingModel "github.com/echochat/backend/app/meeting/model" "github.com/echochat/backend/app/provider" "github.com/echochat/backend/config" "github.com/echochat/backend/pkg/logs" @@ -56,6 +57,8 @@ func main() { &model.Conversation{}, &model.ConversationMember{}, &model.Message{}, + // Phase B 新增:会议录制元数据表 + &meetingModel.MeetingRecording{}, ); err != nil { logs.Fatal(ctx, "main", "IM 表迁移失败", zap.Error(err)) } diff --git a/backend/go-service/router/router.go b/backend/go-service/router/router.go index 8eec19b..f2931f7 100644 --- a/backend/go-service/router/router.go +++ b/backend/go-service/router/router.go @@ -38,12 +38,12 @@ func Setup(engine *gin.Engine, app *provider.App) { // --- 各模块路由注册 --- auth.RegisterRoutes(engine, app.AuthController, app.AdminAuthController, jwtAuth) - admin.RegisterRoutes(engine, app.UserManageController, app.OnlineController, app.ContactManageController, app.GroupManageController, app.MessageManageController, jwtAuth) + admin.RegisterRoutes(engine, app.UserManageController, app.OnlineController, app.ContactManageController, app.GroupManageController, app.MessageManageController, app.MeetingManageController, jwtAuth) wsApp.RegisterRoutes(engine, app.WSHandler) contact.RegisterRoutes(engine, app.ContactController, jwtAuth) imApp.RegisterRoutes(engine, app.IMController, jwtAuth) fileApp.RegisterRoutes(engine, app.FileController, jwtAuth) groupApp.RegisterRoutes(engine, app.GroupController, jwtAuth) notifyApp.RegisterRoutes(engine, app.NotifyController, jwtAuth) - meetingApp.RegisterRoutes(engine, app.MeetingController, jwtAuth) + meetingApp.RegisterRoutes(engine, app.MeetingController, app.MeetingRecordingController, jwtAuth) } diff --git a/docs/plans/2026-05-18-recording-progress.md b/docs/plans/2026-05-18-recording-progress.md new file mode 100644 index 0000000..d79d0a0 --- /dev/null +++ b/docs/plans/2026-05-18-recording-progress.md @@ -0,0 +1,107 @@ +# 会议录制实现进度(Phase B1→B4) + +启动日期:2026-05-18 +负责人:Cascade pair-programming session + +## 全局架构 + +``` +[host 工具栏] → POST /meeting/rooms/{code}/recording/toggle + ↓ +[Go MeetingRecordingService] + ├─ 校验 host 权限 + 会议状态 + ├─ 落库 meeting_recordings(status=recording) + └─ 调 HTTPMediaOrchestrator.StartRecording(roomCode, producerIDs) + ↓ +[media-server /internal/v1/recordings] + ├─ 为每个 Producer 建 PlainTransport(rtcpMux:false, comedia:false) + ├─ consume 该 Producer 拿到 negotiated rtpParameters + ├─ 选未占用 UDP 端口对(rtp+rtcp)给 ffmpeg + ├─ 写 SDP 文件(含 audio/video pt、codec、rtpmap、fmtp) + ├─ spawn ffmpeg: -protocol_whitelist file,udp,rtp -i input.sdp -c copy output.mp4 + └─ transport.connect({ ip, port: ffmpegPort, rtcpPort: ffmpegRtcpPort }) + ↓ +[ffmpeg 写盘 → /tmp/recordings/{recordingId}.mp4] + ↓ +[host 再次 toggle → 停止] + ├─ media-server SIGINT ffmpeg → wait exit → close transports/consumers + ├─ POST /api/v1/meeting/_internal/recordings/{id}/finalize(callback) + └─ Go 流式拉取 mp4 → MinIO → 落库 status=ready + file_url + size + duration + ↓ +[WS 广播 meeting.recording.started/stopped 给所有成员] +``` + +## 阶段拆分 + +### Phase B1 — 媒体层录制管线 + Go 控制端 ✏️ 进行中 +**media-server**: +- [ ] `services/recording.service.ts`:起停录制、SDP 生成、ffmpeg 进程管理、端口池 +- [ ] `routes/recording.route.ts`:`POST /internal/v1/recordings`、`DELETE /internal/v1/recordings/:id`、`GET /internal/v1/recordings/:id` +- [ ] `schemas/recording.schema.ts`:zod 校验 +- [ ] `app.ts`:注册路由 +- [ ] 端口池实现 `utils/port-pool.ts` + +**Go 后端**: +- [ ] `app/meeting/model/meeting_recording.go`:表 ORM +- [ ] `app/meeting/dao/meeting_recording_dao.go` +- [ ] `app/meeting/service/interfaces.go` 扩展 MediaOrchestrator: `StartRecording / StopRecording` +- [ ] `app/meeting/service/http_media_orchestrator.go` 实现 RPC +- [ ] `cmd/server/main.go` AutoMigrate 加入 MeetingRecording + +### Phase B2 — 业务层 + 持久化 +- [ ] `service/meeting_recording_service.go`:StartRecording/StopRecording/ListRecordings/GetRecording +- [ ] `controller/meeting_controller.go` 新增 4 个端点 +- [ ] `router.go` 挂载路由 +- [ ] `provider.go` wire 注册 +- [ ] 内部回调端点 `/_internal/recordings/{id}/finalize`:media-server 录完后回调,Go 拉文件→ MinIO → 更新 DB +- [ ] MinIO 流式上传 + +### Phase B3 — WS 广播 + 自动停录 +- [ ] `constants/meeting.go`:MeetingWSEventRecordingStarted/Stopped +- [ ] `meeting_signal_service.go` / `meeting_broadcaster.go` 广播 +- [ ] EndRoom / HandleEmptyRoomExpired 自动 StopRecording +- [ ] 前端 `constants/meeting.js` + store 监听 + +### Phase B4 — 前端 UI +- [ ] `MeetingToolbar.vue`:录制按钮(仅 host 可见)+ 红点动画 +- [ ] `room.vue`:顶部"REC"指示条(任意成员都可见) +- [ ] 新页面 `pages/meeting/recordings.vue`:列表 + 下载 + +## 已知风险与决策 + +1. **ffmpeg 必须可用**:媒体服务的运行环境(开发:本地 PATH;生产:Docker 镜像 apt install ffmpeg)。Dockerfile 暂不改,由你部署时确认。 +2. **PlainTransport 端口分配**:用 dgram 试探可用 UDP 端口;初始范围 50000-59999,避开 mediasoup `40000-40199`。 +3. **ffmpeg 启停时序**:录制启动顺序 = 起 ffmpeg(listen) → 100ms 延迟 → mediasoup transport.connect → consume;停止顺序 = SIGINT ffmpeg → wait exit → close transports/consumers。 +4. **录制视频编码**:用 `-c:v copy -c:a copy` 直接复制 RTP 载荷为 mp4(不重编码),节省 CPU。注意:浏览器端 VP8/VP9 muxed 进 mp4 可能不被普通播放器识别,必要时改 `-c:v libx264 -preset ultrafast`。本期先用 copy。 +5. **DB 表 schema**:参考下方 §SQL。 +6. **Phase B 整体不强制阻塞**:每完成一个阶段都可以暂停验证。 + +## SQL + +```sql +CREATE TABLE meeting_recordings ( + id BIGINT UNSIGNED PRIMARY KEY AUTO_INCREMENT, + room_id BIGINT UNSIGNED NOT NULL, + room_code VARCHAR(32) NOT NULL, + started_by BIGINT UNSIGNED NOT NULL, + status VARCHAR(16) NOT NULL DEFAULT 'recording', -- recording/uploading/ready/failed + remote_id VARCHAR(64), -- media-server 侧的 recordingId + file_url VARCHAR(512), + file_object VARCHAR(255), + size_bytes BIGINT UNSIGNED DEFAULT 0, + duration_sec INT UNSIGNED DEFAULT 0, + failure_reason VARCHAR(255), + started_at DATETIME NOT NULL, + stopped_at DATETIME, + created_at DATETIME NOT NULL, + updated_at DATETIME NOT NULL, + KEY idx_room (room_id), + KEY idx_room_code (room_code), + KEY idx_started_by (started_by) +); +``` + +## 当前会话进度 + +- ✅ 调研完成:MinIO 已就绪、media-server 架构吃透、HTTPMediaOrchestrator 模式吃透 +- ⏩ 下一步:写 media-server 端口池 + recording.service.ts diff --git a/frontend/src/api/meeting.js b/frontend/src/api/meeting.js index 5c377fb..fe36efe 100644 --- a/frontend/src/api/meeting.js +++ b/frontend/src/api/meeting.js @@ -9,7 +9,7 @@ * - 响应中的 DTO 字段命名与后端 DTO 完全一致(下划线),前端按需再转换 */ -import { get, post } from '@/utils/request' +import { get, post, del } from '@/utils/request' /** * 统一拆包:utils/request.js 返回完整 envelope { code, message, data, trace_id, time }, @@ -145,6 +145,36 @@ const listChats = (roomCode, params = {}) => { return unwrap(get(`/api/v1/meeting/rooms/${roomCode}/chats`, params)) } +// ==================== 会议录制(Phase B 新增,3 接口) ==================== + +/** + * 启动录制(仅 host) + * @param {string} roomCode + * @returns {Promise} + */ +const startRecording = (roomCode) => { + return unwrap(post(`/api/v1/meeting/rooms/${roomCode}/recordings`)) +} + +/** + * 停止录制(仅 host) + * @param {string} roomCode + * @param {number|string} recordingID + * @returns {Promise} + */ +const stopRecording = (roomCode, recordingID) => { + return unwrap(del(`/api/v1/meeting/rooms/${roomCode}/recordings/${recordingID}`)) +} + +/** + * 查询录制列表(host 或参与者均可) + * @param {string} roomCode + * @returns {Promise<{ list: MeetingRecordingDTO[] }>} + */ +const listRecordings = (roomCode) => { + return unwrap(get(`/api/v1/meeting/rooms/${roomCode}/recordings`)) +} + export default { createRoom, getRoom, @@ -157,5 +187,8 @@ export default { inviteUsers, redeemInvite, sendChat, - listChats + listChats, + startRecording, + stopRecording, + listRecordings } diff --git a/frontend/src/components/meeting/MeetingToolbar.vue b/frontend/src/components/meeting/MeetingToolbar.vue index 9228d0d..b3996b1 100644 --- a/frontend/src/components/meeting/MeetingToolbar.vue +++ b/frontend/src/components/meeting/MeetingToolbar.vue @@ -45,6 +45,54 @@ {{ videoEnabled ? '停止视频' : '开启视频' }} + + + + +