package server import ( "context" "time" "sync-live/gen" "sync-live/internal/config" "sync-live/internal/room" "sync-live/internal/sfu/cdn" "sync-live/internal/sfu/cloudflare" "sync-live/internal/sfu/srs" ) // distributionManager handles fanout from SRS trunk to distribution channels. type distributionManager struct { hub *roomHub cf *cloudflare.Provider srsP *srs.Provider cdnP *cdn.Provider cfg *config.Config roomStore *room.Store enabled map[string]bool // global fallback } func newDistributionManager(hub *roomHub, cf *cloudflare.Provider, srsP *srs.Provider, cdnP *cdn.Provider, enabledKinds []string) *distributionManager { enabled := map[string]bool{} for _, k := range enabledKinds { enabled[k] = true } return &distributionManager{hub: hub, cf: cf, srsP: srsP, cdnP: cdnP, enabled: enabled} } func newDistributionManagerWithConfig(hub *roomHub, cf *cloudflare.Provider, srsP *srs.Provider, cdnP *cdn.Provider, cfg *config.Config, rs *room.Store) *distributionManager { enabled := map[string]bool{} if cfg != nil { for _, k := range cfg.DistributorList() { enabled[k] = true } } return &distributionManager{hub: hub, cf: cf, srsP: srsP, cdnP: cdnP, cfg: cfg, roomStore: rs, enabled: enabled} } // effectiveEnabled returns per-room enabled map with fallback to global. func (m *distributionManager) effectiveEnabled(room string) map[string]bool { // try per-room override if m.roomStore != nil && m.roomStore.DB() != nil && room != "" { if rm, err := m.roomStore.GetRoom(context.Background(), room); err == nil && rm != nil && rm.Settings != nil && rm.Settings.HasCustomDistributors() { eff := rm.Settings.EffectiveDistributors(nil) // if room has custom distributors, use exactly that list (even empty means none) // normalize via global parse ensures empty disables all if rm.Settings.Distributors != nil { m2 := map[string]bool{} for _, k := range eff { m2[k] = true } return m2 } } } // fallback to global config or cached enabled if m.cfg != nil { m2 := map[string]bool{} for _, k := range m.cfg.DistributorList() { m2[k] = true } return m2 } return m.enabled } // Ensure fans out the trunk stream to all enabled distribution channels (per-room with global fallback). func (m *distributionManager) Ensure(ctx context.Context, room, stream string) error { if m.hub == nil || stream == "" { return nil } enabled := m.effectiveEnabled(room) now := time.Now().Unix() // SRS HLS is always ready (remux直出) if enabled["srs-hls"] { if res, err := m.srsP.Forward(ctx, room, stream); err == nil { m.hub.setDistribution(room, &gen.DistributionTarget{ Kind: gen.DistributionKind_DISTRIBUTION_KIND_SRS_HLS, Url: res.URL, Status: res.Status, UpdatedAt: now, }) } } if enabled["cf"] && m.cf != nil && m.cf.Configured() { if res, err := m.cf.Forward(ctx, room, stream); err == nil { m.hub.setDistribution(room, &gen.DistributionTarget{ Kind: gen.DistributionKind_DISTRIBUTION_KIND_CF_SFU, Url: res.URL, Status: res.Status, UpdatedAt: now, }) } } if enabled["cdn"] && m.cdnP != nil && m.cdnP.Configured() { if res, err := m.cdnP.Forward(ctx, room, stream); err == nil { m.hub.setDistribution(room, &gen.DistributionTarget{ Kind: gen.DistributionKind_DISTRIBUTION_KIND_CDN, Url: res.URL, Status: res.Status, UpdatedAt: now, }) } } return nil } func (m *distributionManager) UpdateEnabled(kinds []string) { m.enabled = map[string]bool{} for _, k := range kinds { m.enabled[k] = true } } func (m *distributionManager) SetRoomStore(rs *room.Store) { m.roomStore = rs } func (m *distributionManager) SetConfig(cfg *config.Config) { m.cfg = cfg } var _ DistributionManagerIface = (*distributionManager)(nil)