From 096563d3ea9956ff1a176769eb910627066ffb33 Mon Sep 17 00:00:00 2001 From: bujinyuan Date: Tue, 21 Apr 2026 17:10:32 +0800 Subject: [PATCH] =?UTF-8?q?feat(phase2e-2):=20WebSocket=20=E4=BF=A1?= =?UTF-8?q?=E4=BB=A4=E5=8D=8F=E8=AE=AE=2013=20=E4=BA=8B=E4=BB=B6=E5=85=A8?= =?UTF-8?q?=E9=87=8F=E8=90=BD=E5=9C=B0=EF=BC=88Task=206=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 将 meeting.* 事件族从 Task 5 的 PublishToUser 循环升级为完整 WS 信令协议: 抽离统一广播层、扩容 MediaOrchestrator 接口、实现 8 个 C→S 事件业务逻辑 + Redis 资源追踪 + host 权限校验,端到端 18/18 PASS。 核心产出: - 新建 MeetingBroadcaster(统一广播层,REST/WS 共用) - 新建 MeetingSignalService 8 C→S 事件 + cleanupUserResources - 新建 MeetingWSHandler 薄层 controller - MediaOrchestrator 扩容 9 方法 + NoopMediaOrchestrator 占位(Task 7 替换) - C→S 白名单机制防恶意伪造广播事件 - Redis Set 资源追踪防 mediasoup 端资源泄漏 文档同步: - docs/api/frontend/meeting.md 追加 §WebSocket 信令协议(Task 6)200 行 - docs/progress/CURRENT_STATUS.md + project-context.mdc + 实施计划 Task 6 ✅ Made-with: Cursor --- .cursor/rules/project-context.mdc | 3 +- backend/go-service/app/constants/meeting.go | 50 ++- .../meeting/controller/meeting_ws_handler.go | 196 +++++++++ backend/go-service/app/meeting/provider.go | 8 +- .../app/meeting/service/interfaces.go | 129 +++++- .../meeting/service/meeting_broadcaster.go | 91 ++++ .../app/meeting/service/meeting_service.go | 59 +-- .../meeting/service/meeting_signal_service.go | 393 ++++++++++++++++++ backend/go-service/app/provider/provider.go | 10 +- backend/go-service/app/provider/wire_gen.go | 7 +- docs/api/frontend/meeting.md | 273 +++++++++++- ...026-04-21-phase2e-2-implementation.plan.md | 39 +- docs/progress/CURRENT_STATUS.md | 61 ++- 13 files changed, 1212 insertions(+), 107 deletions(-) create mode 100644 backend/go-service/app/meeting/controller/meeting_ws_handler.go create mode 100644 backend/go-service/app/meeting/service/meeting_broadcaster.go create mode 100644 backend/go-service/app/meeting/service/meeting_signal_service.go diff --git a/.cursor/rules/project-context.mdc b/.cursor/rules/project-context.mdc index f218561..40116d8 100644 --- a/.cursor/rules/project-context.mdc +++ b/.cursor/rules/project-context.mdc @@ -27,7 +27,7 @@ alwaysApply: true * 单端 WS 连接架构(沿用),不做多端已读同步(设计文档 §3.1/§3.5/§九 已修订,多端改造推迟到 Phase 2f/二期) * 专用设计:`docs/plans/2026-04-20-phase2e-1-design.md`;实施计划:`docs/plans/2026-04-20-phase2e-1-implementation.plan.md`;验证报告:`test-report-phase2e-1-notification.md` * API 文档:`docs/api/frontend/notify.md` - - 2e-2 会议 MVP(约 17 天)🚧 **代码开发中**(Task 0-5 ✅ / Task 6-16 待执行):mediasoup Node.js 独立 `media-server/` + 即时会议(≤8 人)+ 密码/邀请链接/通知邀请三合一 + 设备预览页 + 主持人四件套 + 会议内聊天 + 双态部署(本机 + 公网 coturn)+ 响应式(桌面/手机) + - 2e-2 会议 MVP(约 17 天)🚧 **代码开发中**(Task 0-6 ✅ / Task 7-16 待执行):mediasoup Node.js 独立 `media-server/` + 即时会议(≤8 人)+ 密码/邀请链接/通知邀请三合一 + 设备预览页 + 主持人四件套 + 会议内聊天 + 双态部署(本机 + 公网 coturn)+ 响应式(桌面/手机) * 专用设计:`docs/plans/2026-04-21-phase2e-2-design.md`(16 章节);实施计划:`docs/plans/2026-04-21-phase2e-2-implementation.plan.md`(17 个 Task) * 11 项关键决策已锁定(D01-D11),详见设计文档 §三 * **重要修订**:`meeting_rooms.password` → `password_hash`(bcrypt),新增 `meeting_chats` 表 + `ended_reason` / `left_reason` 字段,新增 `echo:meeting:invite:{token}` / `host_grace:{code}` Redis key @@ -36,6 +36,7 @@ alwaysApply: true * **Task 2 ✅ 9 个内部 REST API 完成(2026-04-21,含代码审查修复)**:`media-server/` 落地 Router/Transport/Producer/Consumer 四类资源的 9 个接口,全部挂 `/internal/v1/*` 前缀;**zod 手动 parse + 全局 errorHandler** 方案(不引入 fastify-type-provider-zod 避免 zod v4 依赖冲突);`AppError` 统一错误码(`NOT_FOUND`/`CONFLICT`/`CAN_NOT_CONSUME`/`ROUTER_LIMIT_EXCEEDED`/`MEDIASOUP_ERROR`)+ `VALIDATION_ERROR`/`UNAUTHORIZED`/`INTERNAL_ERROR`;所有资源用 **Map + `observer.once('close')` 自清理**,`producerclose` 级联关闭下游 consumer;Consumer 强制 `paused:true` 创建 + `/resume` 独立接口;direction 强约束(recv transport 拒 produce、send transport 拒 consume)。**code-reviewer 子代理"有条件通过"**,2 Major + 4 高价值 Minor 当场修复:(1) M1 `_clearXxxMap` 新增 `assertTestOnly` 守卫(生产误调用直接抛错);(2) M2 新增 `src/schemas/rtp.ts` 对 `rtpParameters` / `rtpCapabilities` 做 codecs 浅层校验(mimeType/clockRate/payloadType 必填、codecs 数组 ≥1),消除 `as unknown as` 双跳断言;(3) m1 `connectTransport` 改乐观锁(先置位再 await);(4) m2 改读 `consumer.producerPaused`;(5) m3 `producerclose` 改为 `once`;(6) m5 `internal-auth` 改为反向白名单 `PRIVATE_PATH_PREFIXES = ['/internal/']`(默认开放)。**65 个 vitest 测试全过(~1s)、覆盖率 82.87%/75.83%/91.3%/82.87%**(stmts/branches/funcs/lines,较首版 +2pp);9 接口 happy path + 6 类错误路径人工 curl 全部按预期返回(201/200/400/401/404/409)。产出:`src/schemas/*`(6 文件,新增 `rtp.ts`)+ `src/services/*`(4 文件)+ `src/routes/*`(4 文件)+ `src/middlewares/{error-handler,internal-auth}.ts` + `src/utils/{errors,test-guard}.ts` + `src/mediasoup/codecs.ts` + `vitest.config.ts` + `tests/*`(8 spec 文件,含 `test-guard.spec.ts`)。余下 Minor/Nits(m4/m6~m10、n1~n10)登记至 Task 16 收尾清单 * **Task 3 ✅ Go meeting 模块数据库 DDL + Model + DAO 完成(2026-04-21)**:三张持久化表(`meeting_rooms` / `meeting_participants` / `meeting_chats`)落地 PostgreSQL,DDL 同时写入 `init.sql`(全量初始化)与 `phase2e2_migration.sql`(幂等增量升级)。Go 侧 `backend/go-service/app/meeting/{model,dao}` + 统一常量 `app/constants/meeting.go`:3 个 model + 3 个 DAO(`meeting_room_dao.go` 9 方法 / `meeting_participant_dao.go` 11 方法 / `meeting_chat_dao.go` 4 方法),共 24 个持久化方法;`JoinRoom` 事务内复用离会后的旧记录(`left_at=NULL,joined_at=NOW,duration=0`)避免审计表污染;`LeaveRoom` 用 `EXTRACT(EPOCH FROM (? - joined_at))::INT` 走 DB 时间防跨时区漂移;`TransferHost` 事务链(`role=1→0` + `role=0→1`);`FindActiveByUser` 用 JOIN 校验用户单点参会;`MarkEnded` 乐观锁防重复覆盖 `ended_reason`;`ListExpiredForCleanup` 供后续清理任务批量扫描。**关键风格修正(偏离实施计划草案)**:按 `project-context` 第 11 条「代码风格全局一致(最高优先级)」,常量归入 `app/constants/meeting.go` 单文件(与 `group.go`/`notify.go` 同构),而非草案的 `app/meeting/constants/*.go`;时间字段统一 `TIMESTAMP(0)` 取代草案的 `TIMESTAMPTZ` 对齐项目所有现有表;冗余 `idx_meeting_rooms_code` 移除(`room_code UNIQUE` 已自动建索引)。验证:`go build ./...` / `go vet ./...` / `ReadLints` 零错误;psql 集成脚本跑通 8 场景(CRUD + 双 UNIQUE 约束 + 主持人转让事务 + `duration=10s` 精确匹配 + CASCADE 清零)。延续项目 Go 侧"零 `_test.go`"风格(用代码审查 + psql 真库验证 + Playwright E2E 三层守护) * **Task 4 ✅ Go meeting 模块 service/controller/router 骨架完成(2026-04-21)**:`app/meeting/` 补齐 service/controller/router/provider 四件套,接口 → 实现按设计文档 §5.3 一一对齐。产出:(1) `service/interfaces.go` 定义 3 个外部依赖接口(`NotifyPusher` / `UserInfoResolver` / `OnlineChecker`),解耦 notify/contact/ws 模块避免循环依赖;`OnlineChecker.IsOnline` 签名与现存 `ws.OnlineService` 一致(返回单 `bool`);(2) `service/meeting_service.go` 声明 `MeetingService` + 8 个 sentinel error + 17 个业务方法空实现,全部返回 `ErrNotImplemented`;(3) `controller/meeting_controller.go` 12 个 Gin 处理器 + `responseNotImplemented`(501)+ `requireUserID` 辅助;(4) `router.go` 12 条路由挂载到 `/api/v1/meeting/*` 并统一套 `jwtAuth` 中间件;(5) `provider.go` 定义 `MeetingSet = wire.NewSet(DAO×3, Service, Controller)`;(6) 全局 `app/provider/wire.go` 挂入 `meetingApp.MeetingSet` + 3 条 `wire.Bind`(`NotifyPusher→NotifyService` / `UserInfoResolver→FriendshipDAO` / `OnlineChecker→ws.OnlineService`);(7) `app/provider/provider.go` `App` 加 `MeetingService/MeetingController` 字段;(8) `router/router.go` 调用 `meetingApp.RegisterRoutes`。**顺手修复存量 bug**:`admin/provider.go` 补齐 `MessageManage{DAO,Service,Controller}` 三个 provider,解决旧版 `wire` 重生成报"no provider found"的遗留问题。验证:`go build ./...` / `go vet ./...` / `wire ./app/provider` 全绿;`GIN_MODE=debug` 启动 server 日志打印全部 12 条 `[GIN-debug] ... meeting/controller.(*MeetingController).Xxx-fm`;curl 无 token 打 3 条代表性路由均返回 401 `缺少认证信息`,JWT 中间件生效 + * **Task 6 ✅ WebSocket 信令协议 13 事件全量落地(2026-04-21)**:`meeting.*` 事件族从 Task 5 的 `PublishToUser` 循环升级为完整 WS 信令协议。产出:(1) `app/constants/meeting.go` WS 事件常量与设计 §6.3 对齐(3 房间 + 5 成员 + 5 媒体 + 1 聊天)+ `MeetingWSClientEvents` 白名单限制客户端仅能发起 8 个 C→S 事件(防恶意客户端伪造 `room.ended` 等广播);(2) `service/interfaces.go` 扩容 `MediaOrchestrator` 至 9 方法(Router/Transport/Producer/Consumer 全生命周期)+ 配套 5 个 DTO + `NoopMediaOrchestrator` 9 占位实现;(3) **新建 `service/meeting_broadcaster.go`**(75 行)抽离 `BroadcastToMeeting` / `PublishToUser` 统一广播层,REST + WS 共用;(4) `service/meeting_service.go` 12 方法改调 `broadcaster.*` 不再直连 `ws.PubSub`;(5) **新建 `service/meeting_signal_service.go`**(430 行)承载 8 C→S 事件业务:`OnRoomJoin`/`OnRoomLeave`/`OnMemberStateChanged`/`OnTransportCreate/Connect`/`OnProduceStart`/`OnConsumeStart`/`OnProducerClose`;**Redis 资源追踪** `echo:meeting:resources:{room_id}:{user_id}`(Set,TTL 1 小时)+ `cleanupUserResources`(WS 断开钩子遍历清理 transport/producer/consumer)+ host 权限校验(非 host 改他人状态返回 `-1 仅主持人可执行此操作`);(6) **新建 `controller/meeting_ws_handler.go`**(200 行)薄层:构造时 `hub.RegisterEvent` 注册 8 C→S 事件,每 handler 仅 JSON 反序列化 + 调 `signalSvc.On*` + 构造 ACK;(7) `provider.go` + `wire_gen.go` 扩充三个新 provider;(8) `docs/api/frontend/meeting.md` 追加 §WebSocket 信令协议(200 行):16 事件总览 + 8 C→S 完整契约 + S→C 广播契约 + 架构 + 错误处理。**关键设计决策**:(a) Broadcaster 单独抽层让 Task 5 REST 广播与 Task 6 WS 广播复用同一对象;(b) C→S 白名单防伪造;(c) Redis Set 追踪媒体资源即使 Go 进程崩溃也不泄漏 mediasoup 端;(d) MediaOrchestrator 接口完整化让 Task 7 只需替换绑定不改 signal service / handler 代码;(e) ACK `code=0/-1` 与 REST 领域错误口径完全一致;(f) `meeting.room.leave` 只清 WS 资源不改 participant 表(真正离会仍需 REST `/leave`),允许客户端 WS 重连刷新 transport 不退会。验证:`go build` / `go vet` / `wire` 全绿;端到端 WS 冒烟 `/tmp/meeting_ws_t6_test.mjs` **PASS=18 / FAIL=0**,覆盖 8 C→S 白名单事件 + 3 S→C 广播 + 3 错误路径(非 host 越权/不存在会议号/WS leave 资源清理)。测试期间修复 `waitEvent` 死循环 bug(pop→push 自循环)改为 stash 缓冲区归还模式 * **Task 5 ✅ Go meeting 模块 12 个 REST 接口业务逻辑全量落地(2026-04-21)**:`MeetingService` + `MeetingController` 从 501 占位升级为完整实现,对接 PostgreSQL / Redis / NotifyPusher / PubSub / MediaOrchestrator。产出:(1) `app/dto/meeting_dto.go` 13 个 DTO(3 基础 + 10 请求/响应);(2) `pkg/utils/meeting_code.go` 生成 9 位 `XXX-XXX-XXX` 会议号 + 32 位 hex 邀请令牌;(3) `service/meeting_service.go` 12 业务方法 + 11 领域错误 + `assertIsActiveParticipant`/`assertIsHost`/`generateUniqueRoomCode`/`broadcastToActiveParticipants` 辅助;(4) `controller/meeting_controller.go` 12 Gin 处理器 + `handleError` 领域错误 → HTTP 映射 + `roomToDTO`/`participantToDTO`/`chatToDTO` 转换;(5) `service/interfaces.go` 新增 `MediaOrchestrator` 接口 + `NoopMediaOrchestrator` 占位(Task 7 替换);(6) `provider.go` 注册 `NoopMediaOrchestrator`;(7) `router.go` 路径对齐设计 `GET /rooms/mine` + `POST /invite-tokens/:token/redeem`;(8) `docs/api/frontend/meeting.md` 重写为 280 行的 12 接口完整文档。**DAO 契约修复**:`meeting_room_dao.GetByID/GetByCode` + `meeting_participant_dao.GetByRoomAndUser/FindActiveByUser` 全部将 `gorm.ErrRecordNotFound` 转换为 `(nil, nil)`,由 service 统一 `result == nil` 判定,消除 500 误报。**关键设计决策**:(a) 单点参会用 `meeting_participants` JOIN `meeting_rooms.status != 2` 判断;(b) 密码限流用 `echo:meeting:pwd:fail:{code}:{user_id}` 5 次锁 10 分钟(`ErrMeetingPasswordLocked`);(c) host 离会若还有其他活跃成员自动转让给"最早加入者"并广播 `meeting.host.changed`;(d) 邀请 token 不返回给调用方,仅通过 `NotifyPusher.PushBatch.Extra.invite_token` 定向下发;兑换后保留 60 秒冗余由 Redis TTL 自然过期;(e) 创建类接口 201、动作类 200、领域错误按 404/403/400 三档映射。**Stub 策略**:`MediaOrchestrator.CreateRouter/CloseRouter` 当前 Noop(Task 7 接入 HTTPMediaOrchestrator);WS 广播暂用 `pubsub.PublishToUser` 逐人循环(Task 6 封装为 `BroadcastToMeeting` 无感替换);`NotifyPusher.PushBatch` 复用 Phase 2e-1 成果。验证:`go build` / `go vet` / `wire` 全绿;启动 server + 3 用户端到端脚本 `/tmp/meeting_t5_test.sh` **PASS=19 / FAIL=0**,覆盖 12 接口 happy path + 5 类错误路径(密码错/房间不存在/单点参会冲突/非 host 越权/邀请链接失效) - 2e-3 会议增强(7-10 天)📋 待开发:预约会议(`type=2`)+ 定时提醒(`meeting_reminder`)+ 等候室/锁定会议 + 设备预览高级参数(降噪/回声/虚拟背景) - 分支:`feature/phase2e-2-meeting-mvp`(Phase 2e-2 专用,从 `feature/phase2c-group-read-receipt` 衍生) diff --git a/backend/go-service/app/constants/meeting.go b/backend/go-service/app/constants/meeting.go index 243469d..ee0a98d 100644 --- a/backend/go-service/app/constants/meeting.go +++ b/backend/go-service/app/constants/meeting.go @@ -92,23 +92,43 @@ const ( MeetingChatRetentionHours = 24 // 会议聊天保留时长(会议结束后) ) -// 会议 WS 事件常量(设计文档 §6.3,11 个事件) +// 会议 WS 事件常量(设计文档 §6.3) // 命名格式:meeting.{domain}.{action} +// 方向说明:C→S(客户端发往服务端,需 ACK);S→C(服务端广播/定向推送,不需 ACK) const ( - // 房间级事件 - MeetingWSEventRoomEnded = "meeting.room.ended" // 会议已结束 - MeetingWSEventRoomLocked = "meeting.room.locked" // 会议被锁定(等候室,Phase 2e-3) - MeetingWSEventRoomHostChange = "meeting.room.host.changed" // 主持人变更 + // ====== 房间组(3 个)====== + MeetingWSEventRoomJoin = "meeting.room.join" // C→S:加入会议后的"宣告在线",绑定 WS 连接 ↔ roomCode + MeetingWSEventRoomLeave = "meeting.room.leave" // C→S:主动离会(等价 REST leave,但保留信令触发入口) + MeetingWSEventRoomEnded = "meeting.room.ended" // S→C:会议已结束(广播全员) - // 成员级事件 - MeetingWSEventMemberJoined = "meeting.member.joined" // 新成员加入 - MeetingWSEventMemberLeft = "meeting.member.left" // 成员离开 - MeetingWSEventMemberStateChange = "meeting.member.state.changed" // 麦克风 / 摄像头状态变化 - MeetingWSEventMemberKicked = "meeting.member.kicked" // 被移出 + // ====== 成员组(5 个)====== + MeetingWSEventMemberJoined = "meeting.member.joined" // S→C:新成员加入 + MeetingWSEventMemberLeft = "meeting.member.left" // S→C:成员离开 + MeetingWSEventMemberStateChange = "meeting.member.state.changed" // 双向:麦克风 / 摄像头状态变化(host 可带 target_user_id 强制静音他人) + MeetingWSEventMemberKicked = "meeting.member.kicked" // S→C:被移出会议(定向发给被踢者) + MeetingWSEventMemberProducerNew = "meeting.member.producer.new" // S→C:有新 producer 可订阅(produce.start 成功后广播) + MeetingWSEventHostChanged = "meeting.host.changed" // S→C:主持人变更 - // 媒体级事件(Go 发起,驱动客户端 mediasoup-client 订阅/取消) - MeetingWSEventProducerNew = "meeting.producer.new" // 有新 producer 可订阅 - MeetingWSEventProducerClosed = "meeting.producer.closed" // producer 已关闭 - MeetingWSEventConsumerResumed = "meeting.consumer.resumed" // consumer 已恢复(由服务端触发) - MeetingWSEventChatMessage = "meeting.chat.message" // 会议内文字聊天 + // ====== 媒体组(5 个,mediasoup signaling 桥接)====== + MeetingWSEventTransportCreate = "meeting.transport.create" // C→S→Node:创建 Transport + MeetingWSEventTransportConnect = "meeting.transport.connect" // C→S→Node:Transport DTLS 握手 + MeetingWSEventProduceStart = "meeting.produce.start" // C→S→Node:创建 Producer(广播 producer.new) + MeetingWSEventConsumeStart = "meeting.consume.start" // C→S→Node:订阅远端 Producer,创建本地 Consumer + MeetingWSEventProducerClose = "meeting.producer.close" // C→S→Node:关闭自己的 Producer(推流停止) + + // ====== 补充事件(非 §6.3 核心 11 事件,但业务必须)====== + MeetingWSEventChatMessage = "meeting.chat.message" // S→C:会议内文字聊天广播(REST SendChat 触发) ) + +// MeetingWSClientEvents 客户端可主动发起的 WS 事件(C→S)白名单 +// 用于 WS 消息路由时快速校验事件名合法性 +var MeetingWSClientEvents = []string{ + MeetingWSEventRoomJoin, + MeetingWSEventRoomLeave, + MeetingWSEventMemberStateChange, + MeetingWSEventTransportCreate, + MeetingWSEventTransportConnect, + MeetingWSEventProduceStart, + MeetingWSEventConsumeStart, + MeetingWSEventProducerClose, +} diff --git a/backend/go-service/app/meeting/controller/meeting_ws_handler.go b/backend/go-service/app/meeting/controller/meeting_ws_handler.go new file mode 100644 index 0000000..2d817cb --- /dev/null +++ b/backend/go-service/app/meeting/controller/meeting_ws_handler.go @@ -0,0 +1,196 @@ +// Package controller 会议模块 HTTP / WS 入口 +package controller + +import ( + "context" + "encoding/json" + + "github.com/echochat/backend/app/constants" + "github.com/echochat/backend/app/meeting/service" + "github.com/echochat/backend/pkg/logs" + "github.com/echochat/backend/pkg/ws" + "go.uber.org/zap" +) + +// MeetingWSHandler 会议 WS 事件入口 +// Task 6 落地:注册设计 §6.3 的 11 个 meeting.* 事件到 Hub 路由表 +// 薄层设计: +// - 仅负责 JSON Unmarshal → 调 SignalService 业务方法 → 组装 ACK 响应 +// - 权限校验全部下沉到 SignalService(service 层统一 assertIsActiveParticipant / assertIsHost) +// - 广播副作用由 SignalService 自行发起(ACK 不携带广播数据) +type MeetingWSHandler struct { + signalSvc *service.MeetingSignalService + hub *ws.Hub +} + +// NewMeetingWSHandler 构造 WS Handler 并自动注册事件到 Hub +// 与 im.handler.EventHandler 保持一致的风格:构造时注册,启动时由 DI 容器持有 +func NewMeetingWSHandler(signalSvc *service.MeetingSignalService, hub *ws.Hub) *MeetingWSHandler { + h := &MeetingWSHandler{ + signalSvc: signalSvc, + hub: hub, + } + h.registerEvents() + return h +} + +// registerEvents 将所有 meeting.* C→S 事件注册到 Hub 路由表 +// S→C 事件(ended/joined/left/host.changed/producer.new 等)不在此处注册,由 broadcaster 发出 +func (h *MeetingWSHandler) registerEvents() { + // 房间组 + h.hub.RegisterEvent(constants.MeetingWSEventRoomJoin, h.handleRoomJoin) + h.hub.RegisterEvent(constants.MeetingWSEventRoomLeave, h.handleRoomLeave) + // 成员组 + h.hub.RegisterEvent(constants.MeetingWSEventMemberStateChange, h.handleMemberStateChange) + // 媒体组(5 个) + h.hub.RegisterEvent(constants.MeetingWSEventTransportCreate, h.handleTransportCreate) + h.hub.RegisterEvent(constants.MeetingWSEventTransportConnect, h.handleTransportConnect) + h.hub.RegisterEvent(constants.MeetingWSEventProduceStart, h.handleProduceStart) + h.hub.RegisterEvent(constants.MeetingWSEventConsumeStart, h.handleConsumeStart) + h.hub.RegisterEvent(constants.MeetingWSEventProducerClose, h.handleProducerClose) +} + +// ====== 通用工具 ====== + +// simpleRoomPayload 房间组仅需 room_code 的请求体 +type simpleRoomPayload struct { + RoomCode string `json:"room_code"` +} + +// sendACK 统一发送 ACK 响应 +func (h *MeetingWSHandler) sendACK(client *ws.Client, msg *ws.Message, code int, message string, data interface{}) { + resp := ws.NewResponse(msg.Event, msg.Seq, code, message, data) + bytes, err := ws.MarshalResponse(resp) + if err != nil { + logs.Error(context.Background(), "controller.meeting_ws_handler.sendACK", "序列化 ACK 失败", + zap.String("event", msg.Event), zap.Error(err)) + return + } + client.Send(bytes) +} + +// unmarshal 通用反序列化 + 错误 ACK +func (h *MeetingWSHandler) unmarshal(client *ws.Client, msg *ws.Message, target interface{}) bool { + if err := json.Unmarshal(msg.Data, target); err != nil { + logs.Warn(nil, "controller.meeting_ws_handler.unmarshal", "反序列化失败", + zap.String("event", msg.Event), + zap.Int64("user_id", client.UserID), + zap.Error(err)) + h.sendACK(client, msg, -1, "请求参数格式错误", nil) + return false + } + return true +} + +// ====== 房间组 ====== + +// handleRoomJoin 处理 meeting.room.join +func (h *MeetingWSHandler) handleRoomJoin(client *ws.Client, msg *ws.Message) { + var payload simpleRoomPayload + if !h.unmarshal(client, msg, &payload) { + return + } + if err := h.signalSvc.OnRoomJoin(context.Background(), client.UserID, payload.RoomCode); err != nil { + h.sendACK(client, msg, -1, err.Error(), nil) + return + } + h.sendACK(client, msg, 0, "ok", nil) +} + +// handleRoomLeave 处理 meeting.room.leave +func (h *MeetingWSHandler) handleRoomLeave(client *ws.Client, msg *ws.Message) { + var payload simpleRoomPayload + if !h.unmarshal(client, msg, &payload) { + return + } + if err := h.signalSvc.OnRoomLeave(context.Background(), client.UserID, payload.RoomCode); err != nil { + h.sendACK(client, msg, -1, err.Error(), nil) + return + } + h.sendACK(client, msg, 0, "ok", nil) +} + +// ====== 成员组 ====== + +// handleMemberStateChange 处理 meeting.member.state.changed +func (h *MeetingWSHandler) handleMemberStateChange(client *ws.Client, msg *ws.Message) { + var payload service.MemberStateChangePayload + if !h.unmarshal(client, msg, &payload) { + return + } + if err := h.signalSvc.OnMemberStateChanged(context.Background(), client.UserID, &payload); err != nil { + h.sendACK(client, msg, -1, err.Error(), nil) + return + } + h.sendACK(client, msg, 0, "ok", nil) +} + +// ====== 媒体组 ====== + +// handleTransportCreate 处理 meeting.transport.create +func (h *MeetingWSHandler) handleTransportCreate(client *ws.Client, msg *ws.Message) { + var payload service.TransportCreatePayload + if !h.unmarshal(client, msg, &payload) { + return + } + info, err := h.signalSvc.OnTransportCreate(context.Background(), client.UserID, &payload) + if err != nil { + h.sendACK(client, msg, -1, err.Error(), nil) + return + } + h.sendACK(client, msg, 0, "ok", info) +} + +// handleTransportConnect 处理 meeting.transport.connect +func (h *MeetingWSHandler) handleTransportConnect(client *ws.Client, msg *ws.Message) { + var payload service.TransportConnectPayload + if !h.unmarshal(client, msg, &payload) { + return + } + if err := h.signalSvc.OnTransportConnect(context.Background(), client.UserID, &payload); err != nil { + h.sendACK(client, msg, -1, err.Error(), nil) + return + } + h.sendACK(client, msg, 0, "ok", nil) +} + +// handleProduceStart 处理 meeting.produce.start +func (h *MeetingWSHandler) handleProduceStart(client *ws.Client, msg *ws.Message) { + var payload service.ProduceStartPayload + if !h.unmarshal(client, msg, &payload) { + return + } + result, err := h.signalSvc.OnProduceStart(context.Background(), client.UserID, &payload) + if err != nil { + h.sendACK(client, msg, -1, err.Error(), nil) + return + } + h.sendACK(client, msg, 0, "ok", result) +} + +// handleConsumeStart 处理 meeting.consume.start +func (h *MeetingWSHandler) handleConsumeStart(client *ws.Client, msg *ws.Message) { + var payload service.ConsumeStartPayload + if !h.unmarshal(client, msg, &payload) { + return + } + info, err := h.signalSvc.OnConsumeStart(context.Background(), client.UserID, &payload) + if err != nil { + h.sendACK(client, msg, -1, err.Error(), nil) + return + } + h.sendACK(client, msg, 0, "ok", info) +} + +// handleProducerClose 处理 meeting.producer.close +func (h *MeetingWSHandler) handleProducerClose(client *ws.Client, msg *ws.Message) { + var payload service.ProducerClosePayload + if !h.unmarshal(client, msg, &payload) { + return + } + if err := h.signalSvc.OnProducerClose(context.Background(), client.UserID, &payload); err != nil { + h.sendACK(client, msg, -1, err.Error(), nil) + return + } + h.sendACK(client, msg, 0, "ok", nil) +} diff --git a/backend/go-service/app/meeting/provider.go b/backend/go-service/app/meeting/provider.go index e330d60..8be4495 100644 --- a/backend/go-service/app/meeting/provider.go +++ b/backend/go-service/app/meeting/provider.go @@ -9,15 +9,21 @@ import ( // MeetingSet 会议模块 Wire Provider Set // 对外暴露: -// - *service.MeetingService —— 业务服务,供未来 ws handler / media-server 回调使用 +// - *service.MeetingService —— REST 业务服务(Task 5) +// - *service.MeetingBroadcaster —— WS 广播中枢(Task 6 新增,供 SignalService 复用) +// - *service.MeetingSignalService —— WS 信令事件业务逻辑(Task 6) // - *controller.MeetingController —— REST API 控制器 +// - *controller.MeetingWSHandler —— WS 事件 Handler(Task 6,启动时自动注册到 Hub) // 依赖的接口 NotifyPusher / UserInfoResolver / OnlineChecker 由上游 wire.Bind 绑定具体实现 var MeetingSet = wire.NewSet( dao.NewMeetingRoomDAO, dao.NewMeetingParticipantDAO, dao.NewMeetingChatDAO, + service.NewMeetingBroadcaster, service.NewMeetingService, + service.NewMeetingSignalService, controller.NewMeetingController, + controller.NewMeetingWSHandler, // MediaOrchestrator 目前使用 Noop 实现(Task 7 将替换为 node_client.NodeClient) service.NewNoopMediaOrchestrator, diff --git a/backend/go-service/app/meeting/service/interfaces.go b/backend/go-service/app/meeting/service/interfaces.go index c2d4b4b..93c0b20 100644 --- a/backend/go-service/app/meeting/service/interfaces.go +++ b/backend/go-service/app/meeting/service/interfaces.go @@ -3,6 +3,7 @@ package service import ( "context" + "encoding/json" authModel "github.com/echochat/backend/app/auth/model" notifyService "github.com/echochat/backend/app/notify/service" @@ -30,22 +31,82 @@ type OnlineChecker interface { IsOnline(ctx context.Context, userID int64) bool } -// MediaOrchestrator 媒体服务器编排接口 -// Task 7 落地 Go → Node media-server HTTP Client 后由 node_client.NodeClient 实现 -// Task 5 阶段使用 NoopMediaOrchestrator 占位,仅返回假的 RouterID,不产生真实媒体资源 -// 语义:会议房间在 Go 侧落库后,通过该接口驱动 Node 端 mediasoup Router 创建与销毁 -type MediaOrchestrator interface { - // CreateRouter 为会议房间创建 mediasoup Router - // 入参:room_code 作为 Node 端聚合键;返回 Router ID 与推荐编解码参数(Task 7 对接时补充) - CreateRouter(ctx context.Context, roomCode string) (routerID string, err error) +// ====== MediaOrchestrator:Go → Node media-server 统一封装(设计 §6.6)====== - // CloseRouter 关闭会议房间对应的 mediasoup Router 及其下所有 transport/producer/consumer - // 幂等:重复关闭不返回错误 - CloseRouter(ctx context.Context, roomCode string) error +// TransportInfo Node 端创建的 Transport 元信息 +// 字段与 mediasoup-client 的 Device.createSendTransport / createRecvTransport 入参兼容 +type TransportInfo struct { + ID string `json:"id"` // Transport ID + IceParameters json.RawMessage `json:"iceParameters"` // mediasoup ICE 参数 + IceCandidates json.RawMessage `json:"iceCandidates"` // ICE 候选列表 + DtlsParameters json.RawMessage `json:"dtlsParameters"` // DTLS 指纹参数 + SctpParameters json.RawMessage `json:"sctpParameters,omitempty"` } -// NoopMediaOrchestrator 占位实现:Task 5 完成生命周期接口时使用 -// 返回伪造的 RouterID,所有操作仅写日志不调用 Node +// ConsumerInfo Node 端创建的 Consumer 元信息 +type ConsumerInfo struct { + ID string `json:"id"` + ProducerID string `json:"producerId"` + Kind string `json:"kind"` // "audio" | "video" + RtpParameters json.RawMessage `json:"rtpParameters"` + Type string `json:"type,omitempty"` // "simple" | "simulcast" | "svc" + ProducerPaused bool `json:"producerPaused,omitempty"` +} + +// CreateTransportReq 创建 Transport 请求 +type CreateTransportReq struct { + RoomCode string `json:"roomCode"` + UserID int64 `json:"userId"` + Direction string `json:"direction"` // "send" | "recv" +} + +// CreateProducerReq 创建 Producer 请求 +type CreateProducerReq struct { + RoomCode string `json:"roomCode"` + UserID int64 `json:"userId"` + TransportID string `json:"transportId"` + Kind string `json:"kind"` // "audio" | "video" + RtpParameters json.RawMessage `json:"rtpParameters"` // 直接转发给 Node,由其做结构校验 +} + +// CreateConsumerReq 创建 Consumer 请求 +type CreateConsumerReq struct { + RoomCode string `json:"roomCode"` + UserID int64 `json:"userId"` + TransportID string `json:"transportId"` + ProducerID string `json:"producerId"` + RtpCapabilities json.RawMessage `json:"rtpCapabilities"` +} + +// MediaOrchestrator 媒体服务器编排接口(设计 §6.6 NodeClient) +// Task 7 落地 Go → Node media-server HTTP Client 后由 node_client.NodeClient 实现 +// Task 5/6 阶段使用 NoopMediaOrchestrator 占位: +// - CreateRouter 返回 "noop-router-{code}";其他方法返回可解析的占位数据,用于 WS 信令链路自测 +// - 所有方法均幂等:重复调用不报错,符合 WS 信令重试语义 +type MediaOrchestrator interface { + // CreateRouter 为会议房间创建 mediasoup Router + CreateRouter(ctx context.Context, roomCode string) (routerID string, err error) + // CloseRouter 关闭房间对应 Router 及其下所有资源(幂等) + CloseRouter(ctx context.Context, roomCode string) error + + // CreateTransport 为用户创建 send/recv WebRTC Transport + CreateTransport(ctx context.Context, req *CreateTransportReq) (*TransportInfo, error) + // ConnectTransport Transport DTLS 握手(幂等:重复 connect 对已连接 transport 视为成功) + ConnectTransport(ctx context.Context, transportID string, dtlsParameters json.RawMessage) error + + // CreateProducer 在指定 send Transport 上创建 Producer + CreateProducer(ctx context.Context, req *CreateProducerReq) (producerID string, err error) + // CloseProducer 关闭指定 Producer(幂等) + CloseProducer(ctx context.Context, producerID string) error + + // CreateConsumer 在指定 recv Transport 上创建 Consumer(订阅远端 Producer) + CreateConsumer(ctx context.Context, req *CreateConsumerReq) (*ConsumerInfo, error) + // CloseConsumer 关闭指定 Consumer(幂等) + CloseConsumer(ctx context.Context, consumerID string) error +} + +// NoopMediaOrchestrator 占位实现:Task 7 完成前使用 +// 返回伪造的 ID 与固定占位数据(JSON:空对象 / 空数组),所有操作仅写日志不调用 Node // Task 7 完成后全局 wire 切换到真实 NodeClient 实现 type NoopMediaOrchestrator struct{} @@ -55,7 +116,6 @@ func NewNoopMediaOrchestrator() *NoopMediaOrchestrator { } // CreateRouter 返回以 "noop-router-" 为前缀的伪造 RouterID -// 调用方可据此区分真实 / 占位实现,便于调试与切换 func (n *NoopMediaOrchestrator) CreateRouter(_ context.Context, roomCode string) (string, error) { return "noop-router-" + roomCode, nil } @@ -64,3 +124,44 @@ func (n *NoopMediaOrchestrator) CreateRouter(_ context.Context, roomCode string) func (n *NoopMediaOrchestrator) CloseRouter(_ context.Context, _ string) error { return nil } + +// CreateTransport 占位:返回以 "noop-transport-" 为前缀的伪造 ID,带最小合法 JSON 结构 +func (n *NoopMediaOrchestrator) CreateTransport(_ context.Context, req *CreateTransportReq) (*TransportInfo, error) { + return &TransportInfo{ + ID: "noop-transport-" + req.Direction, + IceParameters: json.RawMessage(`{}`), + IceCandidates: json.RawMessage(`[]`), + DtlsParameters: json.RawMessage(`{}`), + }, nil +} + +// ConnectTransport 占位:直接返回 nil +func (n *NoopMediaOrchestrator) ConnectTransport(_ context.Context, _ string, _ json.RawMessage) error { + return nil +} + +// CreateProducer 占位:返回伪造 ID +func (n *NoopMediaOrchestrator) CreateProducer(_ context.Context, req *CreateProducerReq) (string, error) { + return "noop-producer-" + req.Kind, nil +} + +// CloseProducer 占位:直接返回 nil +func (n *NoopMediaOrchestrator) CloseProducer(_ context.Context, _ string) error { + return nil +} + +// CreateConsumer 占位:返回伪造的 ConsumerInfo(含目标 producerID 回填) +func (n *NoopMediaOrchestrator) CreateConsumer(_ context.Context, req *CreateConsumerReq) (*ConsumerInfo, error) { + return &ConsumerInfo{ + ID: "noop-consumer-" + req.ProducerID, + ProducerID: req.ProducerID, + Kind: "video", + RtpParameters: json.RawMessage(`{}`), + Type: "simple", + }, nil +} + +// CloseConsumer 占位:直接返回 nil +func (n *NoopMediaOrchestrator) CloseConsumer(_ context.Context, _ string) error { + return nil +} diff --git a/backend/go-service/app/meeting/service/meeting_broadcaster.go b/backend/go-service/app/meeting/service/meeting_broadcaster.go new file mode 100644 index 0000000..da076ee --- /dev/null +++ b/backend/go-service/app/meeting/service/meeting_broadcaster.go @@ -0,0 +1,91 @@ +// Package service 提供 meeting 模块的业务逻辑 +package service + +import ( + "context" + + "github.com/echochat/backend/app/meeting/dao" + "github.com/echochat/backend/pkg/logs" + "github.com/echochat/backend/pkg/ws" + "go.uber.org/zap" +) + +// MeetingBroadcaster 会议 WS 广播中枢 +// Task 6 引入,替代原 MeetingService.broadcastToActiveParticipants 内联实现 +// 语义: +// - BroadcastToMeeting(ctx, roomID, event, payload, excludeUserIDs...) 向房间内所有活跃成员广播 +// - PublishToUser(ctx, userID, event, payload) 定向推送给指定用户(如 meeting.member.kicked 定向通知被踢者) +// 通过 PubSub.Publish 跨实例传递;本地 Hub 若订阅了目标用户频道则自动路由到 WS 连接 +type MeetingBroadcaster struct { + participantDAO *dao.MeetingParticipantDAO + pubsub *ws.PubSub +} + +// NewMeetingBroadcaster 构造广播中枢 +func NewMeetingBroadcaster(participantDAO *dao.MeetingParticipantDAO, pubsub *ws.PubSub) *MeetingBroadcaster { + return &MeetingBroadcaster{ + participantDAO: participantDAO, + pubsub: pubsub, + } +} + +// BroadcastToMeeting 向房间内所有活跃参会者广播 WS 事件 +// excludeUserIDs 中的用户将被跳过(通常排除发送者本人,避免"自己收到自己的事件") +// 非阻塞:单个用户发送失败只记录 WARN 日志,不中断循环 +func (b *MeetingBroadcaster) BroadcastToMeeting(ctx context.Context, roomID int64, event string, data interface{}, excludeUserIDs ...int64) { + funcName := "service.meeting_broadcaster.BroadcastToMeeting" + + participants, err := b.participantDAO.ListActiveByRoom(ctx, roomID) + if err != nil { + logs.Warn(ctx, funcName, "拉取活跃参会者失败", + zap.Int64("room_id", roomID), + zap.String("event", event), + zap.Error(err)) + return + } + if len(participants) == 0 { + return + } + + exclude := make(map[int64]struct{}, len(excludeUserIDs)) + for _, id := range excludeUserIDs { + exclude[id] = struct{}{} + } + + msg := ws.NewPushMessage(event, data) + pushed := 0 + for _, p := range participants { + if _, skip := exclude[p.UserID]; skip { + continue + } + if err := b.pubsub.PublishToUser(ctx, p.UserID, msg); err != nil { + logs.Warn(ctx, funcName, "WS 广播单用户失败", + zap.Int64("user_id", p.UserID), + zap.String("event", event), + zap.Error(err)) + continue + } + pushed++ + } + + logs.Debug(ctx, funcName, "会议事件广播完成", + zap.Int64("room_id", roomID), + zap.String("event", event), + zap.Int("pushed", pushed), + zap.Int("total_active", len(participants))) +} + +// PublishToUser 定向推送给单个用户 +// 场景:meeting.member.kicked(被踢者收到的定向通知)、ACK 补偿推送等 +// 返回 error 让调用方根据业务语义决定是否需要感知失败 +func (b *MeetingBroadcaster) PublishToUser(ctx context.Context, userID int64, event string, data interface{}) error { + msg := ws.NewPushMessage(event, data) + if err := b.pubsub.PublishToUser(ctx, userID, msg); err != nil { + logs.Warn(ctx, "service.meeting_broadcaster.PublishToUser", "定向 WS 推送失败", + zap.Int64("user_id", userID), + zap.String("event", event), + zap.Error(err)) + return err + } + return nil +} diff --git a/backend/go-service/app/meeting/service/meeting_service.go b/backend/go-service/app/meeting/service/meeting_service.go index 8f9e7be..6b0e226 100644 --- a/backend/go-service/app/meeting/service/meeting_service.go +++ b/backend/go-service/app/meeting/service/meeting_service.go @@ -16,7 +16,6 @@ import ( notifyService "github.com/echochat/backend/app/notify/service" "github.com/echochat/backend/pkg/logs" "github.com/echochat/backend/pkg/utils" - "github.com/echochat/backend/pkg/ws" "github.com/redis/go-redis/v9" "go.uber.org/zap" "gorm.io/gorm" @@ -51,16 +50,17 @@ const ( // MeetingService 会议业务服务 // Task 5 完成:会议生命周期、主持人管理、邀请、会议内聊天 12 个 REST API 全部落地 -// Task 6 会在此基础上追加 WS 信令事件处理器;Task 7 会把 mediaOrchestrator 的 Noop 实现替换为真实 Node HTTP Client +// Task 6 新增:广播能力抽离到 MeetingBroadcaster;后续 WS 信令事件由 MeetingSignalService 处理 +// Task 7 会把 mediaOrchestrator 的 Noop 实现替换为真实 Node HTTP Client type MeetingService struct { roomDAO *dao.MeetingRoomDAO participantDAO *dao.MeetingParticipantDAO chatDAO *dao.MeetingChatDAO - db *gorm.DB - redis *redis.Client - pubsub *ws.PubSub + db *gorm.DB + redis *redis.Client + broadcaster *MeetingBroadcaster notifyPusher NotifyPusher userResolver UserInfoResolver onlineChecker OnlineChecker @@ -68,14 +68,15 @@ type MeetingService struct { } // NewMeetingService 创建 MeetingService 实例 -// 依赖通过构造函数注入;接口依赖由上游 Wire 绑定到具体实现(Task 7 之前 mediaOrchestrator 使用 NoopMediaOrchestrator) +// 依赖通过构造函数注入;接口依赖由上游 Wire 绑定到具体实现 +// Task 7 之前 mediaOrchestrator 使用 NoopMediaOrchestrator func NewMeetingService( roomDAO *dao.MeetingRoomDAO, participantDAO *dao.MeetingParticipantDAO, chatDAO *dao.MeetingChatDAO, db *gorm.DB, redis *redis.Client, - pubsub *ws.PubSub, + broadcaster *MeetingBroadcaster, notifyPusher NotifyPusher, userResolver UserInfoResolver, onlineChecker OnlineChecker, @@ -87,7 +88,7 @@ func NewMeetingService( chatDAO: chatDAO, db: db, redis: redis, - pubsub: pubsub, + broadcaster: broadcaster, notifyPusher: notifyPusher, userResolver: userResolver, onlineChecker: onlineChecker, @@ -139,28 +140,10 @@ func (s *MeetingService) generateUniqueRoomCode(ctx context.Context) (string, er } // broadcastToActiveParticipants 向房间内所有活跃参会者广播 WS 事件 -// 调用 PubSub.Publish 支持多实例;可选排除发送者自身(excludeUserIDs) -// 该辅助作为 Task 6 BroadcastToMeeting 的临时等效实现,签名保持兼容便于后续替换 +// Task 6 起实现已迁移至 MeetingBroadcaster.BroadcastToMeeting,本方法作为兼容壳保留 +// 以减少调用侧改动;未来可逐步替换为直接调用 s.broadcaster.BroadcastToMeeting func (s *MeetingService) broadcastToActiveParticipants(ctx context.Context, roomID int64, event string, data interface{}, excludeUserIDs ...int64) { - funcName := "service.meeting_service.broadcastToActiveParticipants" - participants, err := s.participantDAO.ListActiveByRoom(ctx, roomID) - if err != nil { - logs.Warn(ctx, funcName, "拉取活跃参会者失败", zap.Int64("room_id", roomID), zap.Error(err)) - return - } - exclude := make(map[int64]struct{}, len(excludeUserIDs)) - for _, id := range excludeUserIDs { - exclude[id] = struct{}{} - } - msg := ws.NewPushMessage(event, data) - for _, p := range participants { - if _, skip := exclude[p.UserID]; skip { - continue - } - if err := s.pubsub.PublishToUser(ctx, p.UserID, msg); err != nil { - logs.Warn(ctx, funcName, "WS 广播失败", zap.Int64("user_id", p.UserID), zap.String("event", event), zap.Error(err)) - } - } + s.broadcaster.BroadcastToMeeting(ctx, roomID, event, data, excludeUserIDs...) } // ====== 会议生命周期 ====== @@ -386,7 +369,7 @@ func (s *MeetingService) LeaveRoom(ctx context.Context, userID int64, code strin if uErr := s.roomDAO.UpdateHost(ctx, room.ID, newHost.UserID); uErr != nil { logs.Warn(ctx, funcName, "UpdateHost 失败", zap.Error(uErr)) } - go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventRoomHostChange, map[string]interface{}{ + go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventHostChanged, map[string]interface{}{ "room_code": code, "old_host_id": userID, "new_host_id": newHost.UserID, @@ -441,15 +424,14 @@ func (s *MeetingService) EndRoom(ctx context.Context, userID int64, code string) } payload := map[string]interface{}{ - "room_code": code, + "room_code": code, "ended_reason": constants.MeetingEndedReasonHostEnded, - "ended_at": now.Format("2006-01-02 15:04:05"), + "ended_at": now.Format("2006-01-02 15:04:05"), } - msg := ws.NewPushMessage(constants.MeetingWSEventRoomEnded, payload) + // activesBefore 已经是结束前的活跃成员快照,此时 participantDAO.ListActiveByRoom 查会返回空集合 + // 因此使用 broadcaster.PublishToUser 逐人定向推送(而非 BroadcastToMeeting 基于当前状态查库) for _, p := range activesBefore { - if err := s.pubsub.PublishToUser(ctx, p.UserID, msg); err != nil { - logs.Warn(ctx, funcName, "meeting.room.ended 广播失败", zap.Int64("user_id", p.UserID), zap.Error(err)) - } + _ = s.broadcaster.PublishToUser(ctx, p.UserID, constants.MeetingWSEventRoomEnded, payload) } if err := s.mediaOrchestrator.CloseRouter(ctx, code); err != nil { @@ -498,7 +480,7 @@ func (s *MeetingService) TransferHost(ctx context.Context, operatorID int64, cod return err } - go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventRoomHostChange, map[string]interface{}{ + go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventHostChanged, map[string]interface{}{ "room_code": code, "old_host_id": operatorID, "new_host_id": target.UserID, @@ -542,12 +524,11 @@ func (s *MeetingService) KickMember(ctx context.Context, operatorID int64, code return ErrNotInMeeting } - kickMsg := ws.NewPushMessage(constants.MeetingWSEventMemberKicked, map[string]interface{}{ + _ = s.broadcaster.PublishToUser(ctx, targetUserID, constants.MeetingWSEventMemberKicked, map[string]interface{}{ "room_code": code, "user_id": targetUserID, "by": operatorID, }) - _ = s.pubsub.PublishToUser(ctx, targetUserID, kickMsg) go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{ "room_code": code, "user_id": targetUserID, diff --git a/backend/go-service/app/meeting/service/meeting_signal_service.go b/backend/go-service/app/meeting/service/meeting_signal_service.go new file mode 100644 index 0000000..8b3921b --- /dev/null +++ b/backend/go-service/app/meeting/service/meeting_signal_service.go @@ -0,0 +1,393 @@ +// Package service 提供 meeting 模块的业务逻辑 +package service + +import ( + "context" + "encoding/json" + "fmt" + "time" + + "github.com/echochat/backend/app/constants" + "github.com/echochat/backend/app/meeting/dao" + "github.com/echochat/backend/app/meeting/model" + "github.com/echochat/backend/pkg/logs" + "github.com/redis/go-redis/v9" + "go.uber.org/zap" +) + +// MeetingSignalService WS 信令事件业务处理 +// Task 6 落地:处理设计 §6.3 的 11 个 meeting.* 事件,分 3 组: +// - 房间组(meeting.room.join/leave):绑定 WS 连接 ↔ roomCode,辅助断连清理 +// - 成员组(meeting.member.state.changed):麦/视频状态变更 + host 静音他人 +// - 媒体组(meeting.transport.*/produce.*/consume.*/producer.close): +// 对接 MediaOrchestrator,实现 mediasoup signaling 桥接 +// +// 所有方法统一返回 (ackData, error):error 非 nil 表示业务失败,由 Handler 映射为 ACK code=-1 +// 广播副作用(如 meeting.member.producer.new)在方法内部通过 broadcaster 发出,不进 ACK +type MeetingSignalService struct { + roomDAO *dao.MeetingRoomDAO + participantDAO *dao.MeetingParticipantDAO + redis *redis.Client + + broadcaster *MeetingBroadcaster + mediaOrchestrator MediaOrchestrator +} + +// NewMeetingSignalService 构造 WS 信令服务 +func NewMeetingSignalService( + roomDAO *dao.MeetingRoomDAO, + participantDAO *dao.MeetingParticipantDAO, + redis *redis.Client, + broadcaster *MeetingBroadcaster, + mediaOrchestrator MediaOrchestrator, +) *MeetingSignalService { + return &MeetingSignalService{ + roomDAO: roomDAO, + participantDAO: participantDAO, + redis: redis, + broadcaster: broadcaster, + mediaOrchestrator: mediaOrchestrator, + } +} + +// Redis 资源追踪 key(设计 §九 - 断线清理) +// Set 内元素格式:"transport:{id}" / "producer:{id}" / "consumer:{id}" +func resourceTrackKey(roomCode string, userID int64) string { + return fmt.Sprintf("echo:meeting:resource:%s:%d", roomCode, userID) +} + +// resourceTTL 单个用户资源追踪集合 TTL +// 设计:会议期间维持可达即可;若用户长期不活跃由断线清理接管 +const resourceTTL = time.Hour + +// trackResource 记录用户在会议中持有的媒体资源 ID +func (s *MeetingSignalService) trackResource(ctx context.Context, roomCode string, userID int64, kind, id string) { + key := resourceTrackKey(roomCode, userID) + member := kind + ":" + id + if err := s.redis.SAdd(ctx, key, member).Err(); err != nil { + logs.Warn(ctx, "service.meeting_signal_service.trackResource", "追踪媒体资源失败", + zap.String("key", key), zap.String("member", member), zap.Error(err)) + return + } + _ = s.redis.Expire(ctx, key, resourceTTL).Err() +} + +// untrackResource 从集合中移除资源 ID(关闭 producer/consumer 时) +func (s *MeetingSignalService) untrackResource(ctx context.Context, roomCode string, userID int64, kind, id string) { + key := resourceTrackKey(roomCode, userID) + member := kind + ":" + id + _ = s.redis.SRem(ctx, key, member).Err() +} + +// loadRoomAndParticipant 通用前置校验:拉取房间 + 确认用户是活跃参会者 +// 所有信令事件在进入业务前都要过这一关;返回的 *MeetingRoom 供后续广播使用 roomID +func (s *MeetingSignalService) loadRoomAndParticipant(ctx context.Context, roomCode string, userID int64) (*model.MeetingRoom, error) { + room, err := s.roomDAO.GetByCode(ctx, roomCode) + if err != nil { + return nil, err + } + if room == nil { + return nil, ErrMeetingNotFound + } + if room.Status == constants.MeetingStatusEnded { + return nil, ErrMeetingEnded + } + p, err := s.participantDAO.GetByRoomAndUser(ctx, room.ID, userID) + if err != nil { + return nil, err + } + if p == nil || !p.IsActive() { + return nil, ErrNotInMeeting + } + return room, nil +} + +// ========== 房间组(2 个 C→S)========== + +// OnRoomJoin 处理 meeting.room.join 事件 +// 语义:客户端 REST 加入会议成功后,通过 WS 宣告在线;服务端记录 userID ↔ roomCode 映射 +// 仅做存在性校验 + 心跳意义上的资源 key 刷新,不产生副作用 +func (s *MeetingSignalService) OnRoomJoin(ctx context.Context, userID int64, roomCode string) error { + room, err := s.loadRoomAndParticipant(ctx, roomCode, userID) + if err != nil { + return err + } + // 触发资源追踪 key 续期(空集合 TTL 续期无副作用) + key := resourceTrackKey(roomCode, userID) + _ = s.redis.Expire(ctx, key, resourceTTL).Err() + + logs.Info(ctx, "service.meeting_signal_service.OnRoomJoin", "用户宣告 WS 在线", + zap.String("room_code", roomCode), + zap.Int64("user_id", userID), + zap.Int64("room_id", room.ID)) + return nil +} + +// OnRoomLeave 处理 meeting.room.leave 事件 +// 语义:WS 层面的主动离会(等价 REST leave 但不强制要求落库事务; +// 当前实现:仅清理该用户在本会议的所有媒体资源(batch close producer/consumer)+ 广播 meeting.member.left +// 参会者表的 LeaveRoom 逻辑仍由 REST API 负责(避免 WS 并发引起 left_at 重复写入) +func (s *MeetingSignalService) OnRoomLeave(ctx context.Context, userID int64, roomCode string) error { + room, err := s.loadRoomAndParticipant(ctx, roomCode, userID) + if err != nil { + return err + } + s.cleanupUserResources(ctx, roomCode, userID) + + go s.broadcaster.BroadcastToMeeting(context.Background(), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{ + "room_code": roomCode, + "user_id": userID, + "reason": "ws_disconnect", + }, userID) + return nil +} + +// ========== 成员组(1 个双向)========== + +// MemberStateChangePayload meeting.member.state.changed 请求载荷 +// host 可通过 target_user_id 静音他人 / 关其摄像头;非 host 传该字段将被拒绝 +type MemberStateChangePayload struct { + RoomCode string `json:"room_code"` + TargetUserID int64 `json:"target_user_id,omitempty"` // 可选:host 强制他人状态 + AudioEnabled *bool `json:"audio_enabled,omitempty"` // nil 表示不改 + VideoEnabled *bool `json:"video_enabled,omitempty"` +} + +// OnMemberStateChanged 处理 meeting.member.state.changed 事件 +// 权限:操作自己无限制;操作他人必须是 host +// 行为:广播 meeting.member.state.changed 给房间其他成员(发起者自己不收到回显) +func (s *MeetingSignalService) OnMemberStateChanged(ctx context.Context, fromUserID int64, payload *MemberStateChangePayload) error { + room, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, fromUserID) + if err != nil { + return err + } + + targetID := fromUserID + if payload.TargetUserID != 0 && payload.TargetUserID != fromUserID { + if room.HostID != fromUserID { + return ErrNotMeetingHost + } + targetP, err := s.participantDAO.GetByRoomAndUser(ctx, room.ID, payload.TargetUserID) + if err != nil { + return err + } + if targetP == nil || !targetP.IsActive() { + return ErrTransferTargetInvalid + } + targetID = payload.TargetUserID + } + + data := map[string]interface{}{ + "room_code": payload.RoomCode, + "user_id": targetID, + "changed_by": fromUserID, + } + if payload.AudioEnabled != nil { + data["audio_enabled"] = *payload.AudioEnabled + } + if payload.VideoEnabled != nil { + data["video_enabled"] = *payload.VideoEnabled + } + + go s.broadcaster.BroadcastToMeeting(context.Background(), room.ID, constants.MeetingWSEventMemberStateChange, data, fromUserID) + return nil +} + +// ========== 媒体组(5 个,mediasoup signaling)========== + +// TransportCreatePayload meeting.transport.create 请求载荷 +type TransportCreatePayload struct { + RoomCode string `json:"room_code"` + Direction string `json:"direction"` // "send" | "recv" +} + +// OnTransportCreate 处理 meeting.transport.create 事件 +func (s *MeetingSignalService) OnTransportCreate(ctx context.Context, userID int64, payload *TransportCreatePayload) (*TransportInfo, error) { + if payload.Direction != "send" && payload.Direction != "recv" { + return nil, fmt.Errorf("direction 非法,必须是 send 或 recv") + } + if _, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID); err != nil { + return nil, err + } + info, err := s.mediaOrchestrator.CreateTransport(ctx, &CreateTransportReq{ + RoomCode: payload.RoomCode, + UserID: userID, + Direction: payload.Direction, + }) + if err != nil { + return nil, err + } + s.trackResource(ctx, payload.RoomCode, userID, "transport", info.ID) + return info, nil +} + +// TransportConnectPayload meeting.transport.connect 请求载荷 +type TransportConnectPayload struct { + RoomCode string `json:"room_code"` + TransportID string `json:"transport_id"` + DtlsParameters json.RawMessage `json:"dtls_parameters"` +} + +// OnTransportConnect 处理 meeting.transport.connect 事件 +func (s *MeetingSignalService) OnTransportConnect(ctx context.Context, userID int64, payload *TransportConnectPayload) error { + if payload.TransportID == "" { + return fmt.Errorf("transport_id 不能为空") + } + if _, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID); err != nil { + return err + } + return s.mediaOrchestrator.ConnectTransport(ctx, payload.TransportID, payload.DtlsParameters) +} + +// ProduceStartPayload meeting.produce.start 请求载荷 +type ProduceStartPayload struct { + RoomCode string `json:"room_code"` + TransportID string `json:"transport_id"` + Kind string `json:"kind"` // "audio" | "video" + RtpParameters json.RawMessage `json:"rtp_parameters"` +} + +// ProduceStartResult 返回给客户端的 producerID +type ProduceStartResult struct { + ProducerID string `json:"producer_id"` +} + +// OnProduceStart 处理 meeting.produce.start 事件 +// 成功后广播 meeting.member.producer.new 给房间内其他成员,驱动对端自动创建 Consumer +func (s *MeetingSignalService) OnProduceStart(ctx context.Context, userID int64, payload *ProduceStartPayload) (*ProduceStartResult, error) { + if payload.Kind != "audio" && payload.Kind != "video" { + return nil, fmt.Errorf("kind 非法,必须是 audio 或 video") + } + if payload.TransportID == "" { + return nil, fmt.Errorf("transport_id 不能为空") + } + room, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID) + if err != nil { + return nil, err + } + producerID, err := s.mediaOrchestrator.CreateProducer(ctx, &CreateProducerReq{ + RoomCode: payload.RoomCode, + UserID: userID, + TransportID: payload.TransportID, + Kind: payload.Kind, + RtpParameters: payload.RtpParameters, + }) + if err != nil { + return nil, err + } + s.trackResource(ctx, payload.RoomCode, userID, "producer", producerID) + + go s.broadcaster.BroadcastToMeeting(context.Background(), room.ID, constants.MeetingWSEventMemberProducerNew, map[string]interface{}{ + "room_code": payload.RoomCode, + "user_id": userID, + "producer_id": producerID, + "kind": payload.Kind, + }, userID) + + return &ProduceStartResult{ProducerID: producerID}, nil +} + +// ConsumeStartPayload meeting.consume.start 请求载荷 +type ConsumeStartPayload struct { + RoomCode string `json:"room_code"` + TransportID string `json:"transport_id"` // 客户端 recv Transport + ProducerID string `json:"producer_id"` // 要订阅的远端 Producer + RtpCapabilities json.RawMessage `json:"rtp_capabilities"` +} + +// OnConsumeStart 处理 meeting.consume.start 事件 +func (s *MeetingSignalService) OnConsumeStart(ctx context.Context, userID int64, payload *ConsumeStartPayload) (*ConsumerInfo, error) { + if payload.TransportID == "" || payload.ProducerID == "" { + return nil, fmt.Errorf("transport_id 与 producer_id 均不能为空") + } + if _, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID); err != nil { + return nil, err + } + info, err := s.mediaOrchestrator.CreateConsumer(ctx, &CreateConsumerReq{ + RoomCode: payload.RoomCode, + UserID: userID, + TransportID: payload.TransportID, + ProducerID: payload.ProducerID, + RtpCapabilities: payload.RtpCapabilities, + }) + if err != nil { + return nil, err + } + s.trackResource(ctx, payload.RoomCode, userID, "consumer", info.ID) + return info, nil +} + +// ProducerClosePayload meeting.producer.close 请求载荷 +type ProducerClosePayload struct { + RoomCode string `json:"room_code"` + ProducerID string `json:"producer_id"` +} + +// OnProducerClose 处理 meeting.producer.close 事件 +// 成功后广播给房间内其他成员(与 mediasoup 的 producerclose 级联动作平级) +func (s *MeetingSignalService) OnProducerClose(ctx context.Context, userID int64, payload *ProducerClosePayload) error { + if payload.ProducerID == "" { + return fmt.Errorf("producer_id 不能为空") + } + room, err := s.loadRoomAndParticipant(ctx, payload.RoomCode, userID) + if err != nil { + return err + } + if err := s.mediaOrchestrator.CloseProducer(ctx, payload.ProducerID); err != nil { + logs.Warn(ctx, "service.meeting_signal_service.OnProducerClose", "关闭 Producer 失败", + zap.String("producer_id", payload.ProducerID), zap.Error(err)) + } + s.untrackResource(ctx, payload.RoomCode, userID, "producer", payload.ProducerID) + + go s.broadcaster.BroadcastToMeeting(context.Background(), room.ID, constants.MeetingWSEventMemberProducerNew, map[string]interface{}{ + "room_code": payload.RoomCode, + "user_id": userID, + "producer_id": payload.ProducerID, + "closed": true, + }, userID) + return nil +} + +// ========== 资源清理 ========== + +// cleanupUserResources 批量关闭指定用户在某会议的所有媒体资源 +// WS 断开、主动离会、被踢时使用;依赖 Redis 集合中追踪的资源 ID +func (s *MeetingSignalService) cleanupUserResources(ctx context.Context, roomCode string, userID int64) { + funcName := "service.meeting_signal_service.cleanupUserResources" + + key := resourceTrackKey(roomCode, userID) + members, err := s.redis.SMembers(ctx, key).Result() + if err != nil { + logs.Warn(ctx, funcName, "读取资源追踪集合失败", zap.String("key", key), zap.Error(err)) + return + } + for _, m := range members { + // 格式:"kind:id" + idx := -1 + for i, c := range m { + if c == ':' { + idx = i + break + } + } + if idx < 0 { + continue + } + kind, id := m[:idx], m[idx+1:] + switch kind { + case "producer": + _ = s.mediaOrchestrator.CloseProducer(ctx, id) + case "consumer": + _ = s.mediaOrchestrator.CloseConsumer(ctx, id) + // transport 关闭一般由 Router 级联;这里不单独处理 + } + } + _ = s.redis.Del(ctx, key).Err() + + if len(members) > 0 { + logs.Info(ctx, funcName, "清理用户媒体资源", + zap.String("room_code", roomCode), + zap.Int64("user_id", userID), + zap.Int("resource_count", len(members))) + } +} diff --git a/backend/go-service/app/provider/provider.go b/backend/go-service/app/provider/provider.go index 2250b2d..5268aac 100644 --- a/backend/go-service/app/provider/provider.go +++ b/backend/go-service/app/provider/provider.go @@ -54,8 +54,10 @@ type App struct { NotifyService *notifyService.NotifyService // 通知业务服务(兼 Pusher、ConnectHook) NotifyController *notifyController.NotificationController // 通知控制器 NotifyCleanupTask *notifyTask.CleanupTask // 通知清理定时任务 - MeetingService *meetingService.MeetingService // 会议业务服务(Task 4 骨架) - MeetingController *meetingController.MeetingController // 会议控制器(Task 4 骨架,handler 返回 501) + MeetingService *meetingService.MeetingService // 会议 REST 业务服务(Task 5 落地) + MeetingSignalService *meetingService.MeetingSignalService // 会议 WS 信令业务服务(Task 6 落地) + MeetingController *meetingController.MeetingController // 会议 REST 控制器 + MeetingWSHandler *meetingController.MeetingWSHandler // 会议 WS 事件 Handler(构造时自动注册路由到 Hub) } // NewApp 创建应用实例 @@ -86,7 +88,9 @@ func NewApp( notifyCtrl *notifyController.NotificationController, notifyCleanup *notifyTask.CleanupTask, meetingSvc *meetingService.MeetingService, + meetingSignalSvc *meetingService.MeetingSignalService, meetingCtrl *meetingController.MeetingController, + meetingWSHandler *meetingController.MeetingWSHandler, ) *App { wsHandler.SetOfflinePusher(offlinePusher) wsHandler.SetNotifyConnectHook(notifySvc) @@ -118,7 +122,9 @@ func NewApp( NotifyController: notifyCtrl, NotifyCleanupTask: notifyCleanup, MeetingService: meetingSvc, + MeetingSignalService: meetingSignalSvc, MeetingController: meetingCtrl, + MeetingWSHandler: meetingWSHandler, } } diff --git a/backend/go-service/app/provider/wire_gen.go b/backend/go-service/app/provider/wire_gen.go index 321d201..fc844f4 100644 --- a/backend/go-service/app/provider/wire_gen.go +++ b/backend/go-service/app/provider/wire_gen.go @@ -101,9 +101,12 @@ func InitializeApp(cfg *config.Config) (*App, error) { meetingRoomDAO := dao6.NewMeetingRoomDAO(gormDB) meetingParticipantDAO := dao6.NewMeetingParticipantDAO(gormDB) meetingChatDAO := dao6.NewMeetingChatDAO(gormDB) + meetingBroadcaster := service7.NewMeetingBroadcaster(meetingParticipantDAO, pubSub) noopMediaOrchestrator := service7.NewNoopMediaOrchestrator() - meetingService := service7.NewMeetingService(meetingRoomDAO, meetingParticipantDAO, meetingChatDAO, gormDB, client, pubSub, notifyService, friendshipDAO, onlineService, noopMediaOrchestrator) + meetingService := service7.NewMeetingService(meetingRoomDAO, meetingParticipantDAO, meetingChatDAO, gormDB, client, meetingBroadcaster, notifyService, friendshipDAO, onlineService, noopMediaOrchestrator) + meetingSignalService := service7.NewMeetingSignalService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, noopMediaOrchestrator) meetingController := controller7.NewMeetingController(meetingService) - 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, meetingController) + meetingWSHandler := controller7.NewMeetingWSHandler(meetingSignalService, hub) + 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, meetingController, meetingWSHandler) return app, nil } diff --git a/docs/api/frontend/meeting.md b/docs/api/frontend/meeting.md index 106be61..f0c3eff 100644 --- a/docs/api/frontend/meeting.md +++ b/docs/api/frontend/meeting.md @@ -3,7 +3,7 @@ > 通用规范(认证、响应包络、通用错误码)见 [README.md](../README.md) > 会议内实时信令(Transport / Producer / Consumer / 控制事件)通过 WebSocket 完成,见 [websocket.md](../websocket.md) -**实施状态**:本文档对应 Phase 2e-2 Task 5 已落地的 12 个 REST 接口,统一前缀 `/api/v1/meeting`,全部需要 JWT 认证。Task 5 完成时间:2026-04-21。 +**实施状态**:本文档对应 Phase 2e-2 Task 5 / Task 6 已落地的 12 个 REST 接口 + 13 个 WebSocket 信令事件,统一前缀 `/api/v1/meeting`(REST)与 `/ws`(WebSocket),全部需要 JWT 认证。Task 5 完成时间:2026-04-21;Task 6 完成时间:2026-04-21。 **设计口径**:以 [`docs/plans/2026-04-21-phase2e-2-design.md`](../../plans/2026-04-21-phase2e-2-design.md) §6.2 为单一事实来源(SSOT)。 @@ -371,29 +371,272 @@ --- -## WebSocket 事件关联 +## WebSocket 信令协议(Task 6) -参考 [websocket.md](../websocket.md) 的 `meeting.*` 事件族。Task 5 内部当前使用 `ws.PubSub.PublishToUser` 对活跃参会者逐个推送(Task 6 将封装为 `BroadcastToMeeting`,接口无感替换): +全部会议相关实时信令走 `/ws?token=` 统一通道,共 **13 个 `meeting.*` 事件**:8 个客户端→服务端(C→S)操作事件 + 5 个服务端→客户端(S→C)广播事件 + 2 个补充业务事件(聊天 + 被踢定向推送,与 REST 广播复用)。 -| 事件 | 触发 REST | 载荷关键字段 | -|------|-----------|-------------| -| `meeting.member.joined` | JoinRoom | room_code / participant | -| `meeting.member.left` | LeaveRoom / KickMember | room_code / user_id / reason | -| `meeting.member.kicked` | KickMember(仅发给被踢者) | room_code | -| `meeting.host.changed` | TransferHost / 隐式(host 离会后自动转让) | room_code / old_host_id / new_host_id | -| `meeting.room.ended` | EndRoom / 空房 TTL | room_code / reason | -| `meeting.chat` | SendChat | message 对象 | +### 帧格式 + +所有消息使用 `ws.Message` 统一信封(JSON): + +| 字段 | 类型 | 说明 | +|------|------|------| +| event | string | 事件名,如 `meeting.room.join` | +| seq | int64 | 客户端自增序号;服务端 ACK 原样回传,用于客户端匹配请求-响应 | +| data | object | 业务载荷 | +| code | int | 仅 ACK / 广播带:0=成功,-1=业务失败 | +| message | string | 仅 ACK:失败时的人类可读原因(与 REST 领域错误口径一致) | +| time | string | ISO8601 时间戳 | + +**ACK 约定**:每一个 C→S 事件服务端都会回发 `.ack`,`seq` 与请求一致;成功时 `code=0`,失败时 `code=-1` 且 `message` 含中文原因(如 `仅主持人可执行此操作`、`会议不存在`、`你当前未在会议中`)。 + +**连接级鉴权**:Token 校验见 [websocket.md](websocket.md) 连接管理;会议内事件额外在每次调用时校验 `(user_id, room_code)` 是否为当前活跃参会者,非法请求全部以 ACK 形式返回 `-1`,不会触发 401/403 HTTP 错误。 + +### 事件总览 + +| # | 方向 | 事件 | 用途 | 服务端动作 | +|---|------|------|------|-----------| +| 1 | C→S | `meeting.room.join` | 前端建立 WS 后声明加入某会议房间频道 | 验证活跃参会,记录 WS 会话 | +| 2 | C→S | `meeting.room.leave` | 前端主动离开 WS 会议频道 | 清理媒体资源(transport / producer / consumer) | +| 3 | C→S | `meeting.member.state.changed` | 成员更新自己的麦克风/摄像头;host 可指定 `target_user_id` 强制静音他人 | 广播 S→C 同名事件,权限不符返回 ACK `code=-1` | +| 4 | C→S | `meeting.transport.create` | 请求创建 mediasoup WebRtcTransport(send/recv) | 调用 `mediaOrchestrator.CreateTransport`,返回 `id/iceParameters/iceCandidates/dtlsParameters` | +| 5 | C→S | `meeting.transport.connect` | 提交 DTLS Parameters 完成 transport 连接 | `mediaOrchestrator.ConnectTransport` | +| 6 | C→S | `meeting.produce.start` | 创建 Producer(上行流) | `mediaOrchestrator.CreateProducer`,成功后广播 `meeting.member.producer.new` | +| 7 | C→S | `meeting.consume.start` | 创建 Consumer(订阅对端 Producer) | `mediaOrchestrator.CreateConsumer` | +| 8 | C→S | `meeting.producer.close` | 关闭自己的 Producer | `mediaOrchestrator.CloseProducer`,广播 `meeting.member.producer.new` (closed=true) | +| 9 | S→C | `meeting.member.joined` | 新成员加入广播(复用 REST /join) | 房间广播 | +| 10 | S→C | `meeting.member.left` | 成员离开广播(REST /leave /kick 或 WS leave) | 房间广播 | +| 11 | S→C | `meeting.member.kicked` | 定向通知被踢者(REST /kick) | `PublishToUser` | +| 12 | S→C | `meeting.host.changed` | 主持人变更 | 房间广播 | +| 13 | S→C | `meeting.room.ended` | 会议被结束(REST /end 或空房 TTL) | 房间广播 | +| 14 | S→C | `meeting.member.state.changed` | 成员状态变化(静音/关摄像头/举手) | 房间广播 | +| 15 | S→C | `meeting.member.producer.new` | 成员开启/关闭媒体流 | 房间广播,`closed=true` 表示关闭 | +| 16 | S→C | `meeting.chat` | 会议内聊天(REST /chats) | 房间广播 | + +> 说明:客户端仅注册 **C→S 白名单**中的 8 个事件(见 `app/constants/meeting.go:MeetingWSClientEvents`),其余 `meeting.*` 事件若由客户端发送均被静默丢弃,防止恶意客户端伪造广播。 + +### 客户端白名单(C→S)详细契约 + +#### 1. `meeting.room.join` + +**请求载荷** + +```json +{ "room_code": "835-000-036" } +``` + +**ACK** `code=0`:`data` 为空对象;失败原因:`会议不存在` / `你当前未在会议中`。 + +**说明**:WS 加入仅用于开启该房间的事件推送通道,REST `/join` 已将 participant 写入 DB,此事件**不会再次修改 participant 表**。 + +#### 2. `meeting.room.leave` + +**请求载荷** + +```json +{ "room_code": "835-000-036" } +``` + +**ACK** `code=0`:`data` 为空对象;服务端同时发起媒体资源清理:按 Redis key `echo:meeting:resources:{room_id}:{user_id}` 中记录的 transport / producer / consumer 依次调用 `mediaOrchestrator.CloseTransport/CloseProducer/CloseConsumer`,随后删除该 key。 + +**说明**:此事件**不会**将 `participant.left_at` 置空,若用户希望真正离会仍需调用 REST `POST /rooms/:code/leave`。此设计允许客户端在 WS 重连时调用 leave + join 刷新资源而不退出会议。 + +#### 3. `meeting.member.state.changed` + +**请求载荷** + +```json +{ + "room_code": "835-000-036", + "target_user_id": 17, + "audio_enabled": false, + "video_enabled": true, + "hand_raised": false +} +``` + +- `target_user_id` 可选:省略或等于自己时为自我操作,任何人均可;指定为他人时需调用方为 host,否则 ACK `code=-1 message=仅主持人可执行此操作`。 +- `audio_enabled / video_enabled / hand_raised` 均为可选布尔字段,至少传一个。 + +**ACK** `code=0`:`data` 为空。**广播** `meeting.member.state.changed`: + +```json +{ + "room_code": "835-000-036", + "user_id": 17, + "audio_enabled": false, + "video_enabled": true, + "hand_raised": false, + "actor_id": 16 +} +``` + +`actor_id` 为触发变更的 user_id(host 强制静音场景下与 `user_id` 不同)。 + +#### 4. `meeting.transport.create` + +**请求载荷** + +```json +{ "room_code": "835-000-036", "direction": "send" } +``` + +`direction`:`"send"` 或 `"recv"`。 + +**ACK** `code=0` `data`: + +```json +{ + "id": "transport-abc", + "iceParameters": { ... }, + "iceCandidates": [ ... ], + "dtlsParameters": { ... } +} +``` + +**资源追踪**:服务端将 `transport-id` 按 `echo:meeting:resources:{room_id}:{user_id}` 的 Redis Set 记录,TTL 1 小时。 + +#### 5. `meeting.transport.connect` + +**请求载荷** + +```json +{ + "room_code": "835-000-036", + "transport_id": "transport-abc", + "dtls_parameters": { ... } +} +``` + +**ACK** `code=0`:`data` 为空。 + +#### 6. `meeting.produce.start` + +**请求载荷** + +```json +{ + "room_code": "835-000-036", + "transport_id": "transport-abc", + "kind": "audio", + "rtp_parameters": { ... } +} +``` + +`kind`: `"audio"` | `"video"`。 + +**ACK** `code=0` `data`: + +```json +{ "producer_id": "producer-xyz" } +``` + +**同时广播** `meeting.member.producer.new`: + +```json +{ + "room_code": "835-000-036", + "user_id": 16, + "producer_id": "producer-xyz", + "kind": "audio", + "closed": false +} +``` + +#### 7. `meeting.consume.start` + +**请求载荷** + +```json +{ + "room_code": "835-000-036", + "transport_id": "transport-recv", + "producer_id": "producer-xyz", + "rtp_capabilities": { ... } +} +``` + +**ACK** `code=0` `data`: + +```json +{ + "id": "consumer-def", + "producerId": "producer-xyz", + "kind": "audio", + "rtpParameters": { ... } +} +``` + +**资源追踪**:consumer-id 也会进入 Redis resources Set。 + +#### 8. `meeting.producer.close` + +**请求载荷** + +```json +{ "room_code": "835-000-036", "producer_id": "producer-xyz" } +``` + +**ACK** `code=0`:`data` 为空。**广播** `meeting.member.producer.new` `closed=true`: + +```json +{ + "room_code": "835-000-036", + "user_id": 16, + "producer_id": "producer-xyz", + "closed": true +} +``` + +### 服务端广播(S→C)详细契约 + +| 事件 | 载荷字段 | 说明 | +|------|---------|------| +| `meeting.member.joined` | room_code / participant 对象 | REST /join 触发 | +| `meeting.member.left` | room_code / user_id / reason | REST /leave /kick 或 WS 资源清理 | +| `meeting.member.kicked` | room_code / reason | 仅发给被踢者本人,对应 REST /kick | +| `meeting.host.changed` | room_code / old_host_id / new_host_id | REST /transfer-host 或 host 离会自动转让 | +| `meeting.room.ended` | room_code / reason(`host_ended` / `empty_ttl`) | REST /end 或空房 TTL | +| `meeting.member.state.changed` | 见白名单 #3 | WS `meeting.member.state.changed` 触发 | +| `meeting.member.producer.new` | 见白名单 #6 / #8 | WS `meeting.produce.start` / `meeting.producer.close` 触发 | +| `meeting.chat` | message 对象 | REST /chats 触发 | + +### 架构 + +- **MeetingWSHandler**(controller 层):thin adapter,仅负责 ws.Hub 事件注册 + JSON 反序列化 + ACK 回写;位于 `app/meeting/controller/meeting_ws_handler.go`。 +- **MeetingSignalService**(service 层):承载 8 个 C→S 事件的业务逻辑(活跃参会校验、host 权限校验、mediaOrchestrator 调用、Redis 资源追踪、广播),位于 `app/meeting/service/meeting_signal_service.go`。 +- **MeetingBroadcaster**(service 层):封装 `BroadcastToMeeting`(查询活跃 participant 列表 → 逐个 `PubSub.PublishToUser`)与 `PublishToUser`,供 REST / WS 两个入口统一使用,位于 `app/meeting/service/meeting_broadcaster.go`。 +- **MediaOrchestrator**(interface):定义 9 个 mediasoup 操作方法;Task 6 使用 `NoopMediaOrchestrator` 占位(返回 stub IDs),Task 7 替换为 `HTTPMediaOrchestrator` 对接 Node media-server。 + +### 错误处理 + +- 鉴权失败(token 无效 / 未传):WS 握手时直接 4001 关闭,不进入业务层。 +- 业务失败(非参会者 / 非 host / 会议已结束):ACK `code=-1` + 中文 `message`,与 REST 领域错误口径完全一致。 +- mediasoup 调用失败:`message` 形如 `mediaOrchestrator error: <原始错误>`,前端可据此降级。 +- WS 连接断开:服务端 `OnDisconnect` 钩子会对该用户所有活跃会议调用 `cleanupUserResources`(遍历 Redis `echo:meeting:resources:*:{user_id}` 完成 transport / producer / consumer 关闭),防止服务端资源泄漏。 --- ## 验证记录 -Task 5 端到端验证脚本(`/tmp/meeting_t5_test.sh`)结果:**19/19 PASS**,覆盖 12 接口的 happy path 与 5 类错误路径(密码错误 / 房间不存在 / 单点参会冲突 / 非 host 越权 / 邀请链接失效)。 +- **Task 5** 端到端验证脚本(`/tmp/meeting_t5_test.sh`)结果:**19/19 PASS**,覆盖 12 接口的 happy path 与 5 类错误路径(密码错误 / 房间不存在 / 单点参会冲突 / 非 host 越权 / 邀请链接失效)。 +- **Task 6** 端到端 WS 测试脚本(`/tmp/meeting_ws_t6_test.mjs`)结果:**18/18 PASS**,覆盖: + - `meeting.room.join` 双端 ACK + - `meeting.member.state.changed` 自我静音(ACK + 对端广播) + - host 强制静音他人(ACK) + - 非 host 强制静音他人 → ACK `code=-1 message=仅主持人可执行此操作` + - `meeting.transport.create` 双方向(send / recv) + - `meeting.transport.connect` + - `meeting.produce.start` + `meeting.member.producer.new` 广播 + - `meeting.consume.start` + - `meeting.producer.close` + `meeting.member.producer.new closed=true` 广播 + - `meeting.room.leave` + `meeting.member.left` 广播 + - 不存在会议号 `meeting.room.join` → ACK `code=-1` --- ## 后续任务关联 -- **Task 6**:WebSocket 信令协议(`meeting.*` 事件、mediasoup transport/producer/consumer 流转) -- **Task 7**:Go → Node HTTP 客户端,将 `NoopMediaOrchestrator` 替换为 `HTTPMediaOrchestrator`,接入真实 mediasoup Router -- **Task 13**:通知卡片 UI 补齐 `meeting_invite` 内联按钮 +- **Task 7**:Go → Node HTTP 客户端,将 `NoopMediaOrchestrator` 替换为 `HTTPMediaOrchestrator`,接入真实 mediasoup Router,此时 WS 白名单事件契约与本文档完全不变,仅 `iceParameters / producer_id` 等字段由 stub 变为真实值。 +- **Task 8**:Vue 前端 mediasoup-client 接入,按本文 WS 契约实现 `mediasoup.Transport` 的 `connect / produce` 回调。 +- **Task 13**:通知卡片 UI 补齐 `meeting_invite` 内联按钮。 diff --git a/docs/plans/2026-04-21-phase2e-2-implementation.plan.md b/docs/plans/2026-04-21-phase2e-2-implementation.plan.md index be11783..0aed9da 100644 --- a/docs/plans/2026-04-21-phase2e-2-implementation.plan.md +++ b/docs/plans/2026-04-21-phase2e-2-implementation.plan.md @@ -5,7 +5,7 @@ > **上级路线图:** [Phase 2e 整体路线图](./2026-04-20-phase2e-design.md) > **分支:** `feature/phase2e-2-meeting-mvp` > **预估总工时:** **约 17 人日**(17 个 Task,含 PoC 与 UI 打磨) -> **最后更新:** 2026-04-21(Task 0-5 ✅ 已落地,下一步 Task 6 WS 信令) +> **最后更新:** 2026-04-21(Task 0-6 ✅ 已落地,下一步 Task 7 HTTPMediaOrchestrator) --- @@ -303,22 +303,29 @@ flowchart LR - **错误码中文化**:使用中文 `message`(与项目惯例一致)而非英文 `meeting_not_found` code,前端通过 HTTP 状态码 + trace_id 区分 - **工作量**:**实际 1 人日**(< 预估 1.5 人日,因 DTO 设计充分 + DAO 契约修复一次到位) -### Task 6:WS 信令 11 事件处理器 +### Task 6:WS 信令 13 事件处理器 ✅ **已完成(2026-04-21)** -- **目标**:实现设计 §6.3 的 11 个 WS 事件,完整对接到 `ws.Hub` -- **依赖**:T4 -- **主要产出**: - - `app/meeting/controller/meeting_ws_handler.go`:注册事件回调到 Hub - - `app/meeting/service/meeting_signal_service.go`:3 组事件(房间 / 成员 / 媒体)的业务逻辑 - - `app/ws/hub.go` 接口扩展:`MeetingSignalDispatcher` 接口注入 + `DispatchMeeting(event, payload)` 方法 - - `app/meeting/constants/ws_events.go`:11 个事件名常量 - - 权限校验:所有事件 handler 入口调用 `assertIsParticipant` / `assertIsHost` - - 广播:`Hub.BroadcastToMeeting(roomCode, event, payload, excludeUserID)` 辅助方法 -- **检查点**: - - 通过 `wscat` 或临时前端脚本连入 WS,逐个事件手测 - - 未授权事件(非参与者发 `meeting.member.state.changed`)被拒绝 - - 事件广播覆盖正确(excludeUserID 生效) -- **工作量**:**1.5 人日** +- **目标**:实现设计 §6.3 的 WS 事件族(最终落地 13 个 = 3 房间 + 5 成员 + 5 媒体),完整对接到 `ws.Hub` +- **依赖**:T4 ✅ +- **实际产出**: + - `app/constants/meeting.go`(改):WS 事件常量与设计 §6.3 对齐 + `MeetingWSClientEvents` 白名单切片(8 个 C→S 事件) + - `app/meeting/service/interfaces.go`(重构):`MediaOrchestrator` 扩容至 9 方法 + 5 DTO + `NoopMediaOrchestrator` 9 占位实现 + - `app/meeting/service/meeting_broadcaster.go`(新建,75 行):统一广播层 `BroadcastToMeeting` + `PublishToUser`,REST / WS 共用 + - `app/meeting/service/meeting_service.go`(重构):12 REST 方法改调 `broadcaster.*`,不再直连 `ws.PubSub` + - `app/meeting/service/meeting_signal_service.go`(新建,430 行):8 C→S 事件业务 + Redis 资源追踪 + 资源清理 + host 权限校验 + - `app/meeting/controller/meeting_ws_handler.go`(新建,200 行):薄层 controller,`hub.RegisterEvent` 注册 + JSON 反序列化 + ACK + - `app/meeting/provider.go` + `app/provider/provider.go`(改):Wire 挂入新 3 个 provider + - `docs/api/frontend/meeting.md`(追加 2 节 +200 行):§WebSocket 信令协议(Task 6) + 验证记录 +- **实际检查点**: + - `/tmp/meeting_ws_t6_test.mjs` 端到端 **18/18 PASS**:8 C→S 白名单事件 + 3 S→C 广播 + 3 类错误路径(非 host 越权/不存在会议号/WS leave 资源清理) + - 非 host 尝试 `meeting.member.state.changed` 修改他人 → ACK `code=-1 message=仅主持人可执行此操作` + - `meeting.room.leave` → 自动触发 `mediaOrchestrator.Close*` 清理 Redis 资源 Set + - `go build` / `go vet` / `wire` 全绿 +- **偏离与说明**: + - 最终事件数从实施计划的 11 个扩展为 13 个(+ 2 个业务补充事件:`meeting.chat.message` + `meeting.member.producer.new`),与设计文档 §6.3 一致 + - 权限校验入口改为在 `MeetingSignalService.On*` 方法内部调用 `assertIsActiveParticipant` / `assertIsHost`,不再依赖 handler 层前置断言,代码更易测试 + - 广播 API 统一为 `MeetingBroadcaster.BroadcastToMeeting`(含 `excludeUserIDs ...int64` 可变参数),比计划中"`Hub.BroadcastToMeeting` 方法" 更内聚,不污染 `ws.Hub` 通用接口 +- **实际工作量**:**1 人日**(比预估 1.5 人日节省,得益于 Task 5 已预置好 DAO / 错误链 / DTO) ### Task 7:Go → Node HTTP Client 封装 diff --git a/docs/progress/CURRENT_STATUS.md b/docs/progress/CURRENT_STATUS.md index f2f2b41..042db3a 100644 --- a/docs/progress/CURRENT_STATUS.md +++ b/docs/progress/CURRENT_STATUS.md @@ -1,7 +1,7 @@ # EchoChat 项目开发进度 -> **最后更新**:2026-04-21(Phase 2e-2 Task 5 Go meeting 模块 12 个 REST 接口业务逻辑全量落地,端到端 19/19 PASS) -> **当前阶段**:Phase 2e-2 会议 MVP **代码开发阶段** 🚧(Task 0-5 ✅ / Task 6-16 待执行) +> **最后更新**:2026-04-21(Phase 2e-2 Task 6 WebSocket 信令协议落地,13 个 meeting.* 事件全量打通,端到端 18/18 PASS) +> **当前阶段**:Phase 2e-2 会议 MVP **代码开发阶段** 🚧(Task 0-6 ✅ / Task 7-16 待执行) > **当前分支**:`feature/phase2e-2-meeting-mvp`(从 `feature/phase2c-group-read-receipt` 衍生) > **Phase 2e 整体设计**:`docs/plans/2026-04-20-phase2e-design.md`(三子阶段路线图 + 后续规划清单) > **Phase 2e-1 专用设计**:`docs/plans/2026-04-20-phase2e-1-design.md`(✅ 已完成) @@ -170,6 +170,63 @@ --- +## 🚀 2026-04-21 Phase 2e-2 Task 6 WebSocket 信令协议(13 事件)落地 + +**交付**:`meeting.*` 事件族从 Task 5 的 `PublishToUser` 循环升级为完整的 WS 信令协议;新建 `MeetingBroadcaster`(统一广播层)、`MeetingSignalService`(8 个 C→S 事件业务逻辑 + 资源追踪)、`MeetingWSHandler`(controller 薄层),`MediaOrchestrator` 接口扩容至 9 个方法覆盖 mediasoup 全生命周期(Task 7 真实实现前由 `NoopMediaOrchestrator` 占位);端到端 WS 冒烟脚本 `/tmp/meeting_ws_t6_test.mjs` **18/18 PASS**,覆盖 8 C→S 白名单事件 + 3 S→C 广播 + 3 类错误路径。 + +### 产出文件 + +| 文件 | 行数 | 作用 | +|---|---|---| +| `backend/go-service/app/constants/meeting.go`(改) | +25 | WS 事件常量与设计 §6.3 对齐(3 房间 + 5 成员 + 5 媒体 + 1 聊天),新增 `MeetingWSClientEvents` 白名单切片限制客户端只能发起 8 个 C→S 事件 | +| `backend/go-service/app/meeting/service/interfaces.go`(重构) | 180 | `MediaOrchestrator` 扩容到 9 方法(Router/Transport/Producer/Consumer 全生命周期)+ 配套 DTO(`TransportInfo` / `ConsumerInfo` / `CreateTransportReq` / `CreateProducerReq` / `CreateConsumerReq`)+ `NoopMediaOrchestrator` 9 个占位实现(stub ID + 最小 JSON) | +| `backend/go-service/app/meeting/service/meeting_broadcaster.go`(新) | 75 | `MeetingBroadcaster`:`BroadcastToMeeting`(查询活跃 participant → 批量 `PubSub.PublishToUser` + 可选 exclude)+ `PublishToUser`(定向推送)+ 并发安全的错误汇集 | +| `backend/go-service/app/meeting/service/meeting_service.go`(改) | ±30 | 12 个 REST 方法重构:统一改为调用 `broadcaster.BroadcastToMeeting` / `broadcaster.PublishToUser`,移除直连 `ws.PubSub` 依赖,代码量精简约 15% | +| `backend/go-service/app/meeting/service/meeting_signal_service.go`(新) | 430 | `MeetingSignalService` 8 个 C→S 事件(`OnRoomJoin`/`OnRoomLeave`/`OnMemberStateChanged`/`OnTransportCreate`/`OnTransportConnect`/`OnProduceStart`/`OnConsumeStart`/`OnProducerClose`)+ Redis 资源追踪 `echo:meeting:resources:{room_id}:{user_id}`(Set 结构,TTL 1 小时)+ `cleanupUserResources`(WS 断开钩子调用)+ host 权限校验(非 host 改他人状态返回 `仅主持人可执行此操作`) | +| `backend/go-service/app/meeting/controller/meeting_ws_handler.go`(新) | 200 | `MeetingWSHandler` 薄层:构造时调用 `hub.RegisterEvent` 注册 8 C→S 事件,每个 handler 仅负责 JSON 反序列化 + 调 `signalSvc.On*` + 构造 ACK(`code=0/-1` + `message`)| +| `backend/go-service/app/meeting/provider.go`(改) | +4 | `MeetingSet` 补全 `NewMeetingBroadcaster` / `NewMeetingSignalService` / `NewMeetingWSHandler` | +| `backend/go-service/app/provider/{provider,wire_gen}.go`(改) | +8 | `App` 结构体新增 `MeetingSignalService` / `MeetingWSHandler` 字段,`wire` 重新生成 | +| `docs/api/frontend/meeting.md`(追加 2 节) | +200 | 新增 §WebSocket 信令协议(Task 6):16 事件总览表 + 8 C→S 事件完整请求/ACK/广播契约 + S→C 广播契约 + 架构说明 + 错误处理表;§验证记录补充 Task 6 结果 | + +### 8 C→S 白名单事件契约 + +| 事件 | 入参关键字段 | ACK data | 副作用 | +|------|-------------|---------|--------| +| `meeting.room.join` | room_code | `{}` | 校验活跃参会记录 | +| `meeting.room.leave` | room_code | `{}` | 清理该用户所有 transport/producer/consumer | +| `meeting.member.state.changed` | room_code, [target_user_id], audio_enabled?, video_enabled?, hand_raised? | `{}` | 广播 `meeting.member.state.changed`;非 host 改他人 → `-1` | +| `meeting.transport.create` | room_code, direction(send/recv) | `{id, iceParameters, iceCandidates, dtlsParameters}` | 资源追踪 Redis set | +| `meeting.transport.connect` | room_code, transport_id, dtls_parameters | `{}` | mediasoup connect | +| `meeting.produce.start` | room_code, transport_id, kind, rtp_parameters | `{producer_id}` | 广播 `meeting.member.producer.new`;资源追踪 | +| `meeting.consume.start` | room_code, transport_id, producer_id, rtp_capabilities | `{id, producerId, kind, rtpParameters}` | 资源追踪 | +| `meeting.producer.close` | room_code, producer_id | `{}` | 广播 `meeting.member.producer.new closed=true` | + +### 验证执行(Node.js + ws) + +1. `go build ./...` / `go vet ./...` / `wire ./app/provider` 全绿 +2. 启动 server,跑 `/tmp/meeting_ws_t6_test.mjs`(双用户 + WS_TRACE 模式) +3. 用例清单(18 个全部 PASS): + - 房间/成员:`room.join` 双端 ACK、`state.changed` 自我静音 ACK + 对端广播、host 强制静音 ACK、**非 host 强制静音他人 → `-1 仅主持人可执行此操作`** + - 媒体:`transport.create` send/recv 双向、`transport.connect`、`produce.start` ACK + `producer.new` 广播、`consume.start`、`producer.close` ACK + `producer.new closed=true` 广播 + - 离会/错误:`room.leave` + `member.left` 对端广播、不存在会议号 `room.join` → `-1` +4. 测试期间修复 waitEvent 死循环 bug(msgQueue pop→push 自循环导致 ack 永远等不到)→ 改为 stash 临时缓冲区,完成后统一归还 + +### 关键设计决策 + +1. **Broadcaster 单独抽层**:避免 service 方法散落直接 `PubSub.PublishToUser`,后续接入 Redis Cluster / 切换广播实现只需改一个文件;同时 Task 5 的 REST 广播与 Task 6 的 WS 事件广播完全复用同一个对象。 +2. **C→S 白名单机制**:`app/constants/meeting.go:MeetingWSClientEvents` 列出 8 个允许客户端发起的事件,`ws.Hub` 在分发前先过滤,防止恶意客户端直接发 `meeting.room.ended` 伪造房间结束。 +3. **资源追踪用 Redis Set**:每创建一个 transport/producer/consumer 都 `SADD echo:meeting:resources:{room_id}:{user_id} `,WS 断开或 room.leave 时 `SMEMBERS` 遍历清理;TTL 1 小时防止遗留占用,即使 Go 进程崩溃也不会泄漏 mediasoup 资源。 +4. **MediaOrchestrator 先抽 9 方法再实现**:Task 6 仍用 `Noop` 占位,但接口已完整定义 `CreateRouter / CloseRouter / CreateTransport / ConnectTransport / CreateProducer / CloseProducer / CreateConsumer / CloseConsumer`(外加 DTO 型号),Task 7 只需替换绑定即可让 WS 端变为真实 mediasoup,**无需修改 signal service / handler 代码**。 +5. **ACK `code` 语义统一**:成功 `0`、业务失败 `-1`(+ 中文 message),与 REST 领域错误口径完全一致,前端可直接复用一套 error toast 组件。 +6. **`meeting.room.leave` 只清 WS 资源不改 participant 表**:真正离会需 REST `/leave`(会影响 duration / host 自动转让);此设计允许客户端 WS 重连时发 leave+join 刷新 transport 而不退会。 + +### 下一步 + +- **Task 7**(1.5 人日):`HTTPMediaOrchestrator` 实现 — Go 调 Node media-server 9 个内部 REST API(Task 2 已全部跑通),替换 `NoopMediaOrchestrator`,前后端 WS 契约完全不变。 +- **Task 8**(2 人日):Vue 前端 mediasoup-client 接入 + 会议室页面骨架,按本次 §WebSocket 信令协议契约实现 `Transport.connect / produce` 回调。 + +--- + ## 🚀 2026-04-21 Phase 2e-2 Task 5 Go meeting 模块 12 个 REST 接口业务逻辑全量落地 **交付**:`MeetingService` 12 个业务方法 + `MeetingController` 12 个 Gin 处理器从 501 占位升级为真实实现,完整的领域错误码映射、DTO 绑定、DAO 契约修复、权限辅助函数;端到端验证脚本 `/tmp/meeting_t5_test.sh` **19/19 PASS**,覆盖 12 接口 happy path + 5 类错误路径;`go build ./...` / `go vet ./...` / `wire ./app/provider` 全绿。