From 91b4f01310b6056d52e707c6a1b9bd72602974f1 Mon Sep 17 00:00:00 2001 From: noelorin Date: Sun, 23 Aug 2026 14:01:26 +0800 Subject: [PATCH] feat(distro): sfu distributor abstraction + fanout manager + cdn stub --- internal/server/distribution.go | 65 +++++++++++++++++++++++++++++ internal/server/proxy.go | 16 +++++++ internal/server/server.go | 18 ++++++++ internal/sfu/cdn/provider.go | 45 ++++++++++++++++++++ internal/sfu/cloudflare/provider.go | 28 ++++++++++++- internal/sfu/distributor.go | 18 ++++++++ internal/sfu/srs/provider.go | 23 +++++++++- 7 files changed, 211 insertions(+), 2 deletions(-) create mode 100644 internal/server/distribution.go create mode 100644 internal/sfu/cdn/provider.go create mode 100644 internal/sfu/distributor.go diff --git a/internal/server/distribution.go b/internal/server/distribution.go new file mode 100644 index 0000000..5d0efa3 --- /dev/null +++ b/internal/server/distribution.go @@ -0,0 +1,65 @@ +package server + +import ( + "context" + "time" + + "sync-live/gen" + "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 distributionConfig +} + +type distributionConfig struct { + enabled map[string]bool // kind -> enabled (from DistributorList) +} + +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, cfg: distributionConfig{enabled: enabled}} +} + +// Ensure fans out the trunk stream to all enabled distribution channels. +func (m *distributionManager) Ensure(ctx context.Context, room, stream string) error { + if m.hub == nil || stream == "" { + return nil + } + now := time.Now().Unix() + // SRS HLS is always ready (remux直出) + if m.cfg.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 m.cfg.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 m.cfg.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 +} + +var _ DistributionManagerIface = (*distributionManager)(nil) diff --git a/internal/server/proxy.go b/internal/server/proxy.go index 0c90f62..9789bb6 100644 --- a/internal/server/proxy.go +++ b/internal/server/proxy.go @@ -54,6 +54,22 @@ func (s *Server) srsHlsProxyHandler() http.Handler { }) } +// cdnProxyHandler 占位代理第三方 CDN HLS。未配置时 502,便于前端感知与联调。 +func (s *Server) cdnProxyHandler() http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if s.cfg.CDNRTMPURL == "" { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusBadGateway) + w.Write([]byte(`{"error":"cdn not configured: set CDN_RTMP_URL"}`)) + return + } + // 已配置但无实际 HLS 边缘时,返回 502 并附带 CDN 信息 + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusBadGateway) + w.Write([]byte(`{"error":"cdn hls not yet available: forward via RTMP not implemented"}`)) + }) +} + // cfProxyHandler 反向代理 Cloudflare Realtime REST,注入 AppSecret。 // 浏览器只交换 SDP(tracks/new),凭证不出服务端。 func (s *Server) cfProxyHandler() http.Handler { diff --git a/internal/server/server.go b/internal/server/server.go index e3deead..deba595 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -17,6 +17,7 @@ import ( "sync-live/internal/config" "sync-live/internal/db" "sync-live/internal/room" + "sync-live/internal/sfu/cdn" "sync-live/internal/sfu/cloudflare" "sync-live/internal/sfu/srs" ) @@ -52,6 +53,7 @@ func New(cfg *config.Config) *Server { log.Printf("[db] embedded db enabled: dsn=%s remote=%v", cfg.DSN(), cfg.IsRemoteTurso()) rs = room.NewStore(dbConn) _ = room.InitSchema(dbConn) + _ = room.InitPlaylistSchema(dbConn) rsvc = room.NewService(rs) } } else { @@ -59,7 +61,9 @@ func New(cfg *config.Config) *Server { } cf := cloudflare.NewProvider(cfg.CFAppID, cfg.CFAppSecret, cfg.CFBaseURL, cfg.CFStunURL) srsP := srs.NewProvider(cfg.SRSBaseURL, cfg.SRSApp, cfg.SRSSecret, cfg.SRSCandidate) + cdnP := cdn.NewProvider(cfg.CDNRTMPURL, cfg.CDNName) svc := NewService(cfg, cf, srsP, hub) + svc.SetDistributionManager(newDistributionManager(hub, cf, srsP, cdnP, cfg.DistributorList())) s := &Server{cfg: cfg, svc: svc, hub: hub, roomStore: rs, roomSvc: rsvc} s.srsProxy = s.srsProxyHandler() s.cfProxy = s.cfProxyHandler() @@ -80,11 +84,14 @@ func New(cfg *config.Config) *Server { func NewWithHub(cfg *config.Config, hub *roomHub) *Server { cf := cloudflare.NewProvider(cfg.CFAppID, cfg.CFAppSecret, cfg.CFBaseURL, cfg.CFStunURL) srsP := srs.NewProvider(cfg.SRSBaseURL, cfg.SRSApp, cfg.SRSSecret, cfg.SRSCandidate) + cdnP := cdn.NewProvider(cfg.CDNRTMPURL, cfg.CDNName) svc := NewService(cfg, cf, srsP, hub) + svc.SetDistributionManager(newDistributionManager(hub, cf, srsP, cdnP, cfg.DistributorList())) rs := room.NewStore(nil) if hub != nil && hub.DB() != nil { rs = room.NewStore(hub.DB()) _ = room.InitSchema(hub.DB()) + _ = room.InitPlaylistSchema(hub.DB()) } s := &Server{cfg: cfg, svc: svc, hub: hub, roomStore: rs, roomSvc: room.NewService(rs)} s.srsProxy = s.srsProxyHandler() @@ -169,6 +176,11 @@ func (s *Server) Handler() http.Handler { mux.Handle("POST /api/stop", s.authWrap(http.HandlerFunc(s.handleStop), "room", "stop", true)) mux.Handle("GET /api/srs/streams", s.authWrap(http.HandlerFunc(s.handleSRSStreams), "srs", "streams", true)) mux.Handle("GET /api/room/{room}/events", s.authWrap(http.HandlerFunc(s.handleRoomEvents), "room", "watch", true)) + mux.Handle("GET /api/room/{room}/playlist", s.authWrap(http.HandlerFunc(s.handleRoomPlaylist), "room", "watch", true)) + mux.Handle("POST /api/room/{room}/playlist", s.authWrap(http.HandlerFunc(s.handleRoomPlaylist), "room", "watch", true)) + mux.Handle("DELETE /api/room/{room}/playlist", s.authWrap(http.HandlerFunc(s.handleRoomPlaylist), "room", "watch", true)) + mux.Handle("GET /api/room/{room}/playback", s.authWrap(http.HandlerFunc(s.handleRoomPlayback), "room", "watch", true)) + mux.Handle("POST /api/room/{room}/playback", s.authWrap(http.HandlerFunc(s.handleRoomPlayback), "room", "watch", true)) mux.Handle("GET /api/room/{room}/chat", s.authWrap(http.HandlerFunc(s.handleRoomChat), "room", "chat", false)) // 发送与 SSE 订阅一致:guest 亦可发弹幕(RBAC 已授权 room:chat),needAuth=false mux.Handle("POST /api/room/{room}/chat", s.authWrap(http.HandlerFunc(s.handleRoomChatPost), "room", "chat", false)) @@ -188,6 +200,9 @@ func (s *Server) Handler() http.Handler { mux.Handle("PUT /api/manage/rooms/{room}/members/{user}", s.authWrap(http.HandlerFunc(s.handleManageMemberRole), "room", "list", true)) mux.Handle("DELETE /api/manage/rooms/{room}/members/{user}", s.authWrap(http.HandlerFunc(s.handleManageMemberRole), "room", "list", true)) mux.Handle("POST /api/manage/rooms/{room}/members/{user}/perms", s.authWrap(http.HandlerFunc(s.handleManageMemberPerms), "room", "list", true)) + mux.Handle("POST /api/manage/rooms/{room}/members/{user}/approve", s.authWrap(http.HandlerFunc(s.handleManageMemberApprove), "room", "list", true)) + mux.Handle("POST /api/manage/rooms/{room}/members/{user}/reject", s.authWrap(http.HandlerFunc(s.handleManageMemberReject), "room", "list", true)) + mux.Handle("POST /api/manage/rooms/{room}/transfer", s.authWrap(http.HandlerFunc(s.handleManageTransferOwnership), "room", "list", true)) mux.Handle("GET /api/manage/permissions", s.authWrap(http.HandlerFunc(s.handleManagePermissionsMatrix), "config", "read", true)) mux.Handle("GET /api/manage/permissions/{room}", s.authWrap(http.HandlerFunc(s.handleManagePermissionsMatrix), "room", "list", true)) mux.Handle("GET /api/manage/rooms/{room}/settings", s.authWrap(http.HandlerFunc(s.handleManageRoomSettings), "room", "list", true)) @@ -207,6 +222,9 @@ func (s *Server) Handler() http.Handler { // SRS HLS/FLV 静态切片反代(:8080),前端经 /live/* 拉流 mux.Handle("GET /live/", s.authWrap(s.srsHlsProxy, "srs", "streams", false)) mux.Handle("HEAD /live/", s.authWrap(s.srsHlsProxy, "srs", "streams", false)) + // CDN HLS 占位(未配置返回 502 JSON,便于前端区分) + mux.Handle("GET /live/cdn/", s.authWrap(s.cdnProxyHandler(), "srs", "streams", false)) + mux.Handle("HEAD /live/cdn/", s.authWrap(s.cdnProxyHandler(), "srs", "streams", false)) return mux } diff --git a/internal/sfu/cdn/provider.go b/internal/sfu/cdn/provider.go new file mode 100644 index 0000000..496d9a5 --- /dev/null +++ b/internal/sfu/cdn/provider.go @@ -0,0 +1,45 @@ +package cdn + +import ( + "context" + "fmt" + "strings" +) + +// Provider implements third-party CDN distribution via RTMP/SRT forward. +type Provider struct { + rtmpURL string + name string +} + +// NewProvider creates a CDN provider. rtmpURL is the base RTMP URL, e.g. rtmp://cdn.example.com/live +func NewProvider(rtmpURL, name string) *Provider { + if name == "" { + name = "CDN" + } + return &Provider{rtmpURL: rtmpURL, name: name} +} + +func (p *Provider) Kind() string { return "cdn" } +func (p *Provider) Name() string { return p.name } +func (p *Provider) Configured() bool { return strings.TrimSpace(p.rtmpURL) != "" } +func (p *Provider) RTMPURL() string { return p.rtmpURL } + +// Forward returns the RTMP push URL for the given stream. +// In the current implementation this is a stub that just composes the URL and returns ready status. +// Future: actually trigger SRS forward API or external pusher. +func (p *Provider) Forward(ctx context.Context, room, stream string) (*ForwardResult, error) { + if !p.Configured() { + return nil, fmt.Errorf("cdn not configured: CDN_RTMP_URL empty") + } + base := strings.TrimRight(p.rtmpURL, "/") + url := base + "/" + stream + return &ForwardResult{Kind: "cdn", URL: url, Status: "ready"}, nil +} + +// ForwardResult mirrors sfu.ForwardResult to avoid import cycle +type ForwardResult struct { + Kind string + URL string + Status string +} diff --git a/internal/sfu/cloudflare/provider.go b/internal/sfu/cloudflare/provider.go index 15e4509..64152a9 100644 --- a/internal/sfu/cloudflare/provider.go +++ b/internal/sfu/cloudflare/provider.go @@ -1,6 +1,11 @@ package cloudflare -import "sync-live/gen" +import ( + "context" + "fmt" + + "sync-live/gen" +) // Provider 暴露 Cloudflare Realtime 后端的能力信息,并持有底层 REST 客户端。 type Provider struct { @@ -22,6 +27,27 @@ func NewProvider(appID, appSecret, baseURL, stun string) *Provider { func (p *Provider) Name() string { return "cloudflare" } func (p *Provider) Primary() bool { return true } func (p *Provider) Configured() bool { return p.configured } + +func (p *Provider) Kind() string { return "cf" } + +func (p *Provider) Forward(ctx context.Context, room, stream string) (*ForwardResult, error) { + if !p.configured { + return nil, fmt.Errorf("cloudflare not configured") + } + if room == "" || stream == "" { + return nil, fmt.Errorf("room and stream required") + } + // Stub: actual relay would create a CF session and push stream via WHIP/tracks. + // Return forwarding status with a placeholder URL that frontend can use to detect readiness. + return &ForwardResult{Kind: "cf", URL: fmt.Sprintf("/api/cf/forward/%s", stream), Status: "forwarding"}, nil +} + +type ForwardResult struct { + Kind string + URL string + Status string +} + func (p *Provider) Client() *Client { return p.client } func (p *Provider) AppID() string { return p.appID } diff --git a/internal/sfu/distributor.go b/internal/sfu/distributor.go new file mode 100644 index 0000000..f3bf35a --- /dev/null +++ b/internal/sfu/distributor.go @@ -0,0 +1,18 @@ +package sfu + +import "context" + +// Distributor abstracts a distribution channel downstream of the SRS trunk. +type Distributor interface { + Kind() string // srs-hls / cf / cdn + Name() string + Configured() bool + Forward(ctx context.Context, room, stream string) (*ForwardResult, error) +} + +// ForwardResult describes the state of a distribution after Ensure. +type ForwardResult struct { + Kind string + URL string + Status string // ready | forwarding | error +} diff --git a/internal/sfu/srs/provider.go b/internal/sfu/srs/provider.go index 9aad87d..7720cf0 100644 --- a/internal/sfu/srs/provider.go +++ b/internal/sfu/srs/provider.go @@ -1,6 +1,11 @@ package srs -import "sync-live/gen" +import ( + "context" + "fmt" + + "sync-live/gen" +) // Provider 暴露 SRS 后端能力(本地 docker 默认可达)。 type Provider struct { @@ -17,6 +22,22 @@ func NewProvider(baseURL, app, secret, candidate string) *Provider { func (p *Provider) Name() string { return "srs" } func (p *Provider) Primary() bool { return false } func (p *Provider) Configured() bool { return true } + +func (p *Provider) Kind() string { return "srs-hls" } + +func (p *Provider) Forward(ctx context.Context, room, stream string) (*ForwardResult, error) { + if stream == "" { + return nil, fmt.Errorf("stream required") + } + url := fmt.Sprintf("/live/%s.m3u8", stream) + return &ForwardResult{Kind: "srs-hls", URL: url, Status: "ready"}, nil +} + +type ForwardResult struct { + Kind string + URL string + Status string +} func (p *Provider) BaseURL() string { return p.baseURL } func (p *Provider) App() string { return p.app } func (p *Provider) Secret() string { return p.secret }