From c35097a0d87db7e878758147576cf7599add0bdd Mon Sep 17 00:00:00 2001 From: bujinyuan Date: Thu, 23 Apr 2026 17:45:03 +0800 Subject: [PATCH] =?UTF-8?q?fix(meeting):=20Task=2016=20=E4=BF=AE=E5=A4=8D?= =?UTF-8?q?=20code-reviewer=20=E5=AE=A1=E8=AE=A1=20P2=20=E4=B8=83=E9=A1=B9?= =?UTF-8?q?=20+=20Nit=20=E4=B8=83=E9=A1=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 覆盖 Phase 2e-2 代码审查报告(docs/reviews/2026-04-23-phase2e-2-code-review.md) P2/Nit 收尾批次,均在本仓库完成闭环;剩余 5 项登记推迟至 Phase 2f/3。 ===== P2 七项 ===== - P2-1 cleanupUserResources 补 transport 清理 · media-server 新增 DELETE /internal/v1/transports/:id + transport.service.closeTransport · Go MediaOrchestrator 接口新增 CloseTransport;HTTP 实现按 doCloseRequest 走 4xx 幂等 + 指数退避 · meeting_signal_service.cleanupUserResources 新增 "transport:" 分支 - P2-2 preview.vue 快速切设备竞态 · previewSeq 序号 + nextTick 后 stale 判断,丢弃过期结果 · onVideoChange/onAudioChange 走 scheduleRestartPreview 200ms 防抖 · onBeforeUnmount 清理 changeDebounceTimer - P2-3 room.vue onLoad redirectTo 后补 return · 引入 redirectingToJoin 守卫,跳转页不再执行 onMounted 初始化 - P2-6 generateUniqueRoomCode 重试上限监控 · 重试后成功:Warn 日志(码空间健康度告警) · 重试耗尽:Error 日志 + ErrRoomCodeConflict · 修正 logs.Error 调用签名(去掉多余的 nil) - P2-7 SendChatMessage 服务端长度 + 频率限制 · 新增 ErrChatContentEmpty / ErrChatContentTooLong / ErrChatRateLimited · utf8.RuneCountInString 校验 500 字符上限 · Redis INCR + EXPIRE 滑动窗口(30 条/60s,首次写入 EXPIRE 兜底) · controller.handleError 映射为 HTTP 400 - P2-8 MEETING_ENDED_REASON_LABEL 覆盖复核 · 新增前端专属常量 MEETING_ENDED_REASON_KICKED + 文案 · store/meeting.js _onMemberKicked 使用常量 · 同步修复后端 OnWSDisconnect 硬编码 "ws_disconnect" → constants.MeetingLeftReasonDisconnect - P2 已修 P2-1/2/3/6/7/8;P2-4(WS token 迁出 URL query)与 P2-5(Chat 服务拆分)登记推迟 ===== Nit 七项 ===== - Nit: kind:id 解析改用 strings.SplitN · cleanupUserResources / pushExistingRoomState 两处同步 - Nit: resourceTTL 中央化 · 新增 constants.MeetingResourceTrackTTLSeconds(3600) · meeting_signal_service.go resourceTTL 由 const 改 var 并引用常量 - Nit: ws/handler.go CheckOrigin 白名单 · NewHandler 新增 serverCfg 依赖;checkOrigin 支持同源放行 / dev 模式放行 / release 模式白名单严格匹配 · config.ServerConfig 新增 WSAllowedOrigins(逗号分隔)+ AllowedOrigins() / IsRelease() 辅助方法 · provider.go 新增 provideServerConfig,wire_gen.go 同步 - Nit: http_media_orchestrator.go 超时 + CreateRouter 重试 · 默认 TimeoutMS 5000→10000ms 兼容 Worker 冷启动 · 新增 CreateRouterRetry(默认 1 次,300ms 退避),仅对非 404 错误重试 · config.dev.yaml / config.docker.yaml 同步写入显式配置 - Nit: deploy-public.sh REDIS_PASSWORD × redis.conf 联动校验 · 检测 REDIS_PASSWORD 与 redis.conf 的 requirepass 配对一致性 · redis.conf 增加公网部署 requirepass 使用说明 - Nit: media-server internal-auth isPrivatePath 按 path 匹配 · 剔除 query/hash 后再与白名单 startsWith,避免 "?" 语义混淆 - Nit: mediasoup-client.js in-flight 锁走读确认 · finally 分支已覆盖 resolve/reject 两路,追加注释强化语义 - Nit 走读复核:_onMemberLeft 整槽关闭 vs _onProducerNew(closed=true) 精确匹配 producerId · 粒度正确,无需改动(登记结论) ===== 构建验证 ===== - go vet ./... / go build ./... 通过 - frontend npm run build:h5 通过(仅 uni-app legacy warning,无 error) - media-server npx tsc --noEmit 通过 ===== 审查追踪小节 ===== docs/reviews/2026-04-23-phase2e-2-code-review.md 追加 "Task 16 修复追踪(2026-04-24 更新)": - 已修复一览(本批次 14 处 + 历次 commit cdaa39d / ea2bf96 / f5ae095 / 5ed14c2) - 推迟登记表(P2-4 / P2-5 / 端口收敛 / appData 校验 / RFC3339 时间格式,共 5 项) Made-with: Cursor --- backend/go-service/app/constants/meeting.go | 10 +++ .../meeting/controller/meeting_controller.go | 8 +- .../service/http_media_orchestrator.go | 73 ++++++++++++++++--- .../app/meeting/service/interfaces.go | 9 +++ .../app/meeting/service/meeting_service.go | 58 ++++++++++++++- .../meeting/service/meeting_signal_service.go | 39 +++++----- backend/go-service/app/provider/provider.go | 7 ++ backend/go-service/app/provider/wire_gen.go | 3 +- backend/go-service/app/ws/handler.go | 69 +++++++++++++++--- backend/go-service/app/ws/provider.go | 5 +- backend/go-service/config/config.dev.yaml | 3 +- backend/go-service/config/config.docker.yaml | 3 +- backend/go-service/config/config.go | 36 +++++++-- deploy/docker/redis/redis.conf | 7 ++ .../2026-04-23-phase2e-2-code-review.md | 34 +++++++++ frontend/src/constants/meeting.js | 7 +- frontend/src/pages/meeting/preview.vue | 36 ++++++++- frontend/src/pages/meeting/room.vue | 9 +++ frontend/src/store/meeting.js | 3 +- frontend/src/utils/mediasoup-client.js | 3 + media-server/src/middlewares/internal-auth.ts | 25 ++++++- media-server/src/routes/transport.route.ts | 14 +++- .../src/services/transport.service.ts | 23 ++++++ scripts/deploy-public.sh | 23 ++++++ 24 files changed, 444 insertions(+), 63 deletions(-) diff --git a/backend/go-service/app/constants/meeting.go b/backend/go-service/app/constants/meeting.go index 71c34f6..24b03d1 100644 --- a/backend/go-service/app/constants/meeting.go +++ b/backend/go-service/app/constants/meeting.go @@ -90,6 +90,16 @@ const ( MeetingEmptyRoomTTLSeconds = 300 // 空房销毁 TTL(5 分钟) MeetingInviteTokenTTL = 600 // 邀请链接 Token TTL(10 分钟) MeetingChatRetentionHours = 24 // 会议聊天保留时长(会议结束后) + + // Task 16 P2-7:服务端聊天限流(防止客户端校验被绕过) + MeetingChatMaxContentLen = 500 // 单条消息最大字符数(Unicode rune 计数) + MeetingChatRateLimitPerMin = 30 // 单用户单会议每分钟最大条数 + MeetingChatRateLimitWindowS = 60 // 限流滑动窗口秒数(配合 Redis INCR + EXPIRE) + + // Task 16 Nit:单个用户资源追踪集合 TTL(秒) + // 原硬编码在 meeting_signal_service.go::resourceTTL = time.Hour; + // 中央化到 constants 便于统一修改 + 单元测试覆盖 + MeetingResourceTrackTTLSeconds = 3600 // 1 小时 ) // 会议 WS 事件常量(设计文档 §6.3) diff --git a/backend/go-service/app/meeting/controller/meeting_controller.go b/backend/go-service/app/meeting/controller/meeting_controller.go index a3e7d2c..43fd820 100644 --- a/backend/go-service/app/meeting/controller/meeting_controller.go +++ b/backend/go-service/app/meeting/controller/meeting_controller.go @@ -64,7 +64,13 @@ func (ctl *MeetingController) handleError(c *gin.Context, err error, fallbackMsg errors.Is(err, service.ErrRoomCodeConflict), errors.Is(err, service.ErrKickSelfForbidden), errors.Is(err, service.ErrTransferToSelf), - errors.Is(err, service.ErrTransferTargetInvalid): + errors.Is(err, service.ErrTransferTargetInvalid), + // P2-7 会议聊天服务端校验失败:内容为空 / 超长 / 触发限流 + // 当前使用 400(由 ResponseBadRequest 返回),保持与其余业务错误一致; + // 未来若需要精细区分(如 ErrChatRateLimited → 429、ErrChatContentTooLong → 413),可拆分分支 + errors.Is(err, service.ErrChatContentEmpty), + errors.Is(err, service.ErrChatContentTooLong), + errors.Is(err, service.ErrChatRateLimited): utils.ResponseBadRequest(c, err.Error()) default: logs.Warn(c.Request.Context(), "controller.meeting_controller.handleError", 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 36b2d6a..18055c1 100644 --- a/backend/go-service/app/meeting/service/http_media_orchestrator.go +++ b/backend/go-service/app/meeting/service/http_media_orchestrator.go @@ -65,7 +65,10 @@ type HTTPMediaOrchestrator struct { func NewHTTPMediaOrchestrator(cfg *config.Config) *HTTPMediaOrchestrator { mc := cfg.MediaServer if mc.TimeoutMS <= 0 { - mc.TimeoutMS = 5000 + // Task 16 Nit(代码审查 2026-04-23 第 15 条): + // 原默认 5000ms 对 CreateRouter 偏紧(Worker 冷启动 + Router 首次创建在慢机上可达 6~8s), + // 统一将默认超时放宽到 10000ms,显式配置(config.*.yaml)不受影响 + mc.TimeoutMS = 10000 } if mc.CloseTimeoutMS <= 0 { mc.CloseTimeoutMS = 2000 @@ -73,6 +76,14 @@ func NewHTTPMediaOrchestrator(cfg *config.Config) *HTTPMediaOrchestrator { if mc.CloseRetry < 0 { mc.CloseRetry = 0 } + if mc.CreateRouterRetry < 0 { + mc.CreateRouterRetry = 0 + } + if mc.CreateRouterRetry == 0 { + // Task 16 Nit(代码审查 2026-04-23 第 15 条): + // 默认允许 1 次轻量重试,300ms 退避,仅对非 404 错误生效 + mc.CreateRouterRetry = 1 + } // 去除 base_url 末尾斜杠,统一拼接风格 mc.BaseURL = strings.TrimRight(mc.BaseURL, "/") @@ -113,15 +124,42 @@ func (h *HTTPMediaOrchestrator) CreateRouter(ctx context.Context, roomCode strin RouterID string `json:"routerId"` RtpCapabilities json.RawMessage `json:"rtpCapabilities"` } - if err := h.doRequest(ctx, requestOptions{ - method: http.MethodPost, - path: "/internal/v1/routers", - body: reqBody, - timeoutMS: h.cfg.TimeoutMS, - funcName: funcName, - logFields: []zap.Field{zap.String("room_code", roomCode)}, - }, &resp); err != nil { - return "", err + + // Task 16 Nit:CreateRouter 在 5xx / 网络错误时允许 CreateRouterRetry 次重试(退避 300ms) + // - 404 不应出现在 POST /routers,若出现视为 media-server 配置异常,不重试 + // - ctx.Err() 立即终止(上游取消或超时) + attempts := h.cfg.CreateRouterRetry + 1 + var lastErr error + for i := 0; i < attempts; i++ { + if i > 0 { + select { + case <-ctx.Done(): + return "", ctx.Err() + case <-time.After(300 * time.Millisecond): + } + logs.Info(ctx, funcName, "CreateRouter 重试", + zap.String("room_code", roomCode), + zap.Int("attempt", i+1)) + } + err := h.doRequest(ctx, requestOptions{ + method: http.MethodPost, + path: "/internal/v1/routers", + body: reqBody, + timeoutMS: h.cfg.TimeoutMS, + funcName: funcName, + logFields: []zap.Field{zap.String("room_code", roomCode), zap.Int("attempt", i+1)}, + }, &resp) + if err == nil { + lastErr = nil + break + } + lastErr = err + if errors.Is(err, ErrMediaResourceNotFound) { + break + } + } + if lastErr != nil { + return "", lastErr } info := &routerInfoCache{ @@ -257,6 +295,21 @@ func (h *HTTPMediaOrchestrator) ConnectTransport(ctx context.Context, transportI }, nil) } +// CloseTransport 调用 DELETE /internal/v1/transports/:id(Task 16 P2-1 引入) +// 404 映射为 ErrMediaResourceNotFound(上层 cleanupUserResources 视为已清理) +// 与 CloseProducer / CloseConsumer 保持同一重试/超时策略 +func (h *HTTPMediaOrchestrator) CloseTransport(ctx context.Context, transportID string) error { + funcName := "service.http_media_orchestrator.CloseTransport" + + err := h.doCloseRequest(ctx, fmt.Sprintf("/internal/v1/transports/%s", transportID), funcName, []zap.Field{ + zap.String("transport_id", transportID), + }) + if errors.Is(err, ErrMediaResourceNotFound) { + return nil + } + return err +} + // CreateProducer 调用 POST /internal/v1/producers func (h *HTTPMediaOrchestrator) CreateProducer(ctx context.Context, req *CreateProducerReq) (string, error) { funcName := "service.http_media_orchestrator.CreateProducer" diff --git a/backend/go-service/app/meeting/service/interfaces.go b/backend/go-service/app/meeting/service/interfaces.go index 8aa3e6c..d20a18a 100644 --- a/backend/go-service/app/meeting/service/interfaces.go +++ b/backend/go-service/app/meeting/service/interfaces.go @@ -103,6 +103,10 @@ type MediaOrchestrator interface { CreateTransport(ctx context.Context, req *CreateTransportReq) (*TransportInfo, error) // ConnectTransport Transport DTLS 握手(幂等:重复 connect 对已连接 transport 视为成功) ConnectTransport(ctx context.Context, transportID string, dtlsParameters json.RawMessage) error + // CloseTransport 主动关闭指定 Transport(Task 16 P2-1 引入) + // 场景:用户离会 / WS 断连时 cleanupUserResources 精确清理,避免等待 Router 级联 + // 语义:404(Transport 已关闭/不存在)返回 ErrMediaResourceNotFound,上层可视为"已清理"幂等成功 + CloseTransport(ctx context.Context, transportID string) error // CreateProducer 在指定 send Transport 上创建 Producer CreateProducer(ctx context.Context, req *CreateProducerReq) (producerID string, err error) @@ -166,6 +170,11 @@ func (n *NoopMediaOrchestrator) ConnectTransport(_ context.Context, _ string, _ return nil } +// CloseTransport 占位:直接返回 nil(Task 16 P2-1 引入) +func (n *NoopMediaOrchestrator) CloseTransport(_ context.Context, _ string) error { + return nil +} + // CreateProducer 占位:返回伪造 ID func (n *NoopMediaOrchestrator) CreateProducer(_ context.Context, req *CreateProducerReq) (string, error) { return "noop-producer-" + req.Kind, nil diff --git a/backend/go-service/app/meeting/service/meeting_service.go b/backend/go-service/app/meeting/service/meeting_service.go index 15b6697..13366de 100644 --- a/backend/go-service/app/meeting/service/meeting_service.go +++ b/backend/go-service/app/meeting/service/meeting_service.go @@ -7,7 +7,9 @@ import ( "errors" "fmt" "strconv" + "strings" "time" + "unicode/utf8" "github.com/echochat/backend/app/constants" "github.com/echochat/backend/app/dto" @@ -45,6 +47,10 @@ var ( // ErrMediaServiceUnavailable 媒体服务当前不可用(Router 创建失败 / Node 宕机等) // 用于 CreateRoom / JoinRoom 的补偿路径,将前台错误与"用户输入错误"区分开 ErrMediaServiceUnavailable = errors.New("媒体服务暂时不可用,请稍后重试") + // Task 16 P2-7:会议聊天服务端限流 + ErrChatContentEmpty = errors.New("消息内容不能为空") + ErrChatContentTooLong = errors.New("消息长度超过上限") + ErrChatRateLimited = errors.New("发送过于频繁,请稍后再试") ) // Redis key 前缀(设计文档 §5.4 - Redis 数据结构) @@ -52,6 +58,8 @@ const ( redisKeyInvitePrefix = "echo:meeting:invite:" // 邀请 Token redisKeyPasswordLockPrefix = "echo:meeting:lock:" // 密码错误锁(code:user_id) redisPasswordAttemptPrefix = "echo:meeting:pwd_attempt:" // 密码错误计数 + // Task 16 P2-7:会议聊天限流计数键,格式 "echo:meeting:chat_rate::" + redisKeyChatRatePrefix = "echo:meeting:chat_rate:" ) // MeetingService 会议业务服务 @@ -132,7 +140,12 @@ func (s *MeetingService) assertIsHost(ctx context.Context, room *model.MeetingRo } // generateUniqueRoomCode 生成唯一的 XXX-XXX-XXX 会议号,冲突最多重试 MeetingRoomCodeRetryMax 次 +// P2-6(代码审查 2026-04-23): +// - 重试上限由 constants.MeetingRoomCodeRetryMax 控制(当前=3),避免死循环 +// - 只要命中第 2 次及以后就打 Warn(正常情况下几乎不会碰撞,连续碰撞通常是码空间 +// / 生成器被外部污染的信号,需要及早告警) func (s *MeetingService) generateUniqueRoomCode(ctx context.Context) (string, error) { + funcName := "service.meeting_service.generateUniqueRoomCode" for i := 0; i < constants.MeetingRoomCodeRetryMax; i++ { code, err := utils.GenerateMeetingRoomCode() if err != nil { @@ -143,9 +156,16 @@ func (s *MeetingService) generateUniqueRoomCode(ctx context.Context) (string, er return "", err } if !exists { + if i > 0 { + logs.Warn(ctx, funcName, "会议号生成多次冲突后成功,请关注码空间健康度", + zap.Int("retry_count", i), + zap.String("code", code)) + } return code, nil } } + logs.Error(ctx, funcName, "会议号生成连续冲突达到上限,疑似码空间异常或并发洪水", + zap.Int("retry_max", constants.MeetingRoomCodeRetryMax)) return "", ErrRoomCodeConflict } @@ -868,9 +888,24 @@ func (s *MeetingService) RedeemInviteToken(ctx context.Context, userID int64, to // ====== 会议内聊天 ====== // SendChatMessage 会议内发送文本消息 +// Task 16 P2-7:加入服务端校验 + Redis 滑窗限流 +// - 内容去除首尾空白后必须非空 +// - Unicode rune 计数不得超过 MeetingChatMaxContentLen +// - 单用户单会议 MeetingChatRateLimitWindowS 秒内最多 MeetingChatRateLimitPerMin 条 func (s *MeetingService) SendChatMessage(ctx context.Context, userID int64, code, content string) (*model.MeetingChat, error) { funcName := "service.meeting_service.SendChatMessage" + trimmed := strings.TrimSpace(content) + if trimmed == "" { + return nil, ErrChatContentEmpty + } + if runeCount := utf8.RuneCountInString(trimmed); runeCount > constants.MeetingChatMaxContentLen { + logs.Debug(ctx, funcName, "聊天消息超长,拒绝", + zap.String("room_code", code), zap.Int64("user_id", userID), + zap.Int("rune_count", runeCount), zap.Int("limit", constants.MeetingChatMaxContentLen)) + return nil, ErrChatContentTooLong + } + room, err := s.roomDAO.GetByCode(ctx, code) if err != nil { return nil, err @@ -885,10 +920,29 @@ func (s *MeetingService) SendChatMessage(ctx context.Context, userID int64, code return nil, err } + // Redis 滑窗限流(INCR + EXPIRE 组合;首条消息初始化窗口) + // 失败时仅 Warn 并放行,避免 Redis 抖动影响用户发消息 + rateKey := fmt.Sprintf("%s%s:%d", redisKeyChatRatePrefix, code, userID) + cnt, rErr := s.redis.Incr(ctx, rateKey).Result() + if rErr != nil { + logs.Warn(ctx, funcName, "聊天限流 INCR 失败(放行)", + zap.String("key", rateKey), zap.Error(rErr)) + } else { + if cnt == 1 { + _ = s.redis.Expire(ctx, rateKey, time.Duration(constants.MeetingChatRateLimitWindowS)*time.Second).Err() + } + if cnt > int64(constants.MeetingChatRateLimitPerMin) { + logs.Warn(ctx, funcName, "聊天限流触发", + zap.String("room_code", code), zap.Int64("user_id", userID), + zap.Int64("count_in_window", cnt), zap.Int("limit", constants.MeetingChatRateLimitPerMin)) + return nil, ErrChatRateLimited + } + } + chat := &model.MeetingChat{ RoomID: room.ID, UserID: userID, - Content: content, + Content: trimmed, } if err := s.chatDAO.Create(ctx, chat); err != nil { return nil, err @@ -903,7 +957,7 @@ func (s *MeetingService) SendChatMessage(ctx context.Context, userID int64, code "user_id": userID, "user_name": userName, "user_avatar": userAvatar, - "content": content, + "content": trimmed, "created_at": chat.CreatedAt.Format("2006-01-02 15:04:05"), }, userID) 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 8adbf01..bac7639 100644 --- a/backend/go-service/app/meeting/service/meeting_signal_service.go +++ b/backend/go-service/app/meeting/service/meeting_signal_service.go @@ -5,6 +5,7 @@ import ( "context" "encoding/json" "fmt" + "strings" "time" "github.com/echochat/backend/app/constants" @@ -70,7 +71,8 @@ func memberStateKey(roomCode string, userID int64) string { // resourceTTL 单个用户资源追踪集合 TTL // 设计:会议期间维持可达即可;若用户长期不活跃由断线清理接管 -const resourceTTL = time.Hour +// Task 16 Nit:常量已迁出至 constants.MeetingResourceTrackTTLSeconds,此处保留计算式 wrapper 方便调用侧零改动 +var resourceTTL = time.Duration(constants.MeetingResourceTrackTTLSeconds) * time.Second // trackResource 记录用户在会议中持有的媒体资源 ID func (s *MeetingSignalService) trackResource(ctx context.Context, roomCode string, userID int64, kind, id string) { @@ -265,17 +267,12 @@ func (s *MeetingSignalService) pushExistingProducers(ctx context.Context, roomID } for _, m := range members { // 格式:"kind:id";仅关心 producer - idx := -1 - for j, c := range m { - if c == ':' { - idx = j - break - } - } - if idx < 0 { + // Nit(代码审查 2026-04-23):改用 strings.SplitN,避免 rune 解码开销 + parts := strings.SplitN(m, ":", 2) + if len(parts) != 2 { continue } - kind, id := m[:idx], m[idx+1:] + kind, id := parts[0], parts[1] if kind != "producer" || id == "" { continue } @@ -410,10 +407,12 @@ func (s *MeetingSignalService) OnRoomLeave(ctx context.Context, userID int64, ro } s.cleanupUserResources(ctx, roomCode, userID) + // P2-8 修复:使用常量 MeetingLeftReasonDisconnect,避免 "ws_disconnect" 等硬编码 + // 与前端 MEETING_LEFT_REASON_LABEL 字面值不一致 go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{ "room_code": roomCode, "user_id": userID, - "reason": "ws_disconnect", + "reason": constants.MeetingLeftReasonDisconnect, }, userID) return nil } @@ -678,23 +677,21 @@ func (s *MeetingSignalService) cleanupUserResources(ctx context.Context, roomCod } for _, m := range members { // 格式:"kind:id" - idx := -1 - for i, c := range m { - if c == ':' { - idx = i - break - } - } - if idx < 0 { + // Nit(代码审查 2026-04-23):原先使用 for-range + rune 匹配 ':',对 ASCII 过度包装; + // 改用 strings.SplitN 限定 2 段更清晰,且避免 rune 解码开销 + parts := strings.SplitN(m, ":", 2) + if len(parts) != 2 { continue } - kind, id := m[:idx], m[idx+1:] + kind, id := parts[0], parts[1] switch kind { case "producer": _ = s.mediaOrchestrator.CloseProducer(ctx, id) case "consumer": _ = s.mediaOrchestrator.CloseConsumer(ctx, id) - // transport 关闭一般由 Router 级联;这里不单独处理 + case "transport": + // Task 16 P2-1:补全 transport 精确清理,短暂抖动重连场景下 Router 不会级联关闭自己的 transport + _ = s.mediaOrchestrator.CloseTransport(ctx, id) } } _ = s.redis.Del(ctx, key).Err() diff --git a/backend/go-service/app/provider/provider.go b/backend/go-service/app/provider/provider.go index 7a1e6ab..15edea3 100644 --- a/backend/go-service/app/provider/provider.go +++ b/backend/go-service/app/provider/provider.go @@ -156,12 +156,19 @@ func provideMinioConfig(cfg *config.Config) *config.MinioConfig { return &cfg.Minio } +// provideServerConfig 从全局 Config 中提取 ServerConfig +// Task 16 Nit:WebSocket CheckOrigin 白名单收敛需要 server.ws_allowed_origins + server.mode +func provideServerConfig(cfg *config.Config) *config.ServerConfig { + return &cfg.Server +} + // InfraSet 基础设施层 Provider Set var InfraSet = wire.NewSet( provideDBConfig, provideRedisConfig, provideJWTConfig, provideMinioConfig, + provideServerConfig, db.NewPostgres, db.NewRedis, storage.NewMinioClient, diff --git a/backend/go-service/app/provider/wire_gen.go b/backend/go-service/app/provider/wire_gen.go index 355a6c7..b76a0a3 100644 --- a/backend/go-service/app/provider/wire_gen.go +++ b/backend/go-service/app/provider/wire_gen.go @@ -80,7 +80,8 @@ func InitializeApp(cfg *config.Config) (*App, error) { conversationDAO := im.ProvideConversationDAO(gormDB) messageManageService := service2.NewMessageManageService(messageManageDAO, userDAO, conversationDAO, pubSub) messageManageController := controller2.NewMessageManageController(messageManageService) - handler := ws.ProvideWSHandler(hub, pubSub, jwtConfig, onlineService, authService) + serverConfig := provideServerConfig(cfg) + handler := ws.ProvideWSHandler(hub, pubSub, jwtConfig, serverConfig, onlineService, authService) friendGroupDAO := dao3.NewFriendGroupDAO(gormDB) notificationDAO := dao5.NewNotificationDAO(gormDB) notifyService := service3.NewNotifyService(notificationDAO, pubSub, friendshipDAO) diff --git a/backend/go-service/app/ws/handler.go b/backend/go-service/app/ws/handler.go index ddde1e7..d50d57b 100644 --- a/backend/go-service/app/ws/handler.go +++ b/backend/go-service/app/ws/handler.go @@ -5,6 +5,8 @@ package ws import ( "context" "net/http" + "net/url" + "strings" "github.com/echochat/backend/config" "github.com/echochat/backend/pkg/logs" @@ -44,35 +46,80 @@ type MeetingDisconnectHook interface { OnWSDisconnect(ctx context.Context, userID int64) } -var upgrader = websocket.Upgrader{ - ReadBufferSize: 1024, - WriteBufferSize: 1024, - CheckOrigin: func(r *http.Request) bool { - return true // TODO: 生产环境通过配置限制 allowed origins - }, -} - // Handler WebSocket 连接处理器 type Handler struct { hub *ws.Hub pubsub *ws.PubSub jwtCfg *config.JWTConfig + serverCfg *config.ServerConfig // Task 16 Nit:CheckOrigin 白名单需要 onlineService *OnlineService tokenValidator TokenValidator offlinePusher OfflineMessagePusher notifyConnectHook NotifyConnectHook meetingDisconnectHook MeetingDisconnectHook // Task 8 注入 + upgrader websocket.Upgrader } // NewHandler 创建 WebSocket Handler 实例 -func NewHandler(hub *ws.Hub, pubsub *ws.PubSub, jwtCfg *config.JWTConfig, onlineService *OnlineService, tokenValidator TokenValidator) *Handler { - return &Handler{ +// Task 16 Nit:新增 serverCfg 参数,按 server.ws_allowed_origins + server.mode 收敛 CheckOrigin +func NewHandler(hub *ws.Hub, pubsub *ws.PubSub, jwtCfg *config.JWTConfig, serverCfg *config.ServerConfig, onlineService *OnlineService, tokenValidator TokenValidator) *Handler { + h := &Handler{ hub: hub, pubsub: pubsub, jwtCfg: jwtCfg, + serverCfg: serverCfg, onlineService: onlineService, tokenValidator: tokenValidator, } + h.upgrader = websocket.Upgrader{ + ReadBufferSize: 1024, + WriteBufferSize: 1024, + CheckOrigin: h.checkOrigin, + } + return h +} + +// checkOrigin 按配置收敛 WebSocket 握手 Origin: +// - 同源(Origin 为空 或 Origin.Host == Request.Host)→ 放行 +// - dev 模式(server.mode != release)+ 未配置白名单 → 放行全部(便于本地开发调试) +// - release 模式 / 配置了白名单 → 仅放行白名单匹配的 Origin +// +// 精确匹配 scheme+host+port,不做前缀/通配,避免误匹配 +func (h *Handler) checkOrigin(r *http.Request) bool { + origin := strings.TrimSpace(r.Header.Get("Origin")) + // 空 Origin:非浏览器场景(如 curl / server-to-server),默认放行 + if origin == "" { + return true + } + u, err := url.Parse(origin) + if err != nil || u.Host == "" { + logs.Warn(r.Context(), "ws.handler.checkOrigin", "Origin 非法,拒绝", + zap.String("origin", origin)) + return false + } + // 同源(Origin.Host == Request.Host)直接放行 + if strings.EqualFold(u.Host, r.Host) { + return true + } + allowed := h.serverCfg.AllowedOrigins() + // 未配置白名单:dev 放行,release 拒绝 + if len(allowed) == 0 { + if h.serverCfg.IsRelease() { + logs.Warn(r.Context(), "ws.handler.checkOrigin", "release 模式未配置 WSAllowedOrigins,拒绝跨源", + zap.String("origin", origin)) + return false + } + return true + } + // 精确匹配白名单 + for _, ao := range allowed { + if strings.EqualFold(ao, origin) { + return true + } + } + logs.Warn(r.Context(), "ws.handler.checkOrigin", "Origin 不在白名单,拒绝", + zap.String("origin", origin)) + return false } // SetOfflinePusher 设置离线消息推送器(由 IM 模块在初始化时注入) @@ -120,7 +167,7 @@ func (h *Handler) Upgrade(c *gin.Context) { return } - conn, err := upgrader.Upgrade(c.Writer, c.Request, nil) + conn, err := h.upgrader.Upgrade(c.Writer, c.Request, nil) if err != nil { logs.Error(nil, funcName, "WebSocket 升级失败", zap.Int64("user_id", claims.UserID), zap.Error(err)) diff --git a/backend/go-service/app/ws/provider.go b/backend/go-service/app/ws/provider.go index 2f631ba..8f57087 100644 --- a/backend/go-service/app/ws/provider.go +++ b/backend/go-service/app/ws/provider.go @@ -26,8 +26,9 @@ func ProvideOnlineService(rdb *redis.Client, hub *ws.Hub, pubsub *ws.PubSub, fri } // ProvideWSHandler 创建 WebSocket Handler -func ProvideWSHandler(hub *ws.Hub, pubsub *ws.PubSub, cfg *config.JWTConfig, onlineService *OnlineService, tokenValidator TokenValidator) *Handler { - return NewHandler(hub, pubsub, cfg, onlineService, tokenValidator) +// Task 16 Nit:新增 ServerConfig 入参,用于 WS 握手 Origin 白名单 +func ProvideWSHandler(hub *ws.Hub, pubsub *ws.PubSub, jwtCfg *config.JWTConfig, serverCfg *config.ServerConfig, onlineService *OnlineService, tokenValidator TokenValidator) *Handler { + return NewHandler(hub, pubsub, jwtCfg, serverCfg, onlineService, tokenValidator) } // WSSet WebSocket 模块 Wire Provider Set diff --git a/backend/go-service/config/config.dev.yaml b/backend/go-service/config/config.dev.yaml index ad70046..9176782 100644 --- a/backend/go-service/config/config.dev.yaml +++ b/backend/go-service/config/config.dev.yaml @@ -60,9 +60,10 @@ minio: media_server: base_url: "http://localhost:3300" # media-server 基础 URL(不含末尾斜杠) internal_token: "dev-token-abcdef1234567890" # 与 media-server/.env MEDIA_INTERNAL_TOKEN 一致(开发环境) - timeout_ms: 5000 # 创建类接口超时(毫秒) + timeout_ms: 10000 # 创建类接口超时(毫秒);Task 16 Nit:5000 → 10000 兼容 Worker 冷启动 close_timeout_ms: 2000 # 关闭类接口超时(毫秒) close_retry: 2 # 关闭类接口失败重试次数(指数退避 200/500ms) + create_router_retry: 1 # Task 16 Nit:CreateRouter 在 5xx/网络错误时的额外重试次数(默认 1 次,退避 300ms) # 会议模块生命周期参数(Phase 2e-2 Task 8 会议状态机) # E2E 测试可通过环境变量 ECHOCHAT_MEETING_HOST_GRACE_SECONDS 等覆盖为更短值以加速触发 diff --git a/backend/go-service/config/config.docker.yaml b/backend/go-service/config/config.docker.yaml index 94ce7ba..a6ca6fc 100644 --- a/backend/go-service/config/config.docker.yaml +++ b/backend/go-service/config/config.docker.yaml @@ -52,9 +52,10 @@ minio: media_server: base_url: "http://media-server:3300" # Docker 网络内服务名 internal_token: "dev-internal-token-change-me" # 生产环境通过 ECHOCHAT_MEDIA_SERVER_INTERNAL_TOKEN 覆盖 - timeout_ms: 5000 + timeout_ms: 10000 # Task 16 Nit:5000 → 10000 兼容 Worker 冷启动 close_timeout_ms: 2000 close_retry: 2 + create_router_retry: 1 # Task 16 Nit:CreateRouter 在 5xx/网络错误时的额外重试次数 # 会议模块生命周期参数(Phase 2e-2 Task 8) meeting: diff --git a/backend/go-service/config/config.go b/backend/go-service/config/config.go index 7f5c6d2..7ffcd96 100644 --- a/backend/go-service/config/config.go +++ b/backend/go-service/config/config.go @@ -34,11 +34,12 @@ type MeetingConfig struct { // 与 media-server/.env 中的 MEDIA_INTERNAL_TOKEN / HTTP_PORT 成对使用 // BaseURL 需精确到协议与端口:http://host:port,不含末尾斜杠 type MediaServerConfig struct { - BaseURL string `mapstructure:"base_url"` // 如 http://localhost:3300 - InternalToken string `mapstructure:"internal_token"` // 与 Node 共享密钥 - TimeoutMS int `mapstructure:"timeout_ms"` // 创建类接口超时(毫秒),默认 5000 - CloseTimeoutMS int `mapstructure:"close_timeout_ms"` // 关闭类接口超时(毫秒),默认 2000 - CloseRetry int `mapstructure:"close_retry"` // 关闭类接口失败重试次数,默认 2 + BaseURL string `mapstructure:"base_url"` // 如 http://localhost:3300 + InternalToken string `mapstructure:"internal_token"` // 与 Node 共享密钥 + TimeoutMS int `mapstructure:"timeout_ms"` // 创建类接口超时(毫秒),默认 10000 + CloseTimeoutMS int `mapstructure:"close_timeout_ms"` // 关闭类接口超时(毫秒),默认 2000 + CloseRetry int `mapstructure:"close_retry"` // 关闭类接口失败重试次数,默认 2 + CreateRouterRetry int `mapstructure:"create_router_retry"` // Task 16 Nit:CreateRouter 5xx/网络错误时的重试次数,默认 1 } // MinioConfig MinIO 对象存储配置 @@ -54,6 +55,31 @@ type MinioConfig struct { type ServerConfig struct { Port int `mapstructure:"port"` // 监听端口 Mode string `mapstructure:"mode"` // 运行模式: debug/release + // Task 16 Nit:WebSocket 升级握手 Origin 白名单(逗号分隔) + // - 空串 → dev 模式(mode != release)放行全部,release 模式强制拒绝所有跨源(仅同源可建连) + // - 配置示例:"https://app.example.com,http://localhost:5173" + // - 环境变量覆盖:ECHOCHAT_SERVER_WS_ALLOWED_ORIGINS="https://a.com,https://b.com" + WSAllowedOrigins string `mapstructure:"ws_allowed_origins"` +} + +// AllowedOrigins 将逗号分隔的 WSAllowedOrigins 解析为 slice,已去空并 trim +func (s *ServerConfig) AllowedOrigins() []string { + if s.WSAllowedOrigins == "" { + return nil + } + raw := strings.Split(s.WSAllowedOrigins, ",") + out := make([]string, 0, len(raw)) + for _, o := range raw { + if v := strings.TrimSpace(o); v != "" { + out = append(out, v) + } + } + return out +} + +// IsRelease 判断是否为生产运行模式 +func (s *ServerConfig) IsRelease() bool { + return strings.EqualFold(s.Mode, "release") } // DatabaseConfig PostgreSQL 数据库配置 diff --git a/deploy/docker/redis/redis.conf b/deploy/docker/redis/redis.conf index fd466df..19bec75 100644 --- a/deploy/docker/redis/redis.conf +++ b/deploy/docker/redis/redis.conf @@ -9,6 +9,13 @@ protected-mode no # 端口 port 6379 +# 访问鉴权(Task 16 Nit:公网部署必启;开发环境默认关闭) +# 公网部署清单: +# 1) 将下一行取消注释并替换为强密码(或改写为 `requirepass ${REDIS_PASSWORD}` 并通过 envsubst 渲染) +# 2) 同步在 deploy/.env 的 REDIS_PASSWORD 填入相同值 +# 3) scripts/deploy-public.sh 会在 Step 1 做 redis.conf 与 .env 的联动校验 +# requirepass _REPLACE_WITH_STRONG_REDIS_PASSWORD_ + # 持久化:每 60 秒内至少 1 个 key 变更时触发 RDB 快照 save 60 1 diff --git a/docs/reviews/2026-04-23-phase2e-2-code-review.md b/docs/reviews/2026-04-23-phase2e-2-code-review.md index 8b376f6..1cde9b5 100644 --- a/docs/reviews/2026-04-23-phase2e-2-code-review.md +++ b/docs/reviews/2026-04-23-phase2e-2-code-review.md @@ -191,3 +191,37 @@ Phase 2e-2 会议 MVP 的代码质量整体达到了"可交付 demo、内网试 | Nit | media-server | `consumer.service.ts` | 7-9 / 41 | producer.appData 资格校验 TODO | | Nit | frontend/media | `mediasoup-client.js` | - | in-flight 锁 reject 分支未清空 | | Nit | frontend/store | `meeting.js` | - | `_onMemberLeft` vs `_cleanupRemoteProducer` 清理粒度 | + +## Task 16 修复追踪(2026-04-24 更新) + +### 已修复(本次 Phase 2e-2 收尾完成) + +| ID | 级别 | 提交 | 结论 | +|---|---|---|---| +| - | P0 × 4 | cdaa39d | Media 资源归属校验 / CreateRoom 补偿 / 会议密码迁 in-memory store,详见 commit message | +| - | P1 × 8 | ea2bf96 | 主持人转让加事务 + SELECT FOR UPDATE / ListMyMeetings 分页 / goroutine trace_id 透传 / broadcast ack 重试 / Redis TTL 分布式锁兜底 / REST+WS 顺序归一等,详见 commit message | +| - | 资源生命周期 | f5ae095 | MeetingLifecycleService.OnRoomEnded 统一撤销 grace/TTL 定时器;Redis `resourceTrackKey` / `memberStateKey` 显式 DEL;前端 Pinia `_reset` + `_pendingBroadcastTimers` 清理 | +| P2-1 | P2 | 本批次 | `cleanupUserResources` 新增 `transport` 分支 + media-server DELETE /transports/:id + Go 端 `CloseTransport` | +| P2-2 | P2 | 本批次 | `preview.vue` 加 `previewSeq` 序号 + 200ms 切换防抖,杜绝快速换摄像头竞态 | +| P2-3 | P2 | 本批次 | `room.vue` `onLoad` redirectTo 后显式 `return`,避免 onMounted 重复初始化 | +| P2-6 | P2 | 本批次 | `generateUniqueRoomCode` 失败日志 + 达到 retry 上限返回 `ErrRoomCodeConflict` | +| P2-7 | P2 | 本批次 | `SendChatMessage` 服务端长度 500 字符(utf8 rune)+ Redis INCR 滑动窗口 30 条/分钟 | +| P2-8 | P2 | 本批次 | `MEETING_ENDED_REASON_LABEL` 补齐 `kicked`;后端 `OnWSDisconnect` 用 `MeetingLeftReasonDisconnect` 常量 | +| Nit splitn | Nit | 本批次 | `meeting_signal_service.go` 的 `kind:id` 解析统一 `strings.SplitN` | +| Nit ttl | Nit | 本批次 | `resourceTTL` 中央化到 `constants.MeetingResourceTrackTTLSeconds` | +| Nit origin | Nit | 本批次 | `ws/handler.go` CheckOrigin 按 `server.ws_allowed_origins` + `server.mode` 收敛,同源放行,release 模式仅允许白名单 | +| Nit timeout | Nit | 本批次 | `http_media_orchestrator.go` 默认 `TimeoutMS` 5000→10000;CreateRouter 新增 `CreateRouterRetry`(默认 1 次,300ms 退避) | +| Nit redispass | Nit | 本批次 | `deploy-public.sh` 追加 REDIS_PASSWORD × redis.conf requirepass 联动校验;redis.conf 加 TODO 注释 | +| Nit routepath | Nit | 本批次 | `internal-auth.ts` 按 path(剔除 query/hash)匹配白名单,避免 `?` 混淆 | +| Nit inflight | Nit | 本批次 | `mediasoup-client.js` 走读确认 `finally` 已覆盖 resolve/reject 两路,加强注释 | +| Nit review | Nit | 本批次 | `_onMemberLeft` 整槽关闭 vs `_onProducerNew(closed=true)` 精确匹配 producerId,粒度正确,无需改动 | + +### 推迟到独立阶段(登记存档) + +| ID | 级别 | 理由 | 跟进计划 | +|---|---|---|---| +| P2-4 | P2 | WS token 从 URL query 迁到首帧鉴权需要同时改 `ws/handler.go` / 前端 `websocket.js` / 反代日志脱敏,改动面大 | Phase 2f 安全专项批次单独处理 | +| P2-5 | P2 | `MeetingChatService` 拆分属于架构重构 | Phase 2f 服务分层专项批次 | +| Nit ports | Nit | `docker-compose.dev.yml` 200 UDP/TCP 端口暴露收敛 | 正式公网部署清单(Phase 3 前) | +| Nit appdata | Nit | `consumer.service.ts` `producer.appData` 资格校验 | Phase 2e-3 流媒体合流时同步做 | +| Nit rfc3339 | Nit | 广播时间格式统一 RFC3339 | 跨模块低优先级,随 Phase 2f 协议梳理 | diff --git a/frontend/src/constants/meeting.js b/frontend/src/constants/meeting.js index 4c84d54..79abe2f 100644 --- a/frontend/src/constants/meeting.js +++ b/frontend/src/constants/meeting.js @@ -45,16 +45,21 @@ export const MEETING_ROLE_LABEL = { // ==================== 会议结束原因 ==================== +// 注意:下列 4 个值必须与后端 backend/go-service/app/constants/meeting.go::MeetingEndedReason* 保持一致 export const MEETING_ENDED_REASON_HOST_ENDED = 'host_ended' export const MEETING_ENDED_REASON_EMPTY_TTL = 'empty_ttl' export const MEETING_ENDED_REASON_ADMIN_FORCE = 'admin_force' export const MEETING_ENDED_REASON_SYSTEM_ERROR = 'system_error' +// 下列为"前端专用"原因(非后端 meeting_rooms.ended_reason 值), +// 用于 store 将当前用户的"单端终止"归因到 ENDED 本地状态时展示 +export const MEETING_ENDED_REASON_KICKED = 'kicked' // 当前用户被主持人移除(由 member.kicked 触发) export const MEETING_ENDED_REASON_LABEL = { [MEETING_ENDED_REASON_HOST_ENDED]: '主持人结束', [MEETING_ENDED_REASON_EMPTY_TTL]: '空房超时', [MEETING_ENDED_REASON_ADMIN_FORCE]: '管理员强制结束', - [MEETING_ENDED_REASON_SYSTEM_ERROR]: '系统异常' + [MEETING_ENDED_REASON_SYSTEM_ERROR]: '系统异常', + [MEETING_ENDED_REASON_KICKED]: '您已被主持人移出会议' } // ==================== 离会原因 ==================== diff --git a/frontend/src/pages/meeting/preview.vue b/frontend/src/pages/meeting/preview.vue index 2bf40b4..f3123d8 100644 --- a/frontend/src/pages/meeting/preview.vue +++ b/frontend/src/pages/meeting/preview.vue @@ -120,6 +120,12 @@ let audioContext = null let analyser = null let volumeRafId = null +// P2-2 修复:快速切换摄像头/麦克风时 getUserMedia 并发竞态保护 +// - previewSeq 每调用一次 startPreview 递增,保证只有"最后一次"请求结果被挂载 +// - changeDebounceTimer 合并 200ms 内的连续切换(防止 picker 快速滚动触发多次拉流) +let previewSeq = 0 +let changeDebounceTimer = null + const hasPermission = ref(false) const permissionError = ref('') const joining = ref(false) @@ -197,6 +203,7 @@ const startPreview = async () => { } return } + const mySeq = ++previewSeq stopPreviewStream() stopAudioMeter() try { @@ -206,13 +213,21 @@ const startPreview = async () => { ? { deviceId: { exact: selectedVideoId.value }, width: { ideal: 1280 }, height: { ideal: 720 }, frameRate: { ideal: 24, max: 30 } } : { width: { ideal: 1280 }, height: { ideal: 720 }, frameRate: { ideal: 24, max: 30 } } } - previewStream = await navigator.mediaDevices.getUserMedia(constraints) + const stream = await navigator.mediaDevices.getUserMedia(constraints) + if (mySeq !== previewSeq) { + // 期间又触发了新的 startPreview,本次结果已过期,立即关闭 track 丢弃 + stream.getTracks().forEach(t => { try { t.stop() } catch {} }) + return + } + previewStream = stream hasPermission.value = true permissionError.value = '' await nextTick() + if (mySeq !== previewSeq) return mountPreviewVideo() startAudioMeter(previewStream) } catch (err) { + if (mySeq !== previewSeq) return hasPermission.value = false permissionError.value = err?.name === 'NotAllowedError' ? '已拒绝摄像头/麦克风权限,请在浏览器地址栏左侧恢复权限后重试' @@ -222,6 +237,17 @@ const startPreview = async () => { // #endif } +// P2-2 修复:设备切换防抖,避免 picker 快速滚动触发多次 getUserMedia 并发 +const scheduleRestartPreview = () => { + if (changeDebounceTimer) { + clearTimeout(changeDebounceTimer) + } + changeDebounceTimer = setTimeout(() => { + changeDebounceTimer = null + startPreview() + }, 200) +} + const mountPreviewVideo = () => { // #ifdef H5 const parent = resolveDom(videoBox.value) @@ -309,14 +335,14 @@ const onVideoChange = (e) => { const d = videoDevices.value[idx] if (!d) return selectedVideoId.value = d.deviceId - startPreview() + scheduleRestartPreview() } const onAudioChange = (e) => { const idx = Number(e.detail.value) const d = audioDevices.value[idx] if (!d) return selectedAudioId.value = d.deviceId - startPreview() + scheduleRestartPreview() } const onSpeakerChange = (e) => { const idx = Number(e.detail.value) @@ -427,6 +453,10 @@ onLoad(async (query) => { }) onBeforeUnmount(() => { + if (changeDebounceTimer) { + clearTimeout(changeDebounceTimer) + changeDebounceTimer = null + } stopPreviewStream() stopAudioMeter() }) diff --git a/frontend/src/pages/meeting/room.vue b/frontend/src/pages/meeting/room.vue index 592f438..e00a81e 100644 --- a/frontend/src/pages/meeting/room.vue +++ b/frontend/src/pages/meeting/room.vue @@ -565,15 +565,24 @@ const redirectHome = () => { let stateStopWatch = null +// P2-3 修复:redirectTo 后必须 return,避免 onMounted 继续执行闪默认 UI +// 使用 module-scoped 标记,onMounted 读取后判断是否提前退出 +let redirectingToJoin = false + onLoad((query) => { // 刷新页面 / 直链访问 room 时 store 为空,回跳 join 避免白屏 if (!meetingStore.isInMeeting) { const paramCode = (query?.code || '').replace(/\D/g, '') + redirectingToJoin = true uni.redirectTo({ url: `/pages/meeting/join${paramCode ? `?code=${paramCode}` : ''}` }) + return } }) onMounted(() => { + // P2-3 修复:onLoad 已 redirectTo 跳转时,避免本页继续初始化定时器与 watcher + if (redirectingToJoin) return + const startIso = meetingStore.currentRoom?.started_at || meetingStore.currentRoom?.created_at startedAt.value = startIso ? Date.parse(startIso) : Date.now() timerHandle = setInterval(() => { nowTs.value = Date.now() }, 1000) diff --git a/frontend/src/store/meeting.js b/frontend/src/store/meeting.js index d9efd05..d227677 100644 --- a/frontend/src/store/meeting.js +++ b/frontend/src/store/meeting.js @@ -48,6 +48,7 @@ import { MEETING_WS_CHAT_MESSAGE, MEETING_ROLE_HOST, MEETING_ENDED_REASON_LABEL, + MEETING_ENDED_REASON_KICKED, MEETING_HOST_AUTO_REASON_LABEL } from '@/constants/meeting' @@ -408,7 +409,7 @@ export const useMeetingStore = defineStore('meeting', () => { const userStore = useUserStore() if (data.user_id === userStore.userInfo?.id) { _log('warn', '[Meeting] 当前用户被踢出会议') - lastEndedReason.value = 'kicked' + lastEndedReason.value = MEETING_ENDED_REASON_KICKED _cleanupMedia() localState.value = MEETING_LOCAL_STATE_ENDED _unregisterListeners() diff --git a/frontend/src/utils/mediasoup-client.js b/frontend/src/utils/mediasoup-client.js index ae7509f..e59f925 100644 --- a/frontend/src/utils/mediasoup-client.js +++ b/frontend/src/utils/mediasoup-client.js @@ -198,6 +198,8 @@ export function createMediaEngine({ roomCode, userId, sendWithAck, logger = defa if (sendTransportPromise) { return sendTransportPromise } + // Nit(代码审查 2026-04-23 第 15 条):in-flight 锁必须在 resolve/reject 两种结局下均置空, + // 否则失败后的下次调用会复用 rejected Promise 直接 throw。这里使用 finally 同时覆盖两路。 sendTransportPromise = (async () => { try { return await _createSendTransport() @@ -232,6 +234,7 @@ export function createMediaEngine({ roomCode, userId, sendWithAck, logger = defa if (recvTransportPromise) { return recvTransportPromise } + // Nit(同 ensureSendTransport):finally 覆盖 resolve/reject 两路置空 recvTransportPromise = (async () => { try { return await _createRecvTransport() diff --git a/media-server/src/middlewares/internal-auth.ts b/media-server/src/middlewares/internal-auth.ts index caa971e..173f169 100644 --- a/media-server/src/middlewares/internal-auth.ts +++ b/media-server/src/middlewares/internal-auth.ts @@ -14,8 +14,29 @@ const INTERNAL_TOKEN_HEADER = 'x-internal-token'; */ const PRIVATE_PATH_PREFIXES = ['/internal/']; -function isPrivatePath(url: string): boolean { - return PRIVATE_PATH_PREFIXES.some((prefix) => url === prefix.slice(0, -1) || url.startsWith(prefix)); +/** + * 提取 URL 的 pathname,去除 query string / hash。 + * + * Task 16 Nit(代码审查 2026-04-23):之前直接用 `request.url` 做 startsWith 匹配, + * 对 `/healthz?x=/internal/...` 这种带 query 的请求是稳定的(因为 query 在路径之后), + * 但对构造如 `/internal/foo` 开头但含 `?` 的路径日志会出现误判; + * 同时更严格的语义要求按 route path 决策而非 raw URL。 + * + * onRequest 阶段(fastify 路由匹配之前)`request.routerPath` 尚未赋值, + * 因此这里手动按 `?` / `#` 截断再与白名单比较,保证只依赖 path 本身。 + */ +function extractPath(rawUrl: string): string { + const qIdx = rawUrl.indexOf('?'); + const hIdx = rawUrl.indexOf('#'); + let end = rawUrl.length; + if (qIdx !== -1) end = Math.min(end, qIdx); + if (hIdx !== -1) end = Math.min(end, hIdx); + return rawUrl.slice(0, end); +} + +function isPrivatePath(rawUrl: string): boolean { + const path = extractPath(rawUrl); + return PRIVATE_PATH_PREFIXES.some((prefix) => path === prefix.slice(0, -1) || path.startsWith(prefix)); } function safeEqual(a: string, b: string): boolean { diff --git a/media-server/src/routes/transport.route.ts b/media-server/src/routes/transport.route.ts index c6e2d49..0109623 100644 --- a/media-server/src/routes/transport.route.ts +++ b/media-server/src/routes/transport.route.ts @@ -5,7 +5,11 @@ import { createTransportBodySchema, transportIdParamSchema, } from '../schemas/transport.schema.js'; -import { connectTransport, createWebRtcTransport } from '../services/transport.service.js'; +import { + closeTransport, + connectTransport, + createWebRtcTransport, +} from '../services/transport.service.js'; export async function transportRoutes(app: FastifyInstance): Promise { app.post('/transports', async (request, reply) => { @@ -25,4 +29,12 @@ export async function transportRoutes(app: FastifyInstance): Promise { await connectTransport({ transportId: id, dtlsParameters }); return { ok: true as const }; }); + + // Task 16 P2-1:补全 Transport 主动关闭接口,供 Go go-service 在用户离会/断连时 + // 精确清理 orphan transport(避免等待 Router 级联) + app.delete('/transports/:id', async (request) => { + const { id } = transportIdParamSchema.parse(request.params); + closeTransport(id); + return { ok: true as const }; + }); } diff --git a/media-server/src/services/transport.service.ts b/media-server/src/services/transport.service.ts index 9e26690..5a95415 100644 --- a/media-server/src/services/transport.service.ts +++ b/media-server/src/services/transport.service.ts @@ -143,6 +143,29 @@ export function getTransportStats(): { total: number } { return { total: transportMap.size }; } +/** + * 主动关闭指定 Transport(Task 16 P2-1:cleanupUserResources 补全) + * - 已不存在 → 抛 notFound(HTTP 层映射 404,Go orchestrator 视为已清理的幂等成功) + * - 已存在 → 调 transport.close(),observer 'close' 会自动从 transportMap 移除 + * - 幂等:外层可放心重试;内部依赖 mediasoup 的 close() 本身幂等 + */ +export function closeTransport(transportId: string): void { + const entry = transportMap.get(transportId); + if (!entry) { + throw notFound('transport', transportId); + } + try { + entry.transport.close(); + } catch (err) { + // mediasoup 对重复 close 一般不抛,但仍吞错避免"部分成功"语义 + log.warn( + { transportId, err: err instanceof Error ? err.message : String(err) }, + 'transport.close threw (ignored, already closed?)', + ); + } + log.info({ transportId, userId: entry.userId, routerId: entry.routerId }, 'transport closed via API'); +} + /** 测试专用:复位 transportMap,生产环境调用会抛错 */ export function _clearTransportMap(): void { assertTestOnly('_clearTransportMap'); diff --git a/scripts/deploy-public.sh b/scripts/deploy-public.sh index 0dd48bf..7717340 100755 --- a/scripts/deploy-public.sh +++ b/scripts/deploy-public.sh @@ -99,6 +99,29 @@ step_validate_env() { warn "JWT_SECRET 长度仅 ${#JWT_SECRET} 字符,推荐 >= 64 字符(openssl rand -hex 32)" fi + # Task 16 Nit(代码审查 2026-04-23):REDIS_PASSWORD 与 redis.conf requirepass 联动校验 + # 背景:REDIS_PASSWORD 只影响 go-service / media-server 的连接侧, + # 若 deploy/docker/redis/redis.conf 未设置 requirepass,Redis 本身会接受无密码连接(相当于裸奔) + local redis_conf="$DEPLOY_DIR/docker/redis/redis.conf" + if [[ -f "$redis_conf" ]]; then + local has_requirepass + has_requirepass=$(grep -E '^\s*requirepass\s+\S+' "$redis_conf" || true) + if [[ -n "${REDIS_PASSWORD:-}" ]]; then + if [[ -z "$has_requirepass" ]]; then + err "REDIS_PASSWORD 已设置,但 $redis_conf 未启用 requirepass —— Redis 仍允许空密码登录,公网暴露时务必同步开启" + log_info "建议:在 redis.conf 追加 requirepass \$REDIS_PASSWORD(或在 docker-compose command 中传入 --requirepass)" + else + log_ok "Redis requirepass 已在 redis.conf 启用,与 REDIS_PASSWORD 联动" + fi + else + if [[ -n "$has_requirepass" ]]; then + warn "redis.conf 启用了 requirepass,但 deploy/.env 的 REDIS_PASSWORD 为空 —— go-service/media-server 连接会失败" + fi + fi + else + warn "未找到 $redis_conf,跳过 Redis 密码联动校验" + fi + # TURN 开关一致性 if [[ "${TURN_ENABLED:-false}" != "true" ]]; then warn "TURN_ENABLED=${TURN_ENABLED:-false},公网部署下建议开启 TURN 作为对称 NAT fallback"