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

237 lines
11 KiB
Go
Raw 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/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
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)
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)
svc := NewService(cfg, cf, srsP, hub)
s := &Server{cfg: cfg, svc: svc, hub: hub, 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 {
cf := cloudflare.NewProvider(cfg.CFAppID, cfg.CFAppSecret, cfg.CFBaseURL, cfg.CFStunURL)
srsP := srs.NewProvider(cfg.SRSBaseURL, cfg.SRSApp, cfg.SRSSecret, cfg.SRSCandidate)
svc := NewService(cfg, cf, srsP, hub)
rs := room.NewStore(nil)
if hub != nil && hub.DB() != nil {
rs = room.NewStore(hub.DB())
_ = room.InitSchema(hub.DB())
}
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) 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("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("GET /api/room/{room}/events", s.authWrap(http.HandlerFunc(s.handleRoomEvents), "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))
// 主播独立外部推送接口:仅 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("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("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))
mux.Handle("HEAD /live/", s.authWrap(s.srsHlsProxy, "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))
}
}