fix(meeting): Task 16 修复 code-reviewer 审计 P2 七项 + Nit 七项
覆盖 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:<id>" 分支 - 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 处 + 历次 commitcdaa39d/ea2bf96/f5ae095/ 5ed14c2) - 推迟登记表(P2-4 / P2-5 / 端口收敛 / appData 校验 / RFC3339 时间格式,共 5 项) Made-with: Cursor
This commit is contained in:
@@ -90,6 +90,16 @@ const (
|
|||||||
MeetingEmptyRoomTTLSeconds = 300 // 空房销毁 TTL(5 分钟)
|
MeetingEmptyRoomTTLSeconds = 300 // 空房销毁 TTL(5 分钟)
|
||||||
MeetingInviteTokenTTL = 600 // 邀请链接 Token TTL(10 分钟)
|
MeetingInviteTokenTTL = 600 // 邀请链接 Token TTL(10 分钟)
|
||||||
MeetingChatRetentionHours = 24 // 会议聊天保留时长(会议结束后)
|
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)
|
// 会议 WS 事件常量(设计文档 §6.3)
|
||||||
|
|||||||
@@ -64,7 +64,13 @@ func (ctl *MeetingController) handleError(c *gin.Context, err error, fallbackMsg
|
|||||||
errors.Is(err, service.ErrRoomCodeConflict),
|
errors.Is(err, service.ErrRoomCodeConflict),
|
||||||
errors.Is(err, service.ErrKickSelfForbidden),
|
errors.Is(err, service.ErrKickSelfForbidden),
|
||||||
errors.Is(err, service.ErrTransferToSelf),
|
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())
|
utils.ResponseBadRequest(c, err.Error())
|
||||||
default:
|
default:
|
||||||
logs.Warn(c.Request.Context(), "controller.meeting_controller.handleError",
|
logs.Warn(c.Request.Context(), "controller.meeting_controller.handleError",
|
||||||
|
|||||||
@@ -65,7 +65,10 @@ type HTTPMediaOrchestrator struct {
|
|||||||
func NewHTTPMediaOrchestrator(cfg *config.Config) *HTTPMediaOrchestrator {
|
func NewHTTPMediaOrchestrator(cfg *config.Config) *HTTPMediaOrchestrator {
|
||||||
mc := cfg.MediaServer
|
mc := cfg.MediaServer
|
||||||
if mc.TimeoutMS <= 0 {
|
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 {
|
if mc.CloseTimeoutMS <= 0 {
|
||||||
mc.CloseTimeoutMS = 2000
|
mc.CloseTimeoutMS = 2000
|
||||||
@@ -73,6 +76,14 @@ func NewHTTPMediaOrchestrator(cfg *config.Config) *HTTPMediaOrchestrator {
|
|||||||
if mc.CloseRetry < 0 {
|
if mc.CloseRetry < 0 {
|
||||||
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 末尾斜杠,统一拼接风格
|
// 去除 base_url 末尾斜杠,统一拼接风格
|
||||||
mc.BaseURL = strings.TrimRight(mc.BaseURL, "/")
|
mc.BaseURL = strings.TrimRight(mc.BaseURL, "/")
|
||||||
|
|
||||||
@@ -113,15 +124,42 @@ func (h *HTTPMediaOrchestrator) CreateRouter(ctx context.Context, roomCode strin
|
|||||||
RouterID string `json:"routerId"`
|
RouterID string `json:"routerId"`
|
||||||
RtpCapabilities json.RawMessage `json:"rtpCapabilities"`
|
RtpCapabilities json.RawMessage `json:"rtpCapabilities"`
|
||||||
}
|
}
|
||||||
if err := h.doRequest(ctx, requestOptions{
|
|
||||||
method: http.MethodPost,
|
// Task 16 Nit:CreateRouter 在 5xx / 网络错误时允许 CreateRouterRetry 次重试(退避 300ms)
|
||||||
path: "/internal/v1/routers",
|
// - 404 不应出现在 POST /routers,若出现视为 media-server 配置异常,不重试
|
||||||
body: reqBody,
|
// - ctx.Err() 立即终止(上游取消或超时)
|
||||||
timeoutMS: h.cfg.TimeoutMS,
|
attempts := h.cfg.CreateRouterRetry + 1
|
||||||
funcName: funcName,
|
var lastErr error
|
||||||
logFields: []zap.Field{zap.String("room_code", roomCode)},
|
for i := 0; i < attempts; i++ {
|
||||||
}, &resp); err != nil {
|
if i > 0 {
|
||||||
return "", err
|
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{
|
info := &routerInfoCache{
|
||||||
@@ -257,6 +295,21 @@ func (h *HTTPMediaOrchestrator) ConnectTransport(ctx context.Context, transportI
|
|||||||
}, nil)
|
}, 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
|
// CreateProducer 调用 POST /internal/v1/producers
|
||||||
func (h *HTTPMediaOrchestrator) CreateProducer(ctx context.Context, req *CreateProducerReq) (string, error) {
|
func (h *HTTPMediaOrchestrator) CreateProducer(ctx context.Context, req *CreateProducerReq) (string, error) {
|
||||||
funcName := "service.http_media_orchestrator.CreateProducer"
|
funcName := "service.http_media_orchestrator.CreateProducer"
|
||||||
|
|||||||
@@ -103,6 +103,10 @@ type MediaOrchestrator interface {
|
|||||||
CreateTransport(ctx context.Context, req *CreateTransportReq) (*TransportInfo, error)
|
CreateTransport(ctx context.Context, req *CreateTransportReq) (*TransportInfo, error)
|
||||||
// ConnectTransport Transport DTLS 握手(幂等:重复 connect 对已连接 transport 视为成功)
|
// ConnectTransport Transport DTLS 握手(幂等:重复 connect 对已连接 transport 视为成功)
|
||||||
ConnectTransport(ctx context.Context, transportID string, dtlsParameters json.RawMessage) error
|
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 在指定 send Transport 上创建 Producer
|
||||||
CreateProducer(ctx context.Context, req *CreateProducerReq) (producerID string, err error)
|
CreateProducer(ctx context.Context, req *CreateProducerReq) (producerID string, err error)
|
||||||
@@ -166,6 +170,11 @@ func (n *NoopMediaOrchestrator) ConnectTransport(_ context.Context, _ string, _
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// CloseTransport 占位:直接返回 nil(Task 16 P2-1 引入)
|
||||||
|
func (n *NoopMediaOrchestrator) CloseTransport(_ context.Context, _ string) error {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
// CreateProducer 占位:返回伪造 ID
|
// CreateProducer 占位:返回伪造 ID
|
||||||
func (n *NoopMediaOrchestrator) CreateProducer(_ context.Context, req *CreateProducerReq) (string, error) {
|
func (n *NoopMediaOrchestrator) CreateProducer(_ context.Context, req *CreateProducerReq) (string, error) {
|
||||||
return "noop-producer-" + req.Kind, nil
|
return "noop-producer-" + req.Kind, nil
|
||||||
|
|||||||
@@ -7,7 +7,9 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"strconv"
|
"strconv"
|
||||||
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
"unicode/utf8"
|
||||||
|
|
||||||
"github.com/echochat/backend/app/constants"
|
"github.com/echochat/backend/app/constants"
|
||||||
"github.com/echochat/backend/app/dto"
|
"github.com/echochat/backend/app/dto"
|
||||||
@@ -45,6 +47,10 @@ var (
|
|||||||
// ErrMediaServiceUnavailable 媒体服务当前不可用(Router 创建失败 / Node 宕机等)
|
// ErrMediaServiceUnavailable 媒体服务当前不可用(Router 创建失败 / Node 宕机等)
|
||||||
// 用于 CreateRoom / JoinRoom 的补偿路径,将前台错误与"用户输入错误"区分开
|
// 用于 CreateRoom / JoinRoom 的补偿路径,将前台错误与"用户输入错误"区分开
|
||||||
ErrMediaServiceUnavailable = errors.New("媒体服务暂时不可用,请稍后重试")
|
ErrMediaServiceUnavailable = errors.New("媒体服务暂时不可用,请稍后重试")
|
||||||
|
// Task 16 P2-7:会议聊天服务端限流
|
||||||
|
ErrChatContentEmpty = errors.New("消息内容不能为空")
|
||||||
|
ErrChatContentTooLong = errors.New("消息长度超过上限")
|
||||||
|
ErrChatRateLimited = errors.New("发送过于频繁,请稍后再试")
|
||||||
)
|
)
|
||||||
|
|
||||||
// Redis key 前缀(设计文档 §5.4 - Redis 数据结构)
|
// Redis key 前缀(设计文档 §5.4 - Redis 数据结构)
|
||||||
@@ -52,6 +58,8 @@ const (
|
|||||||
redisKeyInvitePrefix = "echo:meeting:invite:" // 邀请 Token
|
redisKeyInvitePrefix = "echo:meeting:invite:" // 邀请 Token
|
||||||
redisKeyPasswordLockPrefix = "echo:meeting:lock:" // 密码错误锁(code:user_id)
|
redisKeyPasswordLockPrefix = "echo:meeting:lock:" // 密码错误锁(code:user_id)
|
||||||
redisPasswordAttemptPrefix = "echo:meeting:pwd_attempt:" // 密码错误计数
|
redisPasswordAttemptPrefix = "echo:meeting:pwd_attempt:" // 密码错误计数
|
||||||
|
// Task 16 P2-7:会议聊天限流计数键,格式 "echo:meeting:chat_rate:<room>:<user_id>"
|
||||||
|
redisKeyChatRatePrefix = "echo:meeting:chat_rate:"
|
||||||
)
|
)
|
||||||
|
|
||||||
// MeetingService 会议业务服务
|
// MeetingService 会议业务服务
|
||||||
@@ -132,7 +140,12 @@ func (s *MeetingService) assertIsHost(ctx context.Context, room *model.MeetingRo
|
|||||||
}
|
}
|
||||||
|
|
||||||
// generateUniqueRoomCode 生成唯一的 XXX-XXX-XXX 会议号,冲突最多重试 MeetingRoomCodeRetryMax 次
|
// generateUniqueRoomCode 生成唯一的 XXX-XXX-XXX 会议号,冲突最多重试 MeetingRoomCodeRetryMax 次
|
||||||
|
// P2-6(代码审查 2026-04-23):
|
||||||
|
// - 重试上限由 constants.MeetingRoomCodeRetryMax 控制(当前=3),避免死循环
|
||||||
|
// - 只要命中第 2 次及以后就打 Warn(正常情况下几乎不会碰撞,连续碰撞通常是码空间
|
||||||
|
// / 生成器被外部污染的信号,需要及早告警)
|
||||||
func (s *MeetingService) generateUniqueRoomCode(ctx context.Context) (string, error) {
|
func (s *MeetingService) generateUniqueRoomCode(ctx context.Context) (string, error) {
|
||||||
|
funcName := "service.meeting_service.generateUniqueRoomCode"
|
||||||
for i := 0; i < constants.MeetingRoomCodeRetryMax; i++ {
|
for i := 0; i < constants.MeetingRoomCodeRetryMax; i++ {
|
||||||
code, err := utils.GenerateMeetingRoomCode()
|
code, err := utils.GenerateMeetingRoomCode()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -143,9 +156,16 @@ func (s *MeetingService) generateUniqueRoomCode(ctx context.Context) (string, er
|
|||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
if !exists {
|
if !exists {
|
||||||
|
if i > 0 {
|
||||||
|
logs.Warn(ctx, funcName, "会议号生成多次冲突后成功,请关注码空间健康度",
|
||||||
|
zap.Int("retry_count", i),
|
||||||
|
zap.String("code", code))
|
||||||
|
}
|
||||||
return code, nil
|
return code, nil
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
logs.Error(ctx, funcName, "会议号生成连续冲突达到上限,疑似码空间异常或并发洪水",
|
||||||
|
zap.Int("retry_max", constants.MeetingRoomCodeRetryMax))
|
||||||
return "", ErrRoomCodeConflict
|
return "", ErrRoomCodeConflict
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -868,9 +888,24 @@ func (s *MeetingService) RedeemInviteToken(ctx context.Context, userID int64, to
|
|||||||
// ====== 会议内聊天 ======
|
// ====== 会议内聊天 ======
|
||||||
|
|
||||||
// SendChatMessage 会议内发送文本消息
|
// 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) {
|
func (s *MeetingService) SendChatMessage(ctx context.Context, userID int64, code, content string) (*model.MeetingChat, error) {
|
||||||
funcName := "service.meeting_service.SendChatMessage"
|
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)
|
room, err := s.roomDAO.GetByCode(ctx, code)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -885,10 +920,29 @@ func (s *MeetingService) SendChatMessage(ctx context.Context, userID int64, code
|
|||||||
return nil, err
|
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{
|
chat := &model.MeetingChat{
|
||||||
RoomID: room.ID,
|
RoomID: room.ID,
|
||||||
UserID: userID,
|
UserID: userID,
|
||||||
Content: content,
|
Content: trimmed,
|
||||||
}
|
}
|
||||||
if err := s.chatDAO.Create(ctx, chat); err != nil {
|
if err := s.chatDAO.Create(ctx, chat); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -903,7 +957,7 @@ func (s *MeetingService) SendChatMessage(ctx context.Context, userID int64, code
|
|||||||
"user_id": userID,
|
"user_id": userID,
|
||||||
"user_name": userName,
|
"user_name": userName,
|
||||||
"user_avatar": userAvatar,
|
"user_avatar": userAvatar,
|
||||||
"content": content,
|
"content": trimmed,
|
||||||
"created_at": chat.CreatedAt.Format("2006-01-02 15:04:05"),
|
"created_at": chat.CreatedAt.Format("2006-01-02 15:04:05"),
|
||||||
}, userID)
|
}, userID)
|
||||||
|
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/echochat/backend/app/constants"
|
"github.com/echochat/backend/app/constants"
|
||||||
@@ -70,7 +71,8 @@ func memberStateKey(roomCode string, userID int64) string {
|
|||||||
|
|
||||||
// resourceTTL 单个用户资源追踪集合 TTL
|
// resourceTTL 单个用户资源追踪集合 TTL
|
||||||
// 设计:会议期间维持可达即可;若用户长期不活跃由断线清理接管
|
// 设计:会议期间维持可达即可;若用户长期不活跃由断线清理接管
|
||||||
const resourceTTL = time.Hour
|
// Task 16 Nit:常量已迁出至 constants.MeetingResourceTrackTTLSeconds,此处保留计算式 wrapper 方便调用侧零改动
|
||||||
|
var resourceTTL = time.Duration(constants.MeetingResourceTrackTTLSeconds) * time.Second
|
||||||
|
|
||||||
// trackResource 记录用户在会议中持有的媒体资源 ID
|
// trackResource 记录用户在会议中持有的媒体资源 ID
|
||||||
func (s *MeetingSignalService) trackResource(ctx context.Context, roomCode string, userID int64, kind, id string) {
|
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 {
|
for _, m := range members {
|
||||||
// 格式:"kind:id";仅关心 producer
|
// 格式:"kind:id";仅关心 producer
|
||||||
idx := -1
|
// Nit(代码审查 2026-04-23):改用 strings.SplitN,避免 rune 解码开销
|
||||||
for j, c := range m {
|
parts := strings.SplitN(m, ":", 2)
|
||||||
if c == ':' {
|
if len(parts) != 2 {
|
||||||
idx = j
|
|
||||||
break
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if idx < 0 {
|
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
kind, id := m[:idx], m[idx+1:]
|
kind, id := parts[0], parts[1]
|
||||||
if kind != "producer" || id == "" {
|
if kind != "producer" || id == "" {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
@@ -410,10 +407,12 @@ func (s *MeetingSignalService) OnRoomLeave(ctx context.Context, userID int64, ro
|
|||||||
}
|
}
|
||||||
s.cleanupUserResources(ctx, roomCode, userID)
|
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{}{
|
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{
|
||||||
"room_code": roomCode,
|
"room_code": roomCode,
|
||||||
"user_id": userID,
|
"user_id": userID,
|
||||||
"reason": "ws_disconnect",
|
"reason": constants.MeetingLeftReasonDisconnect,
|
||||||
}, userID)
|
}, userID)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
@@ -678,23 +677,21 @@ func (s *MeetingSignalService) cleanupUserResources(ctx context.Context, roomCod
|
|||||||
}
|
}
|
||||||
for _, m := range members {
|
for _, m := range members {
|
||||||
// 格式:"kind:id"
|
// 格式:"kind:id"
|
||||||
idx := -1
|
// Nit(代码审查 2026-04-23):原先使用 for-range + rune 匹配 ':',对 ASCII 过度包装;
|
||||||
for i, c := range m {
|
// 改用 strings.SplitN 限定 2 段更清晰,且避免 rune 解码开销
|
||||||
if c == ':' {
|
parts := strings.SplitN(m, ":", 2)
|
||||||
idx = i
|
if len(parts) != 2 {
|
||||||
break
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if idx < 0 {
|
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
kind, id := m[:idx], m[idx+1:]
|
kind, id := parts[0], parts[1]
|
||||||
switch kind {
|
switch kind {
|
||||||
case "producer":
|
case "producer":
|
||||||
_ = s.mediaOrchestrator.CloseProducer(ctx, id)
|
_ = s.mediaOrchestrator.CloseProducer(ctx, id)
|
||||||
case "consumer":
|
case "consumer":
|
||||||
_ = s.mediaOrchestrator.CloseConsumer(ctx, id)
|
_ = 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()
|
_ = s.redis.Del(ctx, key).Err()
|
||||||
|
|||||||
@@ -156,12 +156,19 @@ func provideMinioConfig(cfg *config.Config) *config.MinioConfig {
|
|||||||
return &cfg.Minio
|
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
|
// InfraSet 基础设施层 Provider Set
|
||||||
var InfraSet = wire.NewSet(
|
var InfraSet = wire.NewSet(
|
||||||
provideDBConfig,
|
provideDBConfig,
|
||||||
provideRedisConfig,
|
provideRedisConfig,
|
||||||
provideJWTConfig,
|
provideJWTConfig,
|
||||||
provideMinioConfig,
|
provideMinioConfig,
|
||||||
|
provideServerConfig,
|
||||||
db.NewPostgres,
|
db.NewPostgres,
|
||||||
db.NewRedis,
|
db.NewRedis,
|
||||||
storage.NewMinioClient,
|
storage.NewMinioClient,
|
||||||
|
|||||||
@@ -80,7 +80,8 @@ func InitializeApp(cfg *config.Config) (*App, error) {
|
|||||||
conversationDAO := im.ProvideConversationDAO(gormDB)
|
conversationDAO := im.ProvideConversationDAO(gormDB)
|
||||||
messageManageService := service2.NewMessageManageService(messageManageDAO, userDAO, conversationDAO, pubSub)
|
messageManageService := service2.NewMessageManageService(messageManageDAO, userDAO, conversationDAO, pubSub)
|
||||||
messageManageController := controller2.NewMessageManageController(messageManageService)
|
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)
|
friendGroupDAO := dao3.NewFriendGroupDAO(gormDB)
|
||||||
notificationDAO := dao5.NewNotificationDAO(gormDB)
|
notificationDAO := dao5.NewNotificationDAO(gormDB)
|
||||||
notifyService := service3.NewNotifyService(notificationDAO, pubSub, friendshipDAO)
|
notifyService := service3.NewNotifyService(notificationDAO, pubSub, friendshipDAO)
|
||||||
|
|||||||
@@ -5,6 +5,8 @@ package ws
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"net/url"
|
||||||
|
"strings"
|
||||||
|
|
||||||
"github.com/echochat/backend/config"
|
"github.com/echochat/backend/config"
|
||||||
"github.com/echochat/backend/pkg/logs"
|
"github.com/echochat/backend/pkg/logs"
|
||||||
@@ -44,35 +46,80 @@ type MeetingDisconnectHook interface {
|
|||||||
OnWSDisconnect(ctx context.Context, userID int64)
|
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 连接处理器
|
// Handler WebSocket 连接处理器
|
||||||
type Handler struct {
|
type Handler struct {
|
||||||
hub *ws.Hub
|
hub *ws.Hub
|
||||||
pubsub *ws.PubSub
|
pubsub *ws.PubSub
|
||||||
jwtCfg *config.JWTConfig
|
jwtCfg *config.JWTConfig
|
||||||
|
serverCfg *config.ServerConfig // Task 16 Nit:CheckOrigin 白名单需要
|
||||||
onlineService *OnlineService
|
onlineService *OnlineService
|
||||||
tokenValidator TokenValidator
|
tokenValidator TokenValidator
|
||||||
offlinePusher OfflineMessagePusher
|
offlinePusher OfflineMessagePusher
|
||||||
notifyConnectHook NotifyConnectHook
|
notifyConnectHook NotifyConnectHook
|
||||||
meetingDisconnectHook MeetingDisconnectHook // Task 8 注入
|
meetingDisconnectHook MeetingDisconnectHook // Task 8 注入
|
||||||
|
upgrader websocket.Upgrader
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewHandler 创建 WebSocket Handler 实例
|
// NewHandler 创建 WebSocket Handler 实例
|
||||||
func NewHandler(hub *ws.Hub, pubsub *ws.PubSub, jwtCfg *config.JWTConfig, onlineService *OnlineService, tokenValidator TokenValidator) *Handler {
|
// Task 16 Nit:新增 serverCfg 参数,按 server.ws_allowed_origins + server.mode 收敛 CheckOrigin
|
||||||
return &Handler{
|
func NewHandler(hub *ws.Hub, pubsub *ws.PubSub, jwtCfg *config.JWTConfig, serverCfg *config.ServerConfig, onlineService *OnlineService, tokenValidator TokenValidator) *Handler {
|
||||||
|
h := &Handler{
|
||||||
hub: hub,
|
hub: hub,
|
||||||
pubsub: pubsub,
|
pubsub: pubsub,
|
||||||
jwtCfg: jwtCfg,
|
jwtCfg: jwtCfg,
|
||||||
|
serverCfg: serverCfg,
|
||||||
onlineService: onlineService,
|
onlineService: onlineService,
|
||||||
tokenValidator: tokenValidator,
|
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 模块在初始化时注入)
|
// SetOfflinePusher 设置离线消息推送器(由 IM 模块在初始化时注入)
|
||||||
@@ -120,7 +167,7 @@ func (h *Handler) Upgrade(c *gin.Context) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
conn, err := upgrader.Upgrade(c.Writer, c.Request, nil)
|
conn, err := h.upgrader.Upgrade(c.Writer, c.Request, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logs.Error(nil, funcName, "WebSocket 升级失败",
|
logs.Error(nil, funcName, "WebSocket 升级失败",
|
||||||
zap.Int64("user_id", claims.UserID), zap.Error(err))
|
zap.Int64("user_id", claims.UserID), zap.Error(err))
|
||||||
|
|||||||
@@ -26,8 +26,9 @@ func ProvideOnlineService(rdb *redis.Client, hub *ws.Hub, pubsub *ws.PubSub, fri
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ProvideWSHandler 创建 WebSocket Handler
|
// ProvideWSHandler 创建 WebSocket Handler
|
||||||
func ProvideWSHandler(hub *ws.Hub, pubsub *ws.PubSub, cfg *config.JWTConfig, onlineService *OnlineService, tokenValidator TokenValidator) *Handler {
|
// Task 16 Nit:新增 ServerConfig 入参,用于 WS 握手 Origin 白名单
|
||||||
return NewHandler(hub, pubsub, cfg, onlineService, tokenValidator)
|
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
|
// WSSet WebSocket 模块 Wire Provider Set
|
||||||
|
|||||||
@@ -60,9 +60,10 @@ minio:
|
|||||||
media_server:
|
media_server:
|
||||||
base_url: "http://localhost:3300" # media-server 基础 URL(不含末尾斜杠)
|
base_url: "http://localhost:3300" # media-server 基础 URL(不含末尾斜杠)
|
||||||
internal_token: "dev-token-abcdef1234567890" # 与 media-server/.env MEDIA_INTERNAL_TOKEN 一致(开发环境)
|
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_timeout_ms: 2000 # 关闭类接口超时(毫秒)
|
||||||
close_retry: 2 # 关闭类接口失败重试次数(指数退避 200/500ms)
|
close_retry: 2 # 关闭类接口失败重试次数(指数退避 200/500ms)
|
||||||
|
create_router_retry: 1 # Task 16 Nit:CreateRouter 在 5xx/网络错误时的额外重试次数(默认 1 次,退避 300ms)
|
||||||
|
|
||||||
# 会议模块生命周期参数(Phase 2e-2 Task 8 会议状态机)
|
# 会议模块生命周期参数(Phase 2e-2 Task 8 会议状态机)
|
||||||
# E2E 测试可通过环境变量 ECHOCHAT_MEETING_HOST_GRACE_SECONDS 等覆盖为更短值以加速触发
|
# E2E 测试可通过环境变量 ECHOCHAT_MEETING_HOST_GRACE_SECONDS 等覆盖为更短值以加速触发
|
||||||
|
|||||||
@@ -52,9 +52,10 @@ minio:
|
|||||||
media_server:
|
media_server:
|
||||||
base_url: "http://media-server:3300" # Docker 网络内服务名
|
base_url: "http://media-server:3300" # Docker 网络内服务名
|
||||||
internal_token: "dev-internal-token-change-me" # 生产环境通过 ECHOCHAT_MEDIA_SERVER_INTERNAL_TOKEN 覆盖
|
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_timeout_ms: 2000
|
||||||
close_retry: 2
|
close_retry: 2
|
||||||
|
create_router_retry: 1 # Task 16 Nit:CreateRouter 在 5xx/网络错误时的额外重试次数
|
||||||
|
|
||||||
# 会议模块生命周期参数(Phase 2e-2 Task 8)
|
# 会议模块生命周期参数(Phase 2e-2 Task 8)
|
||||||
meeting:
|
meeting:
|
||||||
|
|||||||
@@ -34,11 +34,12 @@ type MeetingConfig struct {
|
|||||||
// 与 media-server/.env 中的 MEDIA_INTERNAL_TOKEN / HTTP_PORT 成对使用
|
// 与 media-server/.env 中的 MEDIA_INTERNAL_TOKEN / HTTP_PORT 成对使用
|
||||||
// BaseURL 需精确到协议与端口:http://host:port,不含末尾斜杠
|
// BaseURL 需精确到协议与端口:http://host:port,不含末尾斜杠
|
||||||
type MediaServerConfig struct {
|
type MediaServerConfig struct {
|
||||||
BaseURL string `mapstructure:"base_url"` // 如 http://localhost:3300
|
BaseURL string `mapstructure:"base_url"` // 如 http://localhost:3300
|
||||||
InternalToken string `mapstructure:"internal_token"` // 与 Node 共享密钥
|
InternalToken string `mapstructure:"internal_token"` // 与 Node 共享密钥
|
||||||
TimeoutMS int `mapstructure:"timeout_ms"` // 创建类接口超时(毫秒),默认 5000
|
TimeoutMS int `mapstructure:"timeout_ms"` // 创建类接口超时(毫秒),默认 10000
|
||||||
CloseTimeoutMS int `mapstructure:"close_timeout_ms"` // 关闭类接口超时(毫秒),默认 2000
|
CloseTimeoutMS int `mapstructure:"close_timeout_ms"` // 关闭类接口超时(毫秒),默认 2000
|
||||||
CloseRetry int `mapstructure:"close_retry"` // 关闭类接口失败重试次数,默认 2
|
CloseRetry int `mapstructure:"close_retry"` // 关闭类接口失败重试次数,默认 2
|
||||||
|
CreateRouterRetry int `mapstructure:"create_router_retry"` // Task 16 Nit:CreateRouter 5xx/网络错误时的重试次数,默认 1
|
||||||
}
|
}
|
||||||
|
|
||||||
// MinioConfig MinIO 对象存储配置
|
// MinioConfig MinIO 对象存储配置
|
||||||
@@ -54,6 +55,31 @@ type MinioConfig struct {
|
|||||||
type ServerConfig struct {
|
type ServerConfig struct {
|
||||||
Port int `mapstructure:"port"` // 监听端口
|
Port int `mapstructure:"port"` // 监听端口
|
||||||
Mode string `mapstructure:"mode"` // 运行模式: debug/release
|
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 数据库配置
|
// DatabaseConfig PostgreSQL 数据库配置
|
||||||
|
|||||||
@@ -9,6 +9,13 @@ protected-mode no
|
|||||||
# 端口
|
# 端口
|
||||||
port 6379
|
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 快照
|
# 持久化:每 60 秒内至少 1 个 key 变更时触发 RDB 快照
|
||||||
save 60 1
|
save 60 1
|
||||||
|
|
||||||
|
|||||||
@@ -191,3 +191,37 @@ Phase 2e-2 会议 MVP 的代码质量整体达到了"可交付 demo、内网试
|
|||||||
| Nit | media-server | `consumer.service.ts` | 7-9 / 41 | producer.appData 资格校验 TODO |
|
| Nit | media-server | `consumer.service.ts` | 7-9 / 41 | producer.appData 资格校验 TODO |
|
||||||
| Nit | frontend/media | `mediasoup-client.js` | - | in-flight 锁 reject 分支未清空 |
|
| Nit | frontend/media | `mediasoup-client.js` | - | in-flight 锁 reject 分支未清空 |
|
||||||
| Nit | frontend/store | `meeting.js` | - | `_onMemberLeft` vs `_cleanupRemoteProducer` 清理粒度 |
|
| 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 协议梳理 |
|
||||||
|
|||||||
@@ -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_HOST_ENDED = 'host_ended'
|
||||||
export const MEETING_ENDED_REASON_EMPTY_TTL = 'empty_ttl'
|
export const MEETING_ENDED_REASON_EMPTY_TTL = 'empty_ttl'
|
||||||
export const MEETING_ENDED_REASON_ADMIN_FORCE = 'admin_force'
|
export const MEETING_ENDED_REASON_ADMIN_FORCE = 'admin_force'
|
||||||
export const MEETING_ENDED_REASON_SYSTEM_ERROR = 'system_error'
|
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 = {
|
export const MEETING_ENDED_REASON_LABEL = {
|
||||||
[MEETING_ENDED_REASON_HOST_ENDED]: '主持人结束',
|
[MEETING_ENDED_REASON_HOST_ENDED]: '主持人结束',
|
||||||
[MEETING_ENDED_REASON_EMPTY_TTL]: '空房超时',
|
[MEETING_ENDED_REASON_EMPTY_TTL]: '空房超时',
|
||||||
[MEETING_ENDED_REASON_ADMIN_FORCE]: '管理员强制结束',
|
[MEETING_ENDED_REASON_ADMIN_FORCE]: '管理员强制结束',
|
||||||
[MEETING_ENDED_REASON_SYSTEM_ERROR]: '系统异常'
|
[MEETING_ENDED_REASON_SYSTEM_ERROR]: '系统异常',
|
||||||
|
[MEETING_ENDED_REASON_KICKED]: '您已被主持人移出会议'
|
||||||
}
|
}
|
||||||
|
|
||||||
// ==================== 离会原因 ====================
|
// ==================== 离会原因 ====================
|
||||||
|
|||||||
@@ -120,6 +120,12 @@ let audioContext = null
|
|||||||
let analyser = null
|
let analyser = null
|
||||||
let volumeRafId = null
|
let volumeRafId = null
|
||||||
|
|
||||||
|
// P2-2 修复:快速切换摄像头/麦克风时 getUserMedia 并发竞态保护
|
||||||
|
// - previewSeq 每调用一次 startPreview 递增,保证只有"最后一次"请求结果被挂载
|
||||||
|
// - changeDebounceTimer 合并 200ms 内的连续切换(防止 picker 快速滚动触发多次拉流)
|
||||||
|
let previewSeq = 0
|
||||||
|
let changeDebounceTimer = null
|
||||||
|
|
||||||
const hasPermission = ref(false)
|
const hasPermission = ref(false)
|
||||||
const permissionError = ref('')
|
const permissionError = ref('')
|
||||||
const joining = ref(false)
|
const joining = ref(false)
|
||||||
@@ -197,6 +203,7 @@ const startPreview = async () => {
|
|||||||
}
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
const mySeq = ++previewSeq
|
||||||
stopPreviewStream()
|
stopPreviewStream()
|
||||||
stopAudioMeter()
|
stopAudioMeter()
|
||||||
try {
|
try {
|
||||||
@@ -206,13 +213,21 @@ const startPreview = async () => {
|
|||||||
? { deviceId: { exact: selectedVideoId.value }, width: { ideal: 1280 }, height: { ideal: 720 }, frameRate: { ideal: 24, max: 30 } }
|
? { 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 } }
|
: { 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
|
hasPermission.value = true
|
||||||
permissionError.value = ''
|
permissionError.value = ''
|
||||||
await nextTick()
|
await nextTick()
|
||||||
|
if (mySeq !== previewSeq) return
|
||||||
mountPreviewVideo()
|
mountPreviewVideo()
|
||||||
startAudioMeter(previewStream)
|
startAudioMeter(previewStream)
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
|
if (mySeq !== previewSeq) return
|
||||||
hasPermission.value = false
|
hasPermission.value = false
|
||||||
permissionError.value = err?.name === 'NotAllowedError'
|
permissionError.value = err?.name === 'NotAllowedError'
|
||||||
? '已拒绝摄像头/麦克风权限,请在浏览器地址栏左侧恢复权限后重试'
|
? '已拒绝摄像头/麦克风权限,请在浏览器地址栏左侧恢复权限后重试'
|
||||||
@@ -222,6 +237,17 @@ const startPreview = async () => {
|
|||||||
// #endif
|
// #endif
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// P2-2 修复:设备切换防抖,避免 picker 快速滚动触发多次 getUserMedia 并发
|
||||||
|
const scheduleRestartPreview = () => {
|
||||||
|
if (changeDebounceTimer) {
|
||||||
|
clearTimeout(changeDebounceTimer)
|
||||||
|
}
|
||||||
|
changeDebounceTimer = setTimeout(() => {
|
||||||
|
changeDebounceTimer = null
|
||||||
|
startPreview()
|
||||||
|
}, 200)
|
||||||
|
}
|
||||||
|
|
||||||
const mountPreviewVideo = () => {
|
const mountPreviewVideo = () => {
|
||||||
// #ifdef H5
|
// #ifdef H5
|
||||||
const parent = resolveDom(videoBox.value)
|
const parent = resolveDom(videoBox.value)
|
||||||
@@ -309,14 +335,14 @@ const onVideoChange = (e) => {
|
|||||||
const d = videoDevices.value[idx]
|
const d = videoDevices.value[idx]
|
||||||
if (!d) return
|
if (!d) return
|
||||||
selectedVideoId.value = d.deviceId
|
selectedVideoId.value = d.deviceId
|
||||||
startPreview()
|
scheduleRestartPreview()
|
||||||
}
|
}
|
||||||
const onAudioChange = (e) => {
|
const onAudioChange = (e) => {
|
||||||
const idx = Number(e.detail.value)
|
const idx = Number(e.detail.value)
|
||||||
const d = audioDevices.value[idx]
|
const d = audioDevices.value[idx]
|
||||||
if (!d) return
|
if (!d) return
|
||||||
selectedAudioId.value = d.deviceId
|
selectedAudioId.value = d.deviceId
|
||||||
startPreview()
|
scheduleRestartPreview()
|
||||||
}
|
}
|
||||||
const onSpeakerChange = (e) => {
|
const onSpeakerChange = (e) => {
|
||||||
const idx = Number(e.detail.value)
|
const idx = Number(e.detail.value)
|
||||||
@@ -427,6 +453,10 @@ onLoad(async (query) => {
|
|||||||
})
|
})
|
||||||
|
|
||||||
onBeforeUnmount(() => {
|
onBeforeUnmount(() => {
|
||||||
|
if (changeDebounceTimer) {
|
||||||
|
clearTimeout(changeDebounceTimer)
|
||||||
|
changeDebounceTimer = null
|
||||||
|
}
|
||||||
stopPreviewStream()
|
stopPreviewStream()
|
||||||
stopAudioMeter()
|
stopAudioMeter()
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -565,15 +565,24 @@ const redirectHome = () => {
|
|||||||
|
|
||||||
let stateStopWatch = null
|
let stateStopWatch = null
|
||||||
|
|
||||||
|
// P2-3 修复:redirectTo 后必须 return,避免 onMounted 继续执行闪默认 UI
|
||||||
|
// 使用 module-scoped 标记,onMounted 读取后判断是否提前退出
|
||||||
|
let redirectingToJoin = false
|
||||||
|
|
||||||
onLoad((query) => {
|
onLoad((query) => {
|
||||||
// 刷新页面 / 直链访问 room 时 store 为空,回跳 join 避免白屏
|
// 刷新页面 / 直链访问 room 时 store 为空,回跳 join 避免白屏
|
||||||
if (!meetingStore.isInMeeting) {
|
if (!meetingStore.isInMeeting) {
|
||||||
const paramCode = (query?.code || '').replace(/\D/g, '')
|
const paramCode = (query?.code || '').replace(/\D/g, '')
|
||||||
|
redirectingToJoin = true
|
||||||
uni.redirectTo({ url: `/pages/meeting/join${paramCode ? `?code=${paramCode}` : ''}` })
|
uni.redirectTo({ url: `/pages/meeting/join${paramCode ? `?code=${paramCode}` : ''}` })
|
||||||
|
return
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
|
||||||
onMounted(() => {
|
onMounted(() => {
|
||||||
|
// P2-3 修复:onLoad 已 redirectTo 跳转时,避免本页继续初始化定时器与 watcher
|
||||||
|
if (redirectingToJoin) return
|
||||||
|
|
||||||
const startIso = meetingStore.currentRoom?.started_at || meetingStore.currentRoom?.created_at
|
const startIso = meetingStore.currentRoom?.started_at || meetingStore.currentRoom?.created_at
|
||||||
startedAt.value = startIso ? Date.parse(startIso) : Date.now()
|
startedAt.value = startIso ? Date.parse(startIso) : Date.now()
|
||||||
timerHandle = setInterval(() => { nowTs.value = Date.now() }, 1000)
|
timerHandle = setInterval(() => { nowTs.value = Date.now() }, 1000)
|
||||||
|
|||||||
@@ -48,6 +48,7 @@ import {
|
|||||||
MEETING_WS_CHAT_MESSAGE,
|
MEETING_WS_CHAT_MESSAGE,
|
||||||
MEETING_ROLE_HOST,
|
MEETING_ROLE_HOST,
|
||||||
MEETING_ENDED_REASON_LABEL,
|
MEETING_ENDED_REASON_LABEL,
|
||||||
|
MEETING_ENDED_REASON_KICKED,
|
||||||
MEETING_HOST_AUTO_REASON_LABEL
|
MEETING_HOST_AUTO_REASON_LABEL
|
||||||
} from '@/constants/meeting'
|
} from '@/constants/meeting'
|
||||||
|
|
||||||
@@ -408,7 +409,7 @@ export const useMeetingStore = defineStore('meeting', () => {
|
|||||||
const userStore = useUserStore()
|
const userStore = useUserStore()
|
||||||
if (data.user_id === userStore.userInfo?.id) {
|
if (data.user_id === userStore.userInfo?.id) {
|
||||||
_log('warn', '[Meeting] 当前用户被踢出会议')
|
_log('warn', '[Meeting] 当前用户被踢出会议')
|
||||||
lastEndedReason.value = 'kicked'
|
lastEndedReason.value = MEETING_ENDED_REASON_KICKED
|
||||||
_cleanupMedia()
|
_cleanupMedia()
|
||||||
localState.value = MEETING_LOCAL_STATE_ENDED
|
localState.value = MEETING_LOCAL_STATE_ENDED
|
||||||
_unregisterListeners()
|
_unregisterListeners()
|
||||||
|
|||||||
@@ -198,6 +198,8 @@ export function createMediaEngine({ roomCode, userId, sendWithAck, logger = defa
|
|||||||
if (sendTransportPromise) {
|
if (sendTransportPromise) {
|
||||||
return sendTransportPromise
|
return sendTransportPromise
|
||||||
}
|
}
|
||||||
|
// Nit(代码审查 2026-04-23 第 15 条):in-flight 锁必须在 resolve/reject 两种结局下均置空,
|
||||||
|
// 否则失败后的下次调用会复用 rejected Promise 直接 throw。这里使用 finally 同时覆盖两路。
|
||||||
sendTransportPromise = (async () => {
|
sendTransportPromise = (async () => {
|
||||||
try {
|
try {
|
||||||
return await _createSendTransport()
|
return await _createSendTransport()
|
||||||
@@ -232,6 +234,7 @@ export function createMediaEngine({ roomCode, userId, sendWithAck, logger = defa
|
|||||||
if (recvTransportPromise) {
|
if (recvTransportPromise) {
|
||||||
return recvTransportPromise
|
return recvTransportPromise
|
||||||
}
|
}
|
||||||
|
// Nit(同 ensureSendTransport):finally 覆盖 resolve/reject 两路置空
|
||||||
recvTransportPromise = (async () => {
|
recvTransportPromise = (async () => {
|
||||||
try {
|
try {
|
||||||
return await _createRecvTransport()
|
return await _createRecvTransport()
|
||||||
|
|||||||
@@ -14,8 +14,29 @@ const INTERNAL_TOKEN_HEADER = 'x-internal-token';
|
|||||||
*/
|
*/
|
||||||
const PRIVATE_PATH_PREFIXES = ['/internal/'];
|
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 {
|
function safeEqual(a: string, b: string): boolean {
|
||||||
|
|||||||
@@ -5,7 +5,11 @@ import {
|
|||||||
createTransportBodySchema,
|
createTransportBodySchema,
|
||||||
transportIdParamSchema,
|
transportIdParamSchema,
|
||||||
} from '../schemas/transport.schema.js';
|
} 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<void> {
|
export async function transportRoutes(app: FastifyInstance): Promise<void> {
|
||||||
app.post('/transports', async (request, reply) => {
|
app.post('/transports', async (request, reply) => {
|
||||||
@@ -25,4 +29,12 @@ export async function transportRoutes(app: FastifyInstance): Promise<void> {
|
|||||||
await connectTransport({ transportId: id, dtlsParameters });
|
await connectTransport({ transportId: id, dtlsParameters });
|
||||||
return { ok: true as const };
|
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 };
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -143,6 +143,29 @@ export function getTransportStats(): { total: number } {
|
|||||||
return { total: transportMap.size };
|
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,生产环境调用会抛错 */
|
/** 测试专用:复位 transportMap,生产环境调用会抛错 */
|
||||||
export function _clearTransportMap(): void {
|
export function _clearTransportMap(): void {
|
||||||
assertTestOnly('_clearTransportMap');
|
assertTestOnly('_clearTransportMap');
|
||||||
|
|||||||
@@ -99,6 +99,29 @@ step_validate_env() {
|
|||||||
warn "JWT_SECRET 长度仅 ${#JWT_SECRET} 字符,推荐 >= 64 字符(openssl rand -hex 32)"
|
warn "JWT_SECRET 长度仅 ${#JWT_SECRET} 字符,推荐 >= 64 字符(openssl rand -hex 32)"
|
||||||
fi
|
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 开关一致性
|
# TURN 开关一致性
|
||||||
if [[ "${TURN_ENABLED:-false}" != "true" ]]; then
|
if [[ "${TURN_ENABLED:-false}" != "true" ]]; then
|
||||||
warn "TURN_ENABLED=${TURN_ENABLED:-false},公网部署下建议开启 TURN 作为对称 NAT fallback"
|
warn "TURN_ENABLED=${TURN_ENABLED:-false},公网部署下建议开启 TURN 作为对称 NAT fallback"
|
||||||
|
|||||||
Reference in New Issue
Block a user