From c2bcfdbbac6c6b15215a80e65905cf7c4cb0217c Mon Sep 17 00:00:00 2001 From: duoaohui <928970622@qq.com> Date: Mon, 18 May 2026 22:49:09 +0800 Subject: [PATCH] =?UTF-8?q?=E8=A7=86=E9=A2=91=E4=BC=9A=E8=AE=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../service/http_media_orchestrator.go | 166 +++++++++++++++--- backend/go-service/app/provider/wire_gen.go | 2 +- 2 files changed, 138 insertions(+), 30 deletions(-) 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 d0421a9..4cc99d0 100644 --- a/backend/go-service/app/meeting/service/http_media_orchestrator.go +++ b/backend/go-service/app/meeting/service/http_media_orchestrator.go @@ -15,9 +15,20 @@ import ( "github.com/echochat/backend/config" "github.com/echochat/backend/pkg/logs" + "github.com/redis/go-redis/v9" "go.uber.org/zap" ) +// routerCacheKeyPrefix Redis Key 前缀:持久化 roomCode → routerInfoCache +// 只在 CreateRouter / CloseRouter / Resolve* 几条路径上读写,不影响其他业务逻辑 +const routerCacheKeyPrefix = "media:router:" + +// routerCacheTTL Redis Key 生存期:meeting 一般 ≤ 12h,多给一倍容错 +const routerCacheTTL = 24 * time.Hour + +// routerCacheLookupTimeout Redis 查询超时(fallback 路径) +const routerCacheLookupTimeout = 300 * time.Millisecond + // ====== 错误定义(供上层 errors.Is 判定并映射为 WS/HTTP 错误语义)====== // ErrMediaResourceNotFound media-server 明确返回 404 时的错误 @@ -54,15 +65,29 @@ type HTTPMediaOrchestrator struct { cfg config.MediaServerConfig client *http.Client - // roomCode → *routerInfoCache 本地缓存,服务重启后会丢失(此时 Node 侧的 Router 也会随 Node 重启而释放,状态一致) + // roomCode → *routerInfoCache 本地缓存(一级缓存) // Task 9 起由 sync.Map 升级为 sync.Map<*routerInfoCache> + // 2026-05 增强:添加 Redis 二级缓存(践 routerCacheKeyPrefix), + // go-service 重启后能从 Redis 恢复映射,避免处于进行中会议的 Router 丢失 + // ⏬ 导致 录制 / 资源清理 等依赖 routerID 的路径 500 roomRouterInfos sync.Map + + // rdb Redis 二级缓存,仅用于 Router 映射水平容灾;nil 代表未注入时仅走内存 + rdb *redis.Client +} + +// routerCachePayload Redis 序列化有效负载(routerID + rtpCapabilities raw JSON) +// rtpCapabilities 是 Node 返回的不透明 JSON,原样当字符串存起来 +func (h *HTTPMediaOrchestrator) routerCacheKey(roomCode string) string { + return routerCacheKeyPrefix + roomCode } // NewHTTPMediaOrchestrator 构造真实 HTTP 客户端 // 由 wire 注入,在 app/provider/provider.go 中统一绑定为 MediaOrchestrator +// +// rdb 参数隐含 Redis 二级缓存;传 nil 时退化为仅内存缓存(重启丢存) // 创建时做一次性配置校验:base_url 必须非空,internal_token 必须非空 -func NewHTTPMediaOrchestrator(cfg *config.Config) *HTTPMediaOrchestrator { +func NewHTTPMediaOrchestrator(cfg *config.Config, rdb *redis.Client) *HTTPMediaOrchestrator { mc := cfg.MediaServer if mc.TimeoutMS <= 0 { // Task 16 Nit(代码审查 2026-04-23 第 15 条): @@ -97,6 +122,80 @@ func NewHTTPMediaOrchestrator(cfg *config.Config) *HTTPMediaOrchestrator { IdleConnTimeout: 90 * time.Second, }, }, + rdb: rdb, + } +} + +// persistRouterCache 将 Router 映射写入 Redis(失败仅警告,不中断业务) +// 调用者在内存缓存写入后调用,保证最终一致性 +func (h *HTTPMediaOrchestrator) persistRouterCache(ctx context.Context, roomCode string, info *routerInfoCache) { + if h.rdb == nil || info == nil || info.ID == "" { + return + } + funcName := "service.http_media_orchestrator.persistRouterCache" + payload := map[string]any{ + "id": info.ID, + "rtp_capabilities": json.RawMessage(info.RtpCapabilities), + } + data, err := json.Marshal(payload) + if err != nil { + logs.Warn(ctx, funcName, "序列化 router 缓存失败", zap.String("room_code", roomCode), zap.Error(err)) + return + } + cctx, cancel := context.WithTimeout(ctx, routerCacheLookupTimeout) + defer cancel() + if err := h.rdb.Set(cctx, h.routerCacheKey(roomCode), data, routerCacheTTL).Err(); err != nil { + logs.Warn(ctx, funcName, "写入 Redis 失败(仅记录,不影响主路径)", + zap.String("room_code", roomCode), zap.Error(err)) + } +} + +// loadRouterFromRedis 从 Redis 拉取 Router 映射,成功后同步回填内存一级缓存 +// 返回 (info, true) 表示命中;(_, false) 表示不存在 / Redis 不可用 +func (h *HTTPMediaOrchestrator) loadRouterFromRedis(roomCode string) (*routerInfoCache, bool) { + if h.rdb == nil { + return nil, false + } + funcName := "service.http_media_orchestrator.loadRouterFromRedis" + ctx, cancel := context.WithTimeout(context.Background(), routerCacheLookupTimeout) + defer cancel() + raw, err := h.rdb.Get(ctx, h.routerCacheKey(roomCode)).Result() + if err != nil { + if errors.Is(err, redis.Nil) { + return nil, false + } + logs.Warn(ctx, funcName, "读取 Redis 失败", zap.String("room_code", roomCode), zap.Error(err)) + return nil, false + } + var payload struct { + ID string `json:"id"` + RtpCapabilities json.RawMessage `json:"rtp_capabilities"` + } + if err := json.Unmarshal([]byte(raw), &payload); err != nil { + logs.Warn(ctx, funcName, "Redis 返回的 router 缓存 JSON 解析失败", + zap.String("room_code", roomCode), zap.Error(err)) + return nil, false + } + if payload.ID == "" { + return nil, false + } + info := &routerInfoCache{ID: payload.ID, RtpCapabilities: payload.RtpCapabilities} + // 回填内存一级缓存,后续调用不再走 Redis + h.roomRouterInfos.LoadOrStore(roomCode, info) + return info, true +} + +// purgeRouterCache 同时清理内存 + Redis 两级缓存,幂等 +func (h *HTTPMediaOrchestrator) purgeRouterCache(ctx context.Context, roomCode string) { + h.roomRouterInfos.Delete(roomCode) + if h.rdb == nil { + return + } + cctx, cancel := context.WithTimeout(ctx, routerCacheLookupTimeout) + defer cancel() + if err := h.rdb.Del(cctx, h.routerCacheKey(roomCode)).Err(); err != nil { + logs.Warn(ctx, "service.http_media_orchestrator.purgeRouterCache", "删除 Redis Key 失败", + zap.String("room_code", roomCode), zap.Error(err)) } } @@ -178,40 +277,43 @@ func (h *HTTPMediaOrchestrator) CreateRouter(ctx context.Context, roomCode strin } } + // 同步回写 Redis,保证重启后仍能查到 routerID + h.persistRouterCache(ctx, roomCode, info) + logs.Info(ctx, funcName, "Node Router 创建成功", zap.String("room_code", roomCode), zap.String("router_id", resp.RouterID)) return resp.RouterID, nil } -// ResolveRouterID 从本地缓存读取 roomCode 对应的 routerID(不触发 HTTP) -// 返回 (id, true) 表示命中;返回 (_, false) 表示缓存缺失(通常意味着该房间尚未创建 Router) -// 供 MeetingService.JoinRoom 等需复用已有 Router 的场景使用,避免重复调 Node +// ResolveRouterID 读取 roomCode 对应的 routerID(不触发 HTTP) +// 查找顺序:内存 → Redis fallback (命中后回填内存) +// 返回 (id, true) 表示命中;(_, false) 表示缓存缺失 func (h *HTTPMediaOrchestrator) ResolveRouterID(roomCode string) (string, bool) { - v, ok := h.roomRouterInfos.Load(roomCode) - if !ok { - return "", false + if v, ok := h.roomRouterInfos.Load(roomCode); ok { + if info, _ := v.(*routerInfoCache); info != nil && info.ID != "" { + return info.ID, true + } } - info, _ := v.(*routerInfoCache) - if info == nil || info.ID == "" { - return "", false + if info, ok := h.loadRouterFromRedis(roomCode); ok { + return info.ID, true } - return info.ID, true + return "", false } // ResolveRouterInfo 读取 roomCode 对应的 Router 完整信息(routerID + rtpCapabilities) // Task 9 引入:前端 mediasoup-client Device.load 需要 rtpCapabilities,JoinRoom 响应会回填此字段 -// 返回 (id, caps, true) 命中;(_, _, false) 缺失 +// 查找顺序同 ResolveRouterID:内存 → Redis fallback func (h *HTTPMediaOrchestrator) ResolveRouterInfo(roomCode string) (string, json.RawMessage, bool) { - v, ok := h.roomRouterInfos.Load(roomCode) - if !ok { - return "", nil, false + if v, ok := h.roomRouterInfos.Load(roomCode); ok { + if info, _ := v.(*routerInfoCache); info != nil && info.ID != "" { + return info.ID, info.RtpCapabilities, true + } } - info, _ := v.(*routerInfoCache) - if info == nil || info.ID == "" { - return "", nil, false + if info, ok := h.loadRouterFromRedis(roomCode); ok { + return info.ID, info.RtpCapabilities, true } - return info.ID, info.RtpCapabilities, true + return "", nil, false } // CloseRouter 调用 DELETE /internal/v1/routers/:routerId @@ -222,14 +324,17 @@ func (h *HTTPMediaOrchestrator) ResolveRouterInfo(roomCode string) (string, json func (h *HTTPMediaOrchestrator) CloseRouter(ctx context.Context, roomCode string) error { funcName := "service.http_media_orchestrator.CloseRouter" - v, ok := h.roomRouterInfos.Load(roomCode) - if !ok { - logs.Debug(ctx, funcName, "无本地映射,跳过 Close", zap.String("room_code", roomCode)) - return nil + // 优先从内存读,缺失时从 Redis fallback + var info *routerInfoCache + if v, ok := h.roomRouterInfos.Load(roomCode); ok { + info, _ = v.(*routerInfoCache) } - info, _ := v.(*routerInfoCache) if info == nil || info.ID == "" { - h.roomRouterInfos.Delete(roomCode) + info, _ = h.loadRouterFromRedis(roomCode) + } + if info == nil || info.ID == "" { + logs.Debug(ctx, funcName, "无本地/Redis 映射,跳过 Close", zap.String("room_code", roomCode)) + h.purgeRouterCache(ctx, roomCode) // 避免 Redis 残留染色 return nil } routerID := info.ID @@ -238,9 +343,9 @@ func (h *HTTPMediaOrchestrator) CloseRouter(ctx context.Context, roomCode string zap.String("room_code", roomCode), zap.String("router_id", routerID), }) - // 成功或 ResourceNotFound 都从本地缓存删除(幂等) + // 成功或 ResourceNotFound 都从两级缓存删除(幂等) if err == nil || errors.Is(err, ErrMediaResourceNotFound) { - h.roomRouterInfos.Delete(roomCode) + h.purgeRouterCache(ctx, roomCode) return nil } return err @@ -498,9 +603,12 @@ func (h *HTTPMediaOrchestrator) StopRecording(ctx context.Context, recordingID s // ====== 内部工具 ====== -// routerIDByRoomCode 从本地缓存反查 routerID,缺失时返回 ErrMediaResourceNotFound +// routerIDByRoomCode 反查 routerID:内存 → Redis fallback,缺失时返回 ErrMediaResourceNotFound // 此错误语义表达的是"meeting 业务侧尚未(或已清理了)为该房间创建 Router",调用方通常应转为"会议未开始/已结束" func (h *HTTPMediaOrchestrator) routerIDByRoomCode(roomCode string) (string, error) { + if id, ok := h.ResolveRouterID(roomCode); ok { + return id, nil + } v, ok := h.roomRouterInfos.Load(roomCode) if !ok { return "", fmt.Errorf("%w: no router mapped for room_code=%s", ErrMediaResourceNotFound, roomCode) diff --git a/backend/go-service/app/provider/wire_gen.go b/backend/go-service/app/provider/wire_gen.go index f3cf652..92add2f 100644 --- a/backend/go-service/app/provider/wire_gen.go +++ b/backend/go-service/app/provider/wire_gen.go @@ -111,7 +111,7 @@ func InitializeApp(cfg *config.Config) (*App, error) { meetingChatDAO := dao6.NewMeetingChatDAO(gormDB) meetingRecordingDAO := dao6.NewMeetingRecordingDAO(gormDB) meetingBroadcaster := service7.NewMeetingBroadcaster(meetingParticipantDAO, pubSub) - httpMediaOrchestrator := service7.NewHTTPMediaOrchestrator(cfg) + httpMediaOrchestrator := service7.NewHTTPMediaOrchestrator(cfg, client) meetingLifecycleService := service7.NewMeetingLifecycleService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, httpMediaOrchestrator, cfg) meetingService := service7.NewMeetingService(meetingRoomDAO, meetingParticipantDAO, meetingChatDAO, gormDB, client, meetingBroadcaster, notifyService, friendshipDAO, onlineService, httpMediaOrchestrator, meetingLifecycleService) meetingSignalService := service7.NewMeetingSignalService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, httpMediaOrchestrator, meetingLifecycleService)