feat(distro): sfu distributor abstraction + fanout manager + cdn stub

This commit is contained in:
noelorin 2026-08-23 14:01:26 +08:00
parent 478a758263
commit 91b4f01310
7 changed files with 211 additions and 2 deletions

View File

@ -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)

View File

@ -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 {

View File

@ -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
}

View File

@ -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
}

View File

@ -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 }

View File

@ -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
}

View File

@ -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 }