视频会议
This commit is contained in:
@@ -15,9 +15,20 @@ import (
|
|||||||
|
|
||||||
"github.com/echochat/backend/config"
|
"github.com/echochat/backend/config"
|
||||||
"github.com/echochat/backend/pkg/logs"
|
"github.com/echochat/backend/pkg/logs"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
"go.uber.org/zap"
|
"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 错误语义)======
|
// ====== 错误定义(供上层 errors.Is 判定并映射为 WS/HTTP 错误语义)======
|
||||||
|
|
||||||
// ErrMediaResourceNotFound media-server 明确返回 404 时的错误
|
// ErrMediaResourceNotFound media-server 明确返回 404 时的错误
|
||||||
@@ -54,15 +65,29 @@ type HTTPMediaOrchestrator struct {
|
|||||||
cfg config.MediaServerConfig
|
cfg config.MediaServerConfig
|
||||||
client *http.Client
|
client *http.Client
|
||||||
|
|
||||||
// roomCode → *routerInfoCache 本地缓存,服务重启后会丢失(此时 Node 侧的 Router 也会随 Node 重启而释放,状态一致)
|
// roomCode → *routerInfoCache 本地缓存(一级缓存)
|
||||||
// Task 9 起由 sync.Map<string> 升级为 sync.Map<*routerInfoCache>
|
// Task 9 起由 sync.Map<string> 升级为 sync.Map<*routerInfoCache>
|
||||||
|
// 2026-05 增强:添加 Redis 二级缓存(践 routerCacheKeyPrefix),
|
||||||
|
// go-service 重启后能从 Redis 恢复映射,避免处于进行中会议的 Router 丢失
|
||||||
|
// ⏬ 导致 录制 / 资源清理 等依赖 routerID 的路径 500
|
||||||
roomRouterInfos sync.Map
|
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 客户端
|
// NewHTTPMediaOrchestrator 构造真实 HTTP 客户端
|
||||||
// 由 wire 注入,在 app/provider/provider.go 中统一绑定为 MediaOrchestrator
|
// 由 wire 注入,在 app/provider/provider.go 中统一绑定为 MediaOrchestrator
|
||||||
|
//
|
||||||
|
// rdb 参数隐含 Redis 二级缓存;传 nil 时退化为仅内存缓存(重启丢存)
|
||||||
// 创建时做一次性配置校验:base_url 必须非空,internal_token 必须非空
|
// 创建时做一次性配置校验:base_url 必须非空,internal_token 必须非空
|
||||||
func NewHTTPMediaOrchestrator(cfg *config.Config) *HTTPMediaOrchestrator {
|
func NewHTTPMediaOrchestrator(cfg *config.Config, rdb *redis.Client) *HTTPMediaOrchestrator {
|
||||||
mc := cfg.MediaServer
|
mc := cfg.MediaServer
|
||||||
if mc.TimeoutMS <= 0 {
|
if mc.TimeoutMS <= 0 {
|
||||||
// Task 16 Nit(代码审查 2026-04-23 第 15 条):
|
// Task 16 Nit(代码审查 2026-04-23 第 15 条):
|
||||||
@@ -97,6 +122,80 @@ func NewHTTPMediaOrchestrator(cfg *config.Config) *HTTPMediaOrchestrator {
|
|||||||
IdleConnTimeout: 90 * time.Second,
|
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 创建成功",
|
logs.Info(ctx, funcName, "Node Router 创建成功",
|
||||||
zap.String("room_code", roomCode),
|
zap.String("room_code", roomCode),
|
||||||
zap.String("router_id", resp.RouterID))
|
zap.String("router_id", resp.RouterID))
|
||||||
return resp.RouterID, nil
|
return resp.RouterID, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// ResolveRouterID 从本地缓存读取 roomCode 对应的 routerID(不触发 HTTP)
|
// ResolveRouterID 读取 roomCode 对应的 routerID(不触发 HTTP)
|
||||||
// 返回 (id, true) 表示命中;返回 (_, false) 表示缓存缺失(通常意味着该房间尚未创建 Router)
|
// 查找顺序:内存 → Redis fallback (命中后回填内存)
|
||||||
// 供 MeetingService.JoinRoom 等需复用已有 Router 的场景使用,避免重复调 Node
|
// 返回 (id, true) 表示命中;(_, false) 表示缓存缺失
|
||||||
func (h *HTTPMediaOrchestrator) ResolveRouterID(roomCode string) (string, bool) {
|
func (h *HTTPMediaOrchestrator) ResolveRouterID(roomCode string) (string, bool) {
|
||||||
v, ok := h.roomRouterInfos.Load(roomCode)
|
if v, ok := h.roomRouterInfos.Load(roomCode); ok {
|
||||||
if !ok {
|
if info, _ := v.(*routerInfoCache); info != nil && info.ID != "" {
|
||||||
return "", false
|
|
||||||
}
|
|
||||||
info, _ := v.(*routerInfoCache)
|
|
||||||
if info == nil || info.ID == "" {
|
|
||||||
return "", false
|
|
||||||
}
|
|
||||||
return info.ID, true
|
return info.ID, true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if info, ok := h.loadRouterFromRedis(roomCode); ok {
|
||||||
|
return info.ID, true
|
||||||
|
}
|
||||||
|
return "", false
|
||||||
}
|
}
|
||||||
|
|
||||||
// ResolveRouterInfo 读取 roomCode 对应的 Router 完整信息(routerID + rtpCapabilities)
|
// ResolveRouterInfo 读取 roomCode 对应的 Router 完整信息(routerID + rtpCapabilities)
|
||||||
// Task 9 引入:前端 mediasoup-client Device.load 需要 rtpCapabilities,JoinRoom 响应会回填此字段
|
// 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) {
|
func (h *HTTPMediaOrchestrator) ResolveRouterInfo(roomCode string) (string, json.RawMessage, bool) {
|
||||||
v, ok := h.roomRouterInfos.Load(roomCode)
|
if v, ok := h.roomRouterInfos.Load(roomCode); ok {
|
||||||
if !ok {
|
if info, _ := v.(*routerInfoCache); info != nil && info.ID != "" {
|
||||||
return "", nil, false
|
|
||||||
}
|
|
||||||
info, _ := v.(*routerInfoCache)
|
|
||||||
if info == nil || info.ID == "" {
|
|
||||||
return "", nil, false
|
|
||||||
}
|
|
||||||
return info.ID, info.RtpCapabilities, true
|
return info.ID, info.RtpCapabilities, true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if info, ok := h.loadRouterFromRedis(roomCode); ok {
|
||||||
|
return info.ID, info.RtpCapabilities, true
|
||||||
|
}
|
||||||
|
return "", nil, false
|
||||||
}
|
}
|
||||||
|
|
||||||
// CloseRouter 调用 DELETE /internal/v1/routers/:routerId
|
// 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 {
|
func (h *HTTPMediaOrchestrator) CloseRouter(ctx context.Context, roomCode string) error {
|
||||||
funcName := "service.http_media_orchestrator.CloseRouter"
|
funcName := "service.http_media_orchestrator.CloseRouter"
|
||||||
|
|
||||||
v, ok := h.roomRouterInfos.Load(roomCode)
|
// 优先从内存读,缺失时从 Redis fallback
|
||||||
if !ok {
|
var info *routerInfoCache
|
||||||
logs.Debug(ctx, funcName, "无本地映射,跳过 Close", zap.String("room_code", roomCode))
|
if v, ok := h.roomRouterInfos.Load(roomCode); ok {
|
||||||
return nil
|
info, _ = v.(*routerInfoCache)
|
||||||
}
|
}
|
||||||
info, _ := v.(*routerInfoCache)
|
|
||||||
if info == nil || info.ID == "" {
|
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
|
return nil
|
||||||
}
|
}
|
||||||
routerID := info.ID
|
routerID := info.ID
|
||||||
@@ -238,9 +343,9 @@ func (h *HTTPMediaOrchestrator) CloseRouter(ctx context.Context, roomCode string
|
|||||||
zap.String("room_code", roomCode),
|
zap.String("room_code", roomCode),
|
||||||
zap.String("router_id", routerID),
|
zap.String("router_id", routerID),
|
||||||
})
|
})
|
||||||
// 成功或 ResourceNotFound 都从本地缓存删除(幂等)
|
// 成功或 ResourceNotFound 都从两级缓存删除(幂等)
|
||||||
if err == nil || errors.Is(err, ErrMediaResourceNotFound) {
|
if err == nil || errors.Is(err, ErrMediaResourceNotFound) {
|
||||||
h.roomRouterInfos.Delete(roomCode)
|
h.purgeRouterCache(ctx, roomCode)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
return err
|
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",调用方通常应转为"会议未开始/已结束"
|
// 此错误语义表达的是"meeting 业务侧尚未(或已清理了)为该房间创建 Router",调用方通常应转为"会议未开始/已结束"
|
||||||
func (h *HTTPMediaOrchestrator) routerIDByRoomCode(roomCode string) (string, error) {
|
func (h *HTTPMediaOrchestrator) routerIDByRoomCode(roomCode string) (string, error) {
|
||||||
|
if id, ok := h.ResolveRouterID(roomCode); ok {
|
||||||
|
return id, nil
|
||||||
|
}
|
||||||
v, ok := h.roomRouterInfos.Load(roomCode)
|
v, ok := h.roomRouterInfos.Load(roomCode)
|
||||||
if !ok {
|
if !ok {
|
||||||
return "", fmt.Errorf("%w: no router mapped for room_code=%s", ErrMediaResourceNotFound, roomCode)
|
return "", fmt.Errorf("%w: no router mapped for room_code=%s", ErrMediaResourceNotFound, roomCode)
|
||||||
|
|||||||
@@ -111,7 +111,7 @@ func InitializeApp(cfg *config.Config) (*App, error) {
|
|||||||
meetingChatDAO := dao6.NewMeetingChatDAO(gormDB)
|
meetingChatDAO := dao6.NewMeetingChatDAO(gormDB)
|
||||||
meetingRecordingDAO := dao6.NewMeetingRecordingDAO(gormDB)
|
meetingRecordingDAO := dao6.NewMeetingRecordingDAO(gormDB)
|
||||||
meetingBroadcaster := service7.NewMeetingBroadcaster(meetingParticipantDAO, pubSub)
|
meetingBroadcaster := service7.NewMeetingBroadcaster(meetingParticipantDAO, pubSub)
|
||||||
httpMediaOrchestrator := service7.NewHTTPMediaOrchestrator(cfg)
|
httpMediaOrchestrator := service7.NewHTTPMediaOrchestrator(cfg, client)
|
||||||
meetingLifecycleService := service7.NewMeetingLifecycleService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, httpMediaOrchestrator, cfg)
|
meetingLifecycleService := service7.NewMeetingLifecycleService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, httpMediaOrchestrator, cfg)
|
||||||
meetingService := service7.NewMeetingService(meetingRoomDAO, meetingParticipantDAO, meetingChatDAO, gormDB, client, meetingBroadcaster, notifyService, friendshipDAO, onlineService, httpMediaOrchestrator, meetingLifecycleService)
|
meetingService := service7.NewMeetingService(meetingRoomDAO, meetingParticipantDAO, meetingChatDAO, gormDB, client, meetingBroadcaster, notifyService, friendshipDAO, onlineService, httpMediaOrchestrator, meetingLifecycleService)
|
||||||
meetingSignalService := service7.NewMeetingSignalService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, httpMediaOrchestrator, meetingLifecycleService)
|
meetingSignalService := service7.NewMeetingSignalService(meetingRoomDAO, meetingParticipantDAO, client, meetingBroadcaster, httpMediaOrchestrator, meetingLifecycleService)
|
||||||
|
|||||||
Reference in New Issue
Block a user