live-sfu-demo/internal/server/server.go

292 lines
14 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package server
import (
"bytes"
"embed"
"encoding/json"
"io/fs"
"log"
"net"
"net/http"
"strings"
"time"
"google.golang.org/grpc"
"sync-live/gen"
"sync-live/internal/auth"
"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"
)
//go:embed static
var staticFS embed.FS
// Server 聚合控制面(gRPC + JSON 网关)、媒体面反向代理与静态 UI,并集成登录与 Casbin 鉴权。
type Server struct {
cfg *config.Config
svc *Service
hub *roomHub
cdnP *cdn.Provider
srsProxy http.Handler
cfProxy http.Handler
srsHlsProxy http.Handler
auth *auth.Manager
chat *chatHub
chatRL *chatRateLimiter
roomSvc *room.Service
roomStore *room.Store
}
func New(cfg *config.Config) *Server {
var hub *roomHub
var rs *room.Store
var rsvc *room.Service
if cfg.DatabaseURL != "" {
if dbConn, err := db.Open(cfg.DatabaseURL); err != nil {
log.Printf("[warn] embedded db open failed (%v), falling back to memory hub: %s", err, cfg.DatabaseURL)
hub = newRoomHub()
} else {
hub = newRoomHubWithDB(dbConn)
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 {
hub = newRoomHub()
}
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.cdnP = cdnP
svc.SetDistributionManager(newDistributionManagerWithConfig(hub, cf, srsP, cdnP, cfg, rs))
s := &Server{cfg: cfg, svc: svc, hub: hub, cdnP: cdnP, roomStore: rs, roomSvc: rsvc}
s.srsProxy = s.srsProxyHandler()
s.cfProxy = s.cfProxyHandler()
s.srsHlsProxy = s.srsHlsProxyHandler()
mgr, err := auth.NewManager(cfg.JWTSecret, cfg.JWTTTL, cfg.AuthModel, cfg.AuthPolicy, cfg.AuthUserFile)
if err != nil {
log.Printf("[warn] auth manager init failed (%v), running without auth", err)
} else {
s.auth = mgr
log.Printf("[auth] casbin enabled: model=%s policy=%s users=%s", cfg.AuthModel, cfg.AuthPolicy, cfg.AuthUserFile)
svc.auth = mgr
s.chat = newChatHub()
s.chatRL = newChatRateLimiter(5, 3*time.Second)
}
return s
}
func NewWithHub(cfg *config.Config, hub *roomHub) *Server {
rs := room.NewStore(nil)
if hub != nil && hub.DB() != nil {
rs = room.NewStore(hub.DB())
_ = room.InitSchema(hub.DB())
_ = room.InitPlaylistSchema(hub.DB())
}
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.cdnP = cdnP
svc.SetDistributionManager(newDistributionManagerWithConfig(hub, cf, srsP, cdnP, cfg, rs))
s := &Server{cfg: cfg, svc: svc, hub: hub, roomStore: rs, roomSvc: room.NewService(rs)}
s.srsProxy = s.srsProxyHandler()
s.cfProxy = s.cfProxyHandler()
s.srsHlsProxy = s.srsHlsProxyHandler()
mgr, err := auth.NewManager(cfg.JWTSecret, cfg.JWTTTL, cfg.AuthModel, cfg.AuthPolicy, cfg.AuthUserFile)
if err == nil {
s.auth = mgr
svc.auth = mgr
s.chat = newChatHub()
s.chatRL = newChatRateLimiter(5, 3*time.Second)
}
return s
}
func (s *Server) ApplyRuntimeConfig() {
if s.svc == nil || s.cfg == nil {
return
}
if s.svc.cf != nil {
s.svc.cf.ApplyConfig(s.cfg.CFAppID, s.cfg.CFAppSecret, s.cfg.CFBaseURL, s.cfg.CFStunURL)
}
if s.svc.srsP != nil {
s.svc.srsP.ApplyConfig(s.cfg.SRSBaseURL, s.cfg.SRSHttpURL, s.cfg.SRSApp, s.cfg.SRSSecret, s.cfg.SRSCandidate)
}
if s.svc.cdnP != nil {
s.svc.cdnP.ApplyConfig(s.cfg.CDNRTMPURL, s.cfg.CDNName)
} else if s.cdnP != nil {
s.cdnP.ApplyConfig(s.cfg.CDNRTMPURL, s.cfg.CDNName)
}
if dm, ok := s.svc.distributionManager.(*distributionManager); ok {
dm.UpdateEnabled(s.cfg.DistributorList())
dm.SetConfig(s.cfg)
dm.SetRoomStore(s.roomStore)
}
s.srsProxy = s.srsProxyHandler()
s.cfProxy = s.cfProxyHandler()
s.srsHlsProxy = s.srsHlsProxyHandler()
}
func (s *Server) Auth() *auth.Manager { return s.auth }
func (s *Server) StartGRPC() error {
lis, err := net.Listen("tcp", ":"+s.cfg.GRPCPort)
if err != nil {
return err
}
var opts []grpc.ServerOption
if s.auth != nil {
opts = append(opts,
grpc.UnaryInterceptor(s.auth.UnaryAuthInterceptor()),
grpc.StreamInterceptor(s.auth.StreamAuthInterceptor()),
)
}
gs := grpc.NewServer(opts...)
gen.RegisterSyncLiveServer(gs, NewGRPCServer(s.svc))
go func() { _ = gs.Serve(lis) }()
return nil
}
// handleCreateRoom 创建一个房间占位(发布前即可出现在房间列表)。
func (s *Server) handleCreateRoom(w http.ResponseWriter, r *http.Request) {
var body struct {
Name string `json:"name"`
}
_ = json.NewDecoder(r.Body).Decode(&body)
name := strings.TrimSpace(body.Name)
if name == "" {
http.Error(w, "room name required", http.StatusBadRequest)
return
}
s.hub.createRoom(name)
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]any{"ok": true, "room": name})
}
func (s *Server) Handler() http.Handler {
mux := http.NewServeMux()
sub, _ := fs.Sub(staticFS, "static")
index := s.fileServer("static/index.html")
// SPA 入口:发布/观看/登录均为前端路由,统一回退 index.html
mux.Handle("GET /", index)
mux.Handle("GET /publish", index)
mux.Handle("GET /watch", index)
mux.Handle("GET /login", index)
mux.Handle("GET /manage", index)
mux.Handle("GET /manage/", index)
mux.Handle("GET /static/", http.StripPrefix("/static/", http.FileServer(http.FS(sub))))
mux.HandleFunc("POST /api/auth/login", s.handleLogin)
mux.HandleFunc("POST /api/auth/register", s.handleRegister)
mux.HandleFunc("POST /api/auth/logout", s.handleLogout)
mux.HandleFunc("POST /api/auth/refresh", s.handleRefresh)
mux.Handle("GET /api/auth/me", s.authWrap(http.HandlerFunc(s.handleMe), "user", "list", true))
mux.Handle("GET /api/auth/users", s.authWrap(http.HandlerFunc(s.handleListUsers), "user", "list", true))
mux.Handle("POST /api/auth/users/role", s.authWrap(http.HandlerFunc(s.handleUpdateRole), "user", "manage", true))
mux.Handle("POST /api/auth/users/ban", s.authWrap(http.HandlerFunc(s.handleBanUser), "user", "manage", true))
mux.Handle("POST /api/auth/users/unban", s.authWrap(http.HandlerFunc(s.handleUnbanUser), "user", "manage", true))
mux.Handle("GET /api/auth/check", s.authWrap(http.HandlerFunc(s.handleAuthCheck), "config", "read", false))
// 对标 SyncTV 的分页查询:GET /api/rooms/query?page=1&page_size=20&search=&status=active
mux.Handle("GET /api/rooms/query", s.authWrap(http.HandlerFunc(s.handleRoomsQuery), "room", "list", true))
mux.Handle("POST /api/rooms/events", s.authWrap(http.HandlerFunc(s.handleRoomsEvents), "room", "list", true))
mux.Handle("GET /api/config", s.authWrap(http.HandlerFunc(s.handleConfig), "config", "read", false))
mux.Handle("GET /api/rooms", s.authWrap(http.HandlerFunc(s.handleRooms), "room", "list", true))
mux.Handle("POST /api/rooms", s.authWrap(http.HandlerFunc(s.handleCreateRoom), "room", "list", true))
mux.Handle("POST /api/publish", s.authWrap(http.HandlerFunc(s.handlePublish), "room", "publish", true))
mux.Handle("POST /api/subscribe", s.authWrap(http.HandlerFunc(s.handleSubscribe), "room", "subscribe", true))
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("POST /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("POST /api/room/{room}/chat/subscribe", 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))
// 主播独立外部推送接口:仅 publisher/admin(room:publish),供 OBS/机器人/管理工具以 Bearer token 调用,消息标记 host
mux.Handle("POST /api/room/{room}/broadcast", s.authWrap(http.HandlerFunc(s.handleRoomBroadcast), "room", "publish", true))
// ---- Manage 统一管理面 ----
mux.Handle("GET /api/manage/overview", s.authWrap(http.HandlerFunc(s.handleManageOverview), "config", "read", true))
mux.Handle("GET /api/manage/rooms", s.authWrap(http.HandlerFunc(s.handleManageRooms), "room", "list", true))
mux.Handle("POST /api/manage/rooms", s.authWrap(http.HandlerFunc(s.handleManageRooms), "room", "list", true))
mux.Handle("GET /api/manage/rooms/{room}", s.authWrap(http.HandlerFunc(s.handleManageRoomDetail), "room", "list", true))
mux.Handle("PUT /api/manage/rooms/{room}", s.authWrap(http.HandlerFunc(s.handleManageRoomDetail), "room", "list", true))
mux.Handle("PATCH /api/manage/rooms/{room}", s.authWrap(http.HandlerFunc(s.handleManageRoomDetail), "room", "list", true))
mux.Handle("DELETE /api/manage/rooms/{room}", s.authWrap(http.HandlerFunc(s.handleManageRoomDetail), "room", "list", true))
mux.Handle("GET /api/manage/rooms/{room}/members", s.authWrap(http.HandlerFunc(s.handleManageMembers), "room", "list", true))
mux.Handle("POST /api/manage/rooms/{room}/members", s.authWrap(http.HandlerFunc(s.handleManageMembers), "room", "list", true))
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}/members/{user}/ban", s.authWrap(http.HandlerFunc(s.handleManageMemberBan), "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/config", s.authWrap(http.HandlerFunc(s.handleManageRuntimeConfig), "config", "read", true))
mux.Handle("PUT /api/manage/config", s.authWrap(http.HandlerFunc(s.handleManageRuntimeConfig), "config", "write", true))
mux.Handle("POST /api/manage/config", s.authWrap(http.HandlerFunc(s.handleManageRuntimeConfig), "config", "write", true))
mux.Handle("POST /api/manage/config/reset", s.authWrap(http.HandlerFunc(s.handleManageRuntimeConfigReset), "config", "write", 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))
mux.Handle("PUT /api/manage/rooms/{room}/settings", s.authWrap(http.HandlerFunc(s.handleManageRoomSettings), "room", "list", true))
mux.Handle("GET /api/manage/rooms/{room}/distribution", s.authWrap(http.HandlerFunc(s.handleManageRoomDistribution), "room", "list", true))
mux.Handle("PUT /api/manage/rooms/{room}/distribution", s.authWrap(http.HandlerFunc(s.handleManageRoomDistribution), "room", "list", true))
mux.Handle("POST /api/manage/rooms/{room}/distribution", s.authWrap(http.HandlerFunc(s.handleManageRoomDistribution), "room", "list", true))
mux.Handle("PATCH /api/manage/rooms/{room}/distribution", s.authWrap(http.HandlerFunc(s.handleManageRoomDistribution), "room", "list", true))
mux.Handle("DELETE /api/manage/rooms/{room}/distribution", s.authWrap(http.HandlerFunc(s.handleManageRoomDistribution), "room", "list", true))
mux.Handle("POST /api/manage/rooms/{room}/join", s.authWrap(http.HandlerFunc(s.handleManageJoin), "room", "subscribe", true))
mux.Handle("POST /api/manage/rooms/{room}/leave", s.authWrap(http.HandlerFunc(s.handleManageLeave), "room", "subscribe", true))
mux.Handle("GET /rtc/v1/", s.authWrap(s.srsProxy, "room", "publish", false))
mux.Handle("POST /rtc/v1/", s.authWrap(s.srsProxy, "room", "publish", false))
mux.Handle("PUT /rtc/v1/", s.authWrap(s.srsProxy, "room", "publish", false))
mux.Handle("DELETE /rtc/v1/", s.authWrap(s.srsProxy, "room", "publish", false))
mux.Handle("GET /api/cf/", s.authWrap(s.cfProxy, "room", "watch", false))
mux.Handle("POST /api/cf/", s.authWrap(s.cfProxy, "room", "publish", false))
mux.Handle("PUT /api/cf/", s.authWrap(s.cfProxy, "room", "publish", false))
mux.Handle("DELETE /api/cf/", s.authWrap(s.cfProxy, "room", "publish", false))
// SRS HLS/FLV 静态切片反代(:8080),前端经 /live/* 拉流
mux.Handle("GET /live/", s.authWrap(s.srsHlsProxy, "srs", "streams", false))
// CDN HLS 占位(未配置返回 502 JSON,便于前端区分)
mux.Handle("GET /live/cdn/", s.authWrap(s.cdnProxyHandler(), "srs", "streams", false))
return mux
}
func (s *Server) authWrap(next http.Handler, obj, act string, needAuth bool) http.Handler {
if s.auth == nil {
return next
}
return s.auth.AuthorizeMiddleware(obj, act, needAuth)(next)
}
func (s *Server) Close() error {
if s.hub != nil && s.hub.DB() != nil {
return s.hub.DB().Close()
}
return nil
}
func (s *Server) fileServer(name string) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
data, err := staticFS.ReadFile(name)
if err != nil {
http.Error(w, "not found", http.StatusNotFound)
return
}
http.ServeContent(w, r, name, time.Time{}, bytes.NewReader(data))
}
}