feat: live SFU fan-out demo (Cloudflare Realtime + SRS)

- protobuf/gRPC control plane (LiveSFU service, buf generated)
- Cloudflare Realtime as primary SFU, SRS as secondary backend
- WHIP/WHEP + Cloudflare Realtime media proxies (server injects credentials)
- room fan-out state with SSE; publish/watch frontend
This commit is contained in:
xuhongyuan 2026-08-17 15:14:51 +08:00
commit e46db4d365
32 changed files with 3511 additions and 0 deletions

22
.env.example Normal file
View File

@ -0,0 +1,22 @@
# HTTP / gRPC 控制面端口
HTTP_PORT=8088
GRPC_PORT=9090
# 后端顺序(主 SFU 在前):cloudflare 优先,srs 本地对照
SFU_PROVIDER=cloudflare,srs
# SRS(本地对照后端:docker compose up -d srs)
SRS_API_BASE=http://localhost:1985
SRS_APP=live
SRS_SECRET=
SRS_CANDIDATE=127.0.0.1
# Cloudflare Realtime(主 SFU)
CF_APP_ID=
CF_APP_SECRET=
CF_BASE_URL=https://rtc.live.cloudflare.com/v1
CF_STUN_URL=stun:stun.cloudflare.com:3478
# 推流 JWT(SRS 入口鉴权,可选;设为 1 后 WHIP 必须带有效 token)
DEMO_TOKEN_SECRET=
SFU_TOKEN_REQUIRED=0

3
.gitignore vendored Normal file
View File

@ -0,0 +1,3 @@
.DS_Store
.env
*.local.env

76
README.md Normal file
View File

@ -0,0 +1,76 @@
# 直播 SFU 分流 Demo(Cloudflare Realtime · SRS)
参考 GOSpeak 技术栈的直播 SFU 分流演示。控制面用 **protobuf / gRPC** 定义,媒体面走
**Cloudflare Realtime**(主 SFU)与 **SRS**(本地对照)的 WebRTC 扇出。
- 一路推流 → SFU 扇出 → 多路拉流。
- 同一路发布可同时落到 **Cloudflare Realtime** 与 **SRS** 两条分发链路(分流)。
- 浏览器经服务端反向代理直连 SFU,凭证不出服务端。
## 目录结构
```
app/live-sfu-demo/
api/live_sfu.proto # protobuf 控制面契约
gen/ # buf generate 产出(Go)
internal/
config/ # 运行参数(对齐 GOSpeak 的 env 布局)
sfu/cloudflare/ # Cloudflare Realtime REST 客户端 + Provider
sfu/srs/ # SRS Provider
server/ # gRPC 服务 + JSON 网关 + SRS/CF 媒体反代 + 房间扇出状态
static/ # 浏览器 UI(发布 / 观看 / 分发面板)
cmd/server/main.go
deploy/ # docker-compose.yml + srs.conf
```
## 运行(SRS 链路,开箱即跑)
```bash
# 1) 起本地 SRS(WHIP/WHEP 后端)
cd app/live-sfu-demo/deploy && SRS_CANDIDATE=127.0.0.1 docker compose up -d srs
# 2) 起 Demo 控制面
cd app/live-sfu-demo
cp .env.example .env # 可选;SRS 链路无需任何凭证
go run ./cmd/server
# 3) 打开浏览器
# 发布: http://localhost:8088/publish?room=demo
# 观看: http://localhost:8088/watch?room=demo
```
两个标签页用同一房间名即可配对。发布页勾选的分发后端会在「分发状态」中实时显示。
## 运行(Cloudflare Realtime 链路,主 SFU)
在 `.env` 填入 Cloudflare Realtime 凭证后,`go run ./cmd/server` 即启用主 SFU:
```
CF_APP_ID=xxxxxxxxxxxx
CF_APP_SECRET=xxxxxxxxxxxx
```
- 发布:服务端用 `CF_APP_SECRET` 创建 Cloudflare session,浏览器只交换 SDP(`/api/cf/...` 反代注入 Bearer)。
- 观看:服务端为观众创建 viewer session,并订阅发布者 session 的轨道(`location=remote`)。
Cloudflare 未配置时,UI 会标注「未配置」,SRS 链路不受影响。
## protobuf / gRPC
```bash
# 修改 api/live_sfu.proto 后重新生成
export PATH="$HOME/go/bin:$PATH"
buf generate api
```
契约(`LiveSFU` 服务):`GetConfig` / `ListRooms` / `Publish` / `Subscribe` / `StopStream` /
`WatchRoom`(服务端流式推送房间分发拓扑)。gRPC 监听 `GRPC_PORT`(默认 9090),浏览器走同端口
的 JSON 网关(`protojson`)。
## 分流拓扑
```
publisher ──WHIP/tracks.new──▶ Cloudflare Realtime SFU ──▶ viewer(s)
└────WHIP────────────▶ SRS SFU (WHEP) ──▶ viewer(s)
房间分发目标由 WatchRoom 实时广播,观众任选后端拉流
```

7
api/buf.yaml Normal file
View File

@ -0,0 +1,7 @@
version: v2
lint:
use:
- DEFAULT
breaking:
use:
- FILE

107
api/live_sfu.proto Normal file
View File

@ -0,0 +1,107 @@
syntax = "proto3";
package gospeak.livedemo.v1;
option go_package = "gospeak-live-sfu-demo/gen";
// LiveSFU 是直播 SFU 分流 Demo 的控制面契约(protobuf / gRPC)。
//
// 媒体面(WebRTC SDP 交换)不在此契约内:浏览器与 Cloudflare Realtime / SRS
// 直连,经本服务反向代理转发(服务端注入 AppSecret / stream token)。
// 本契约只负责房间、后端、发布/订阅拓扑的协调与扇出状态广播。
service LiveSFU {
rpc GetConfig(GetConfigRequest) returns (GetConfigResponse);
rpc ListRooms(ListRoomsRequest) returns (ListRoomsResponse);
rpc Publish(PublishRequest) returns (PublishResponse);
rpc Subscribe(SubscribeRequest) returns (SubscribeResponse);
rpc StopStream(StopStreamRequest) returns (StopStreamResponse);
rpc WatchRoom(WatchRoomRequest) returns (stream RoomEvent);
}
enum BackendKind {
BACKEND_KIND_UNSPECIFIED = 0;
BACKEND_KIND_CLOUDFLARE = 1; // Cloudflare Realtime(主 SFU)
BACKEND_KIND_SRS = 2; // SRS(次级 / 本地对照)
}
message IceServer {
repeated string urls = 1;
string username = 2;
string credential = 3;
}
message GetConfigRequest {}
message BackendInfo {
BackendKind kind = 1;
string name = 2;
bool configured = 3;
bool primary = 4;
}
message GetConfigResponse {
repeated BackendInfo backends = 1;
string candidate = 2;
bool token_required = 3;
}
// StreamTarget 表示一个房间在某后端的一条分发目标(SFU 扇出出口)。
message StreamTarget {
BackendKind backend = 1;
string session_id = 2; // Cloudflare: 发布者 sessionId
string stream = 3; // SRS: stream 名
string publish_token = 4;// SRS: 推流 JWT(可选)
string url = 5; // 拉流播放地址(SRS WHEP,可选)
int64 published_at = 6;
}
message Room {
string name = 1;
repeated StreamTarget targets = 2;
}
message ListRoomsRequest {}
message ListRoomsResponse {
repeated Room rooms = 1;
}
message PublishRequest {
string room = 1;
BackendKind backend = 2;
string identity = 3;
}
message PublishResponse {
string session_id = 1;
string stream = 2;
string publish_token = 3;
repeated IceServer ice_servers = 4;
StreamTarget target = 5;
}
message SubscribeRequest {
string room = 1;
BackendKind backend = 2;
string identity = 3;
}
message SubscribeResponse {
string session_id = 1;
string stream = 2;
string publisher_session_id = 3;
repeated IceServer ice_servers = 4;
}
message StopStreamRequest {
string room = 1;
BackendKind backend = 2;
}
message StopStreamResponse {
bool ok = 1;
}
message WatchRoomRequest {
string room = 1;
}
message RoomEvent {
string room = 1;
repeated StreamTarget targets = 2;
}

13
buf.gen.yaml Normal file
View File

@ -0,0 +1,13 @@
version: v1
managed:
enabled: false
plugins:
- plugin: go
out: gen
opt:
- paths=source_relative
- plugin: go-grpc
out: gen
opt:
- paths=source_relative
- require_unimplemented_servers=false

24
cmd/server/main.go Normal file
View File

@ -0,0 +1,24 @@
package main
import (
"log"
"net/http"
"gospeak-live-sfu-demo/internal/config"
"gospeak-live-sfu-demo/internal/server"
)
func main() {
cfg := config.Load()
srv := server.New(cfg)
if err := srv.StartGRPC(); err != nil {
log.Printf("[warn] grpc control plane failed to start: %v", err)
}
addr := ":" + cfg.HTTPPort
log.Printf("live-sfu demo ready: http://localhost:%s (grpc :%s)", cfg.HTTPPort, cfg.GRPCPort)
log.Printf("backends (primary first): %v", cfg.ProviderList())
log.Printf("cloudflare configured: %v", cfg.CFAppID != "" && cfg.CFAppSecret != "")
if err := http.ListenAndServe(addr, srv.Handler()); err != nil {
log.Fatalf("http server error: %v", err)
}
}

16
deploy/docker-compose.yml Normal file
View File

@ -0,0 +1,16 @@
services:
# SRS:本地对照 SFU(WHIP/WHEP)。媒体走 :8000,信令/HTTP API 走 :1985。
srs:
image: ossrs/srs:6
container_name: live-sfu-srs
ports:
- "1935:1935"
- "1985:1985"
- "8080:8080"
- "8000:8000/udp"
- "8000:8000/tcp"
environment:
CANDIDATE: "${SRS_CANDIDATE:-127.0.0.1}"
volumes:
- ./srs.conf:/usr/local/srs/conf/srs.conf:ro
command: ./objs/srs -c conf/srs.conf

38
deploy/srs.conf Normal file
View File

@ -0,0 +1,38 @@
listen 1935;
max_connections 1000;
daemon off;
http_server {
enabled on;
listen 8080;
dir ./objs/nginx/html;
}
http_api {
enabled on;
listen 1985;
crossdomain on;
}
rtc_server {
enabled on;
listen 8000;
tcp {
enabled on;
listen 8000;
}
protocol all;
candidate $CANDIDATE;
}
vhost __defaultVhost__ {
rtc {
enabled on;
nack on;
twcc on;
}
http_remux {
enabled on;
mount [vhost]/[app]/[stream].flv;
}
}

1141
gen/live_sfu.pb.go Normal file

File diff suppressed because it is too large Load Diff

325
gen/live_sfu_grpc.pb.go Normal file
View File

@ -0,0 +1,325 @@
// Code generated by protoc-gen-go-grpc. DO NOT EDIT.
// versions:
// - protoc-gen-go-grpc v1.5.1
// - protoc (unknown)
// source: live_sfu.proto
package gen
import (
context "context"
grpc "google.golang.org/grpc"
codes "google.golang.org/grpc/codes"
status "google.golang.org/grpc/status"
)
// This is a compile-time assertion to ensure that this generated file
// is compatible with the grpc package it is being compiled against.
// Requires gRPC-Go v1.64.0 or later.
const _ = grpc.SupportPackageIsVersion9
const (
LiveSFU_GetConfig_FullMethodName = "/gospeak.livedemo.v1.LiveSFU/GetConfig"
LiveSFU_ListRooms_FullMethodName = "/gospeak.livedemo.v1.LiveSFU/ListRooms"
LiveSFU_Publish_FullMethodName = "/gospeak.livedemo.v1.LiveSFU/Publish"
LiveSFU_Subscribe_FullMethodName = "/gospeak.livedemo.v1.LiveSFU/Subscribe"
LiveSFU_StopStream_FullMethodName = "/gospeak.livedemo.v1.LiveSFU/StopStream"
LiveSFU_WatchRoom_FullMethodName = "/gospeak.livedemo.v1.LiveSFU/WatchRoom"
)
// LiveSFUClient is the client API for LiveSFU service.
//
// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
//
// LiveSFU 是直播 SFU 分流 Demo 的控制面契约(protobuf / gRPC)。
//
// 媒体面(WebRTC SDP 交换)不在此契约内:浏览器与 Cloudflare Realtime / SRS
// 直连,经本服务反向代理转发(服务端注入 AppSecret / stream token)。
// 本契约只负责房间、后端、发布/订阅拓扑的协调与扇出状态广播。
type LiveSFUClient interface {
GetConfig(ctx context.Context, in *GetConfigRequest, opts ...grpc.CallOption) (*GetConfigResponse, error)
ListRooms(ctx context.Context, in *ListRoomsRequest, opts ...grpc.CallOption) (*ListRoomsResponse, error)
Publish(ctx context.Context, in *PublishRequest, opts ...grpc.CallOption) (*PublishResponse, error)
Subscribe(ctx context.Context, in *SubscribeRequest, opts ...grpc.CallOption) (*SubscribeResponse, error)
StopStream(ctx context.Context, in *StopStreamRequest, opts ...grpc.CallOption) (*StopStreamResponse, error)
WatchRoom(ctx context.Context, in *WatchRoomRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[RoomEvent], error)
}
type liveSFUClient struct {
cc grpc.ClientConnInterface
}
func NewLiveSFUClient(cc grpc.ClientConnInterface) LiveSFUClient {
return &liveSFUClient{cc}
}
func (c *liveSFUClient) GetConfig(ctx context.Context, in *GetConfigRequest, opts ...grpc.CallOption) (*GetConfigResponse, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(GetConfigResponse)
err := c.cc.Invoke(ctx, LiveSFU_GetConfig_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *liveSFUClient) ListRooms(ctx context.Context, in *ListRoomsRequest, opts ...grpc.CallOption) (*ListRoomsResponse, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(ListRoomsResponse)
err := c.cc.Invoke(ctx, LiveSFU_ListRooms_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *liveSFUClient) Publish(ctx context.Context, in *PublishRequest, opts ...grpc.CallOption) (*PublishResponse, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(PublishResponse)
err := c.cc.Invoke(ctx, LiveSFU_Publish_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *liveSFUClient) Subscribe(ctx context.Context, in *SubscribeRequest, opts ...grpc.CallOption) (*SubscribeResponse, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(SubscribeResponse)
err := c.cc.Invoke(ctx, LiveSFU_Subscribe_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *liveSFUClient) StopStream(ctx context.Context, in *StopStreamRequest, opts ...grpc.CallOption) (*StopStreamResponse, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(StopStreamResponse)
err := c.cc.Invoke(ctx, LiveSFU_StopStream_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *liveSFUClient) WatchRoom(ctx context.Context, in *WatchRoomRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[RoomEvent], error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
stream, err := c.cc.NewStream(ctx, &LiveSFU_ServiceDesc.Streams[0], LiveSFU_WatchRoom_FullMethodName, cOpts...)
if err != nil {
return nil, err
}
x := &grpc.GenericClientStream[WatchRoomRequest, RoomEvent]{ClientStream: stream}
if err := x.ClientStream.SendMsg(in); err != nil {
return nil, err
}
if err := x.ClientStream.CloseSend(); err != nil {
return nil, err
}
return x, nil
}
// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name.
type LiveSFU_WatchRoomClient = grpc.ServerStreamingClient[RoomEvent]
// LiveSFUServer is the server API for LiveSFU service.
// All implementations should embed UnimplementedLiveSFUServer
// for forward compatibility.
//
// LiveSFU 是直播 SFU 分流 Demo 的控制面契约(protobuf / gRPC)。
//
// 媒体面(WebRTC SDP 交换)不在此契约内:浏览器与 Cloudflare Realtime / SRS
// 直连,经本服务反向代理转发(服务端注入 AppSecret / stream token)。
// 本契约只负责房间、后端、发布/订阅拓扑的协调与扇出状态广播。
type LiveSFUServer interface {
GetConfig(context.Context, *GetConfigRequest) (*GetConfigResponse, error)
ListRooms(context.Context, *ListRoomsRequest) (*ListRoomsResponse, error)
Publish(context.Context, *PublishRequest) (*PublishResponse, error)
Subscribe(context.Context, *SubscribeRequest) (*SubscribeResponse, error)
StopStream(context.Context, *StopStreamRequest) (*StopStreamResponse, error)
WatchRoom(*WatchRoomRequest, grpc.ServerStreamingServer[RoomEvent]) error
}
// UnimplementedLiveSFUServer should be embedded to have
// forward compatible implementations.
//
// NOTE: this should be embedded by value instead of pointer to avoid a nil
// pointer dereference when methods are called.
type UnimplementedLiveSFUServer struct{}
func (UnimplementedLiveSFUServer) GetConfig(context.Context, *GetConfigRequest) (*GetConfigResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method GetConfig not implemented")
}
func (UnimplementedLiveSFUServer) ListRooms(context.Context, *ListRoomsRequest) (*ListRoomsResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method ListRooms not implemented")
}
func (UnimplementedLiveSFUServer) Publish(context.Context, *PublishRequest) (*PublishResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method Publish not implemented")
}
func (UnimplementedLiveSFUServer) Subscribe(context.Context, *SubscribeRequest) (*SubscribeResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method Subscribe not implemented")
}
func (UnimplementedLiveSFUServer) StopStream(context.Context, *StopStreamRequest) (*StopStreamResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method StopStream not implemented")
}
func (UnimplementedLiveSFUServer) WatchRoom(*WatchRoomRequest, grpc.ServerStreamingServer[RoomEvent]) error {
return status.Errorf(codes.Unimplemented, "method WatchRoom not implemented")
}
func (UnimplementedLiveSFUServer) testEmbeddedByValue() {}
// UnsafeLiveSFUServer may be embedded to opt out of forward compatibility for this service.
// Use of this interface is not recommended, as added methods to LiveSFUServer will
// result in compilation errors.
type UnsafeLiveSFUServer interface {
mustEmbedUnimplementedLiveSFUServer()
}
func RegisterLiveSFUServer(s grpc.ServiceRegistrar, srv LiveSFUServer) {
// If the following call pancis, it indicates UnimplementedLiveSFUServer was
// embedded by pointer and is nil. This will cause panics if an
// unimplemented method is ever invoked, so we test this at initialization
// time to prevent it from happening at runtime later due to I/O.
if t, ok := srv.(interface{ testEmbeddedByValue() }); ok {
t.testEmbeddedByValue()
}
s.RegisterService(&LiveSFU_ServiceDesc, srv)
}
func _LiveSFU_GetConfig_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(GetConfigRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(LiveSFUServer).GetConfig(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: LiveSFU_GetConfig_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(LiveSFUServer).GetConfig(ctx, req.(*GetConfigRequest))
}
return interceptor(ctx, in, info, handler)
}
func _LiveSFU_ListRooms_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(ListRoomsRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(LiveSFUServer).ListRooms(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: LiveSFU_ListRooms_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(LiveSFUServer).ListRooms(ctx, req.(*ListRoomsRequest))
}
return interceptor(ctx, in, info, handler)
}
func _LiveSFU_Publish_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(PublishRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(LiveSFUServer).Publish(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: LiveSFU_Publish_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(LiveSFUServer).Publish(ctx, req.(*PublishRequest))
}
return interceptor(ctx, in, info, handler)
}
func _LiveSFU_Subscribe_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(SubscribeRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(LiveSFUServer).Subscribe(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: LiveSFU_Subscribe_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(LiveSFUServer).Subscribe(ctx, req.(*SubscribeRequest))
}
return interceptor(ctx, in, info, handler)
}
func _LiveSFU_StopStream_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(StopStreamRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(LiveSFUServer).StopStream(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: LiveSFU_StopStream_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(LiveSFUServer).StopStream(ctx, req.(*StopStreamRequest))
}
return interceptor(ctx, in, info, handler)
}
func _LiveSFU_WatchRoom_Handler(srv interface{}, stream grpc.ServerStream) error {
m := new(WatchRoomRequest)
if err := stream.RecvMsg(m); err != nil {
return err
}
return srv.(LiveSFUServer).WatchRoom(m, &grpc.GenericServerStream[WatchRoomRequest, RoomEvent]{ServerStream: stream})
}
// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name.
type LiveSFU_WatchRoomServer = grpc.ServerStreamingServer[RoomEvent]
// LiveSFU_ServiceDesc is the grpc.ServiceDesc for LiveSFU service.
// It's only intended for direct use with grpc.RegisterService,
// and not to be introspected or modified (even as a copy)
var LiveSFU_ServiceDesc = grpc.ServiceDesc{
ServiceName: "gospeak.livedemo.v1.LiveSFU",
HandlerType: (*LiveSFUServer)(nil),
Methods: []grpc.MethodDesc{
{
MethodName: "GetConfig",
Handler: _LiveSFU_GetConfig_Handler,
},
{
MethodName: "ListRooms",
Handler: _LiveSFU_ListRooms_Handler,
},
{
MethodName: "Publish",
Handler: _LiveSFU_Publish_Handler,
},
{
MethodName: "Subscribe",
Handler: _LiveSFU_Subscribe_Handler,
},
{
MethodName: "StopStream",
Handler: _LiveSFU_StopStream_Handler,
},
},
Streams: []grpc.StreamDesc{
{
StreamName: "WatchRoom",
Handler: _LiveSFU_WatchRoom_Handler,
ServerStreams: true,
},
},
Metadata: "live_sfu.proto",
}

17
go.mod Normal file
View File

@ -0,0 +1,17 @@
module gospeak-live-sfu-demo
go 1.24.0
toolchain go1.24.5
require (
google.golang.org/grpc v1.80.0
google.golang.org/protobuf v1.36.11
)
require (
golang.org/x/net v0.49.0 // indirect
golang.org/x/sys v0.40.0 // indirect
golang.org/x/text v0.33.0 // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20260120221211-b8f7ae30c516 // indirect
)

38
go.sum Normal file
View File

@ -0,0 +1,38 @@
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI=
github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY=
github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag=
github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE=
github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek=
github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps=
github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=
github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64=
go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y=
go.opentelemetry.io/otel v1.39.0 h1:8yPrr/S0ND9QEfTfdP9V+SiwT4E0G7Y5MO7p85nis48=
go.opentelemetry.io/otel v1.39.0/go.mod h1:kLlFTywNWrFyEdH0oj2xK0bFYZtHRYUdv1NklR/tgc8=
go.opentelemetry.io/otel/metric v1.39.0 h1:d1UzonvEZriVfpNKEVmHXbdf909uGTOQjA0HF0Ls5Q0=
go.opentelemetry.io/otel/metric v1.39.0/go.mod h1:jrZSWL33sD7bBxg1xjrqyDjnuzTUB0x1nBERXd7Ftcs=
go.opentelemetry.io/otel/sdk v1.39.0 h1:nMLYcjVsvdui1B/4FRkwjzoRVsMK8uL/cj0OyhKzt18=
go.opentelemetry.io/otel/sdk v1.39.0/go.mod h1:vDojkC4/jsTJsE+kh+LXYQlbL8CgrEcwmt1ENZszdJE=
go.opentelemetry.io/otel/sdk/metric v1.39.0 h1:cXMVVFVgsIf2YL6QkRF4Urbr/aMInf+2WKg+sEJTtB8=
go.opentelemetry.io/otel/sdk/metric v1.39.0/go.mod h1:xq9HEVH7qeX69/JnwEfp6fVq5wosJsY1mt4lLfYdVew=
go.opentelemetry.io/otel/trace v1.39.0 h1:2d2vfpEDmCJ5zVYz7ijaJdOF59xLomrvj7bjt6/qCJI=
go.opentelemetry.io/otel/trace v1.39.0/go.mod h1:88w4/PnZSazkGzz/w84VHpQafiU4EtqqlVdxWy+rNOA=
golang.org/x/net v0.49.0 h1:eeHFmOGUTtaaPSGNmjBKpbng9MulQsJURQUAfUwY++o=
golang.org/x/net v0.49.0/go.mod h1:/ysNB2EvaqvesRkuLAyjI1ycPZlQHM3q01F02UY/MV8=
golang.org/x/sys v0.40.0 h1:DBZZqJ2Rkml6QMQsZywtnjnnGvHza6BTfYFWY9kjEWQ=
golang.org/x/sys v0.40.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks=
golang.org/x/text v0.33.0 h1:B3njUFyqtHDUI5jMn1YIr5B0IE2U0qck04r6d4KPAxE=
golang.org/x/text v0.33.0/go.mod h1:LuMebE6+rBincTi9+xWTY8TztLzKHc/9C1uBCG27+q8=
gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4=
gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E=
google.golang.org/genproto/googleapis/rpc v0.0.0-20260120221211-b8f7ae30c516 h1:sNrWoksmOyF5bvJUcnmbeAmQi8baNhqg5IWaI3llQqU=
google.golang.org/genproto/googleapis/rpc v0.0.0-20260120221211-b8f7ae30c516/go.mod h1:j9x/tPzZkyxcgEFkiKEEGxfvyumM01BEtsW8xzOahRQ=
google.golang.org/grpc v1.80.0 h1:Xr6m2WmWZLETvUNvIUmeD5OAagMw3FiKmMlTdViWsHM=
google.golang.org/grpc v1.80.0/go.mod h1:ho/dLnxwi3EDJA4Zghp7k2Ec1+c2jqup0bFkw07bwF4=
google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE=
google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=

85
internal/config/config.go Normal file
View File

@ -0,0 +1,85 @@
package config
import (
"os"
"strings"
)
// Config 承载直播 SFU 分流 Demo 的运行参数。
// 优先级:环境变量 > 默认值。Cloudflare Realtime 为主 SFU,SRS 为本地对照后端。
type Config struct {
HTTPPort string
GRPCPort string
ProviderOrder string // 逗号分隔的后端顺序,如 "cloudflare,srs"
SRSBaseURL string
SRSApp string
SRSSecret string
SRSCandidate string
CFAppID string
CFAppSecret string
CFBaseURL string
CFStunURL string
TokenSecret string
TokenRequired bool
}
func Load() *Config {
c := &Config{
HTTPPort: getenv("HTTP_PORT", "8088"),
GRPCPort: getenv("GRPC_PORT", "9090"),
ProviderOrder: getenv("SFU_PROVIDER", "cloudflare,srs"),
SRSBaseURL: getenv("SRS_API_BASE", "http://localhost:1985"),
SRSApp: getenv("SRS_APP", "live"),
SRSSecret: getenv("SRS_SECRET", ""),
SRSCandidate: getenv("SRS_CANDIDATE", "127.0.0.1"),
CFAppID: getenv("CF_APP_ID", ""),
CFAppSecret: getenv("CF_APP_SECRET", ""),
CFBaseURL: getenv("CF_BASE_URL", "https://rtc.live.cloudflare.com/v1"),
CFStunURL: getenv("CF_STUN_URL", "stun:stun.cloudflare.com:3478"),
TokenSecret: getenv("DEMO_TOKEN_SECRET", ""),
TokenRequired: getenv("SFU_TOKEN_REQUIRED", "0") == "1",
}
if c.TokenSecret == "" {
c.TokenSecret = "insecure-demo-secret-change-me"
}
return c
}
func getenv(k, def string) string {
if v := os.Getenv(k); v != "" {
return v
}
return def
}
// ProviderList 返回去重、保序的后端名列表。
func (c *Config) ProviderList() []string {
parts := strings.Split(c.ProviderOrder, ",")
out := make([]string, 0, len(parts))
seen := make(map[string]struct{}, len(parts))
for _, p := range parts {
p = strings.TrimSpace(strings.ToLower(p))
if p == "" || p == "both" || p == "all" {
if p == "both" || p == "all" {
for _, b := range []string{"cloudflare", "srs"} {
if _, ok := seen[b]; !ok {
seen[b] = struct{}{}
out = append(out, b)
}
}
}
continue
}
if _, ok := seen[p]; !ok {
seen[p] = struct{}{}
out = append(out, p)
}
}
return out
}

160
internal/server/gateway.go Normal file
View File

@ -0,0 +1,160 @@
package server
import (
"encoding/json"
"fmt"
"net/http"
"time"
"google.golang.org/protobuf/encoding/protojson"
"google.golang.org/protobuf/proto"
"gospeak-live-sfu-demo/gen"
)
func (s *Server) writeProto(w http.ResponseWriter, m proto.Message, err error) {
if err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
w.Header().Set("Content-Type", "application/json")
data, merr := protojson.Marshal(m)
if merr != nil {
http.Error(w, merr.Error(), http.StatusInternalServerError)
return
}
w.Write(data)
}
func (s *Server) handleConfig(w http.ResponseWriter, r *http.Request) {
resp, err := s.svc.GetConfig(r.Context(), &gen.GetConfigRequest{})
s.writeProto(w, resp, err)
}
func (s *Server) handleRooms(w http.ResponseWriter, r *http.Request) {
resp, err := s.svc.ListRooms(r.Context(), &gen.ListRoomsRequest{})
s.writeProto(w, resp, err)
}
func (s *Server) handlePublish(w http.ResponseWriter, r *http.Request) {
var req gen.PublishRequest
if err := protojson.Unmarshal(readBody(r), &req); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
resp, err := s.svc.Publish(r.Context(), &req)
s.writeProto(w, resp, err)
}
func (s *Server) handleSubscribe(w http.ResponseWriter, r *http.Request) {
var req gen.SubscribeRequest
if err := protojson.Unmarshal(readBody(r), &req); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
resp, err := s.svc.Subscribe(r.Context(), &req)
s.writeProto(w, resp, err)
}
func (s *Server) handleStop(w http.ResponseWriter, r *http.Request) {
var req gen.StopStreamRequest
if err := protojson.Unmarshal(readBody(r), &req); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
resp, err := s.svc.StopStream(r.Context(), &req)
s.writeProto(w, resp, err)
}
// handleRoomEvents 以 SSE 推送房间分发拓扑(分流状态),浏览器据此自动拉流。
func (s *Server) handleRoomEvents(w http.ResponseWriter, r *http.Request) {
room := r.PathValue("room")
flusher, ok := w.(http.Flusher)
if !ok {
http.Error(w, "streaming unsupported", http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
w.Header().Set("Connection", "keep-alive")
ch, unsub := s.hub.subscribe(room)
defer unsub()
sendRoom := func(targets []*gen.StreamTarget) {
data, _ := protojson.Marshal(&gen.RoomEvent{Room: room, Targets: targets})
fmt.Fprintf(w, "event: room\ndata: %s\n\n", data)
flusher.Flush()
}
sendRoom(s.hub.targets(room))
ticker := time.NewTicker(20 * time.Second)
defer ticker.Stop()
for {
select {
case <-r.Context().Done():
return
case msg := <-ch:
sendRoom(decodeTargets(msg))
case <-ticker.C:
sendRoom(s.hub.targets(room))
}
}
}
func (s *Server) handleSRSStreams(w http.ResponseWriter, r *http.Request) {
resp, err := http.Get(s.cfg.SRSBaseURL + "/api/v1/streams/")
if err != nil {
http.Error(w, "srs unreachable: "+err.Error(), http.StatusBadGateway)
return
}
defer resp.Body.Close()
var parsed struct {
Code int `json:"code"`
Streams []struct {
App string `json:"app"`
Name string `json:"name"`
} `json:"streams"`
}
if err := json.NewDecoder(resp.Body).Decode(&parsed); err != nil {
http.Error(w, "srs decode: "+err.Error(), http.StatusBadGateway)
return
}
names := make([]string, 0, len(parsed.Streams))
for _, st := range parsed.Streams {
names = append(names, st.Name)
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]interface{}{"code": parsed.Code, "streams": names})
}
func readBody(r *http.Request) []byte {
buf := make([]byte, 0, 4096)
tmp := make([]byte, 4096)
for {
n, err := r.Body.Read(tmp)
if n > 0 {
buf = append(buf, tmp[:n]...)
}
if err != nil {
break
}
if len(buf) > 1<<20 {
break
}
}
return buf
}
func decodeTargets(msg []byte) []*gen.StreamTarget {
var m struct {
Targets map[string]*gen.StreamTarget `json:"targets"`
}
if err := json.Unmarshal(msg, &m); err != nil {
return nil
}
out := make([]*gen.StreamTarget, 0, len(m.Targets))
for _, t := range m.Targets {
out = append(out, t)
}
return out
}

41
internal/server/grpc.go Normal file
View File

@ -0,0 +1,41 @@
package server
import (
"context"
"gospeak-live-sfu-demo/gen"
)
// grpcServer 把 Service 适配为 protobuf 生成的 gRPC 服务端接口。
type grpcServer struct {
gen.UnimplementedLiveSFUServer
svc *Service
}
func NewGRPCServer(svc *Service) *grpcServer {
return &grpcServer{svc: svc}
}
func (g *grpcServer) GetConfig(ctx context.Context, req *gen.GetConfigRequest) (*gen.GetConfigResponse, error) {
return g.svc.GetConfig(ctx, req)
}
func (g *grpcServer) ListRooms(ctx context.Context, req *gen.ListRoomsRequest) (*gen.ListRoomsResponse, error) {
return g.svc.ListRooms(ctx, req)
}
func (g *grpcServer) Publish(ctx context.Context, req *gen.PublishRequest) (*gen.PublishResponse, error) {
return g.svc.Publish(ctx, req)
}
func (g *grpcServer) Subscribe(ctx context.Context, req *gen.SubscribeRequest) (*gen.SubscribeResponse, error) {
return g.svc.Subscribe(ctx, req)
}
func (g *grpcServer) StopStream(ctx context.Context, req *gen.StopStreamRequest) (*gen.StopStreamResponse, error) {
return g.svc.StopStream(ctx, req)
}
func (g *grpcServer) WatchRoom(req *gen.WatchRoomRequest, stream gen.LiveSFU_WatchRoomServer) error {
return g.svc.WatchRoom(req, stream)
}

View File

@ -0,0 +1,55 @@
package server
import (
"context"
"net"
"testing"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
"gospeak-live-sfu-demo/gen"
"gospeak-live-sfu-demo/internal/config"
"gospeak-live-sfu-demo/internal/sfu/cloudflare"
"gospeak-live-sfu-demo/internal/sfu/srs"
)
// TestGRPCGetConfig 校验 protobuf/gRPC 控制面:生成的服务端 + 生成的客户端能互通。
func TestGRPCGetConfig(t *testing.T) {
cfg := &config.Config{
ProviderOrder: "cloudflare,srs",
SRSBaseURL: "http://localhost:1985",
SRSCandidate: "127.0.0.1",
TokenSecret: "x",
}
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)
gs := grpc.NewServer()
gen.RegisterLiveSFUServer(gs, NewGRPCServer(svc))
lis, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen: %v", err)
}
go func() { _ = gs.Serve(lis) }()
defer gs.Stop()
conn, err := grpc.NewClient(lis.Addr().String(), grpc.WithTransportCredentials(insecure.NewCredentials()))
if err != nil {
t.Fatalf("dial: %v", err)
}
defer conn.Close()
cli := gen.NewLiveSFUClient(conn)
resp, err := cli.GetConfig(context.Background(), &gen.GetConfigRequest{})
if err != nil {
t.Fatalf("GetConfig: %v", err)
}
if len(resp.Backends) == 0 {
t.Fatal("expected backends in config")
}
if resp.Backends[0].Kind != gen.BackendKind_BACKEND_KIND_CLOUDFLARE {
t.Fatalf("cloudflare should be primary, got %v", resp.Backends[0].Kind)
}
}

52
internal/server/proxy.go Normal file
View File

@ -0,0 +1,52 @@
package server
import (
"net/http"
"net/http/httputil"
"net/url"
"strings"
)
// srsProxyHandler 反向代理 SRS HTTP API(WHIP/WHEP 在 /rtc/v1/)。
// 信令可反代,媒体由浏览器直连 SRS :8000。推流可按需校验服务端下发的 JWT。
func (s *Server) srsProxyHandler() http.Handler {
target, err := url.Parse(s.cfg.SRSBaseURL)
if err != nil {
target, _ = url.Parse("http://localhost:1985")
}
rp := httputil.NewSingleHostReverseProxy(target)
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if s.cfg.TokenRequired && strings.HasPrefix(r.URL.Path, "/rtc/v1/whip/") {
room := r.URL.Query().Get("stream")
tok := r.URL.Query().Get("token")
claims, verr := verifyStreamToken(s.cfg.TokenSecret, tok)
if verr != nil || claims.Room != room || claims.Role != "publish" {
http.Error(w, "forbidden: valid publish token required", http.StatusForbidden)
return
}
}
rp.ServeHTTP(w, r)
})
}
// cfProxyHandler 反向代理 Cloudflare Realtime REST,注入 AppSecret。
// 浏览器只交换 SDP(tracks/new),凭证不出服务端。
func (s *Server) cfProxyHandler() http.Handler {
base := strings.TrimRight(s.cfg.CFBaseURL, "/") + "/apps/" + s.cfg.CFAppID
target, err := url.Parse(base)
if err != nil {
target, _ = url.Parse("https://rtc.live.cloudflare.com/v1/apps/" + s.cfg.CFAppID)
}
rp := httputil.NewSingleHostReverseProxy(target)
rp.Director = func(r *http.Request) {
rest := strings.TrimPrefix(r.URL.Path, "/api/cf")
r.URL.Path = target.Path + rest
r.URL.Scheme = target.Scheme
r.URL.Host = target.Host
r.Host = target.Host
if s.cfg.CFAppSecret != "" {
r.Header.Set("Authorization", "Bearer "+s.cfg.CFAppSecret)
}
}
return rp
}

112
internal/server/rooms.go Normal file
View File

@ -0,0 +1,112 @@
package server
import (
"encoding/json"
"sync"
"gospeak-live-sfu-demo/gen"
)
// roomHub 维护房间 -> 各后端分发目标(StreamTarget)的内存状态,
// 并向订阅者广播拓扑变化(实现「分流」状态实时可见)。
type roomHub struct {
mu sync.RWMutex
rooms map[string]*roomEntry
}
type roomEntry struct {
targets map[string]*gen.StreamTarget // key: 后端名(cloudflare / srs)
subs map[chan []byte]struct{}
}
func newRoomHub() *roomHub {
return &roomHub{rooms: map[string]*roomEntry{}}
}
func (h *roomHub) get(name string) *roomEntry {
e, ok := h.rooms[name]
if !ok {
e = &roomEntry{targets: map[string]*gen.StreamTarget{}, subs: map[chan []byte]struct{}{}}
h.rooms[name] = e
}
return e
}
func (h *roomHub) setTarget(room, backend string, t *gen.StreamTarget) {
h.mu.Lock()
defer h.mu.Unlock()
e := h.get(room)
e.targets[backend] = t
h.broadcast(room, e)
}
func (h *roomHub) removeTarget(room, backend string) {
h.mu.Lock()
defer h.mu.Unlock()
e, ok := h.rooms[room]
if !ok {
return
}
delete(e.targets, backend)
if len(e.targets) == 0 && len(e.subs) == 0 {
delete(h.rooms, room)
return
}
h.broadcast(room, e)
}
func (h *roomHub) targets(room string) []*gen.StreamTarget {
h.mu.RLock()
defer h.mu.RUnlock()
e, ok := h.rooms[room]
if !ok {
return nil
}
out := make([]*gen.StreamTarget, 0, len(e.targets))
for _, t := range e.targets {
out = append(out, t)
}
return out
}
func (h *roomHub) list() []*gen.Room {
h.mu.RLock()
defer h.mu.RUnlock()
out := make([]*gen.Room, 0, len(h.rooms))
for name, e := range h.rooms {
room := &gen.Room{Name: name}
for _, t := range e.targets {
room.Targets = append(room.Targets, t)
}
out = append(out, room)
}
return out
}
func (h *roomHub) subscribe(room string) (chan []byte, func()) {
h.mu.Lock()
defer h.mu.Unlock()
e := h.get(room)
ch := make(chan []byte, 8)
e.subs[ch] = struct{}{}
return ch, func() {
h.mu.Lock()
defer h.mu.Unlock()
if e2, ok := h.rooms[room]; ok {
delete(e2.subs, ch)
}
}
}
func (h *roomHub) broadcast(room string, e *roomEntry) {
payload, _ := json.Marshal(map[string]interface{}{
"room": room,
"targets": e.targets,
})
for ch := range e.subs {
select {
case ch <- payload:
default:
}
}
}

91
internal/server/server.go Normal file
View File

@ -0,0 +1,91 @@
package server
import (
"bytes"
"embed"
"io/fs"
"net"
"net/http"
"time"
"google.golang.org/grpc"
"gospeak-live-sfu-demo/gen"
"gospeak-live-sfu-demo/internal/config"
"gospeak-live-sfu-demo/internal/sfu/cloudflare"
"gospeak-live-sfu-demo/internal/sfu/srs"
)
//go:embed static
var staticFS embed.FS
// Server 聚合控制面(gRPC + JSON 网关)、媒体面反向代理与静态 UI。
type Server struct {
cfg *config.Config
svc *Service
hub *roomHub
srsProxy http.Handler
cfProxy http.Handler
}
func New(cfg *config.Config) *Server {
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}
s.srsProxy = s.srsProxyHandler()
s.cfProxy = s.cfProxyHandler()
return s
}
// StartGRPC 启动 gRPC 控制面(protobuf 契约的服务端实现)。
func (s *Server) StartGRPC() error {
lis, err := net.Listen("tcp", ":"+s.cfg.GRPCPort)
if err != nil {
return err
}
gs := grpc.NewServer()
gen.RegisterLiveSFUServer(gs, NewGRPCServer(s.svc))
go func() { _ = gs.Serve(lis) }()
return nil
}
// Handler 返回 HTTP 路由:静态 UI、JSON 网关、SSE、SRS/Cloudflare 媒体反代。
func (s *Server) Handler() http.Handler {
mux := http.NewServeMux()
sub, _ := fs.Sub(staticFS, "static")
mux.Handle("GET /", s.fileServer("static/index.html"))
mux.Handle("GET /publish", s.fileServer("static/publish.html"))
mux.Handle("GET /watch", s.fileServer("static/watch.html"))
mux.Handle("GET /static/", http.StripPrefix("/static/", http.FileServer(http.FS(sub))))
mux.HandleFunc("GET /api/config", s.handleConfig)
mux.HandleFunc("GET /api/rooms", s.handleRooms)
mux.HandleFunc("POST /api/publish", s.handlePublish)
mux.HandleFunc("POST /api/subscribe", s.handleSubscribe)
mux.HandleFunc("POST /api/stop", s.handleStop)
mux.HandleFunc("GET /api/srs/streams", s.handleSRSStreams)
mux.Handle("GET /api/room/{room}/events", http.HandlerFunc(s.handleRoomEvents))
mux.Handle("GET /rtc/v1/", s.srsProxy)
mux.Handle("POST /rtc/v1/", s.srsProxy)
mux.Handle("PUT /rtc/v1/", s.srsProxy)
mux.Handle("DELETE /rtc/v1/", s.srsProxy)
mux.Handle("GET /api/cf/", s.cfProxy)
mux.Handle("POST /api/cf/", s.cfProxy)
mux.Handle("PUT /api/cf/", s.cfProxy)
mux.Handle("DELETE /api/cf/", s.cfProxy)
return mux
}
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))
}
}

217
internal/server/service.go Normal file
View File

@ -0,0 +1,217 @@
package server
import (
"context"
"fmt"
"time"
"gospeak-live-sfu-demo/gen"
"gospeak-live-sfu-demo/internal/config"
"gospeak-live-sfu-demo/internal/sfu/cloudflare"
"gospeak-live-sfu-demo/internal/sfu/srs"
)
// Service 实现 LiveSFU 控制面逻辑:协调房间在各 SFU 后端的发布/订阅拓扑。
// 媒体面(SDP 交换)由浏览器经反向代理直连 Cloudflare / SRS,本服务只持有凭证。
type Service struct {
cfg *config.Config
cf *cloudflare.Provider
srsP *srs.Provider
hub *roomHub
}
func NewService(cfg *config.Config, cf *cloudflare.Provider, srsP *srs.Provider, hub *roomHub) *Service {
return &Service{cfg: cfg, cf: cf, srsP: srsP, hub: hub}
}
func backendName(k gen.BackendKind) string {
switch k {
case gen.BackendKind_BACKEND_KIND_CLOUDFLARE:
return "cloudflare"
case gen.BackendKind_BACKEND_KIND_SRS:
return "srs"
default:
return ""
}
}
// GetConfig 返回后端能力与全局拓扑信息(浏览器据此渲染分发面板)。
func (s *Service) GetConfig(ctx context.Context, req *gen.GetConfigRequest) (*gen.GetConfigResponse, error) {
resp := &gen.GetConfigResponse{
Candidate: s.srsP.Candidate(),
TokenRequired: s.cfg.TokenRequired,
}
for _, name := range s.cfg.ProviderList() {
switch name {
case "cloudflare":
resp.Backends = append(resp.Backends, s.cf.BackendInfo())
case "srs":
resp.Backends = append(resp.Backends, s.srsP.BackendInfo())
}
}
return resp, nil
}
func (s *Service) ListRooms(ctx context.Context, req *gen.ListRoomsRequest) (*gen.ListRoomsResponse, error) {
return &gen.ListRoomsResponse{Rooms: s.hub.list()}, nil
}
// Publish 开始向某后端发布:Cloudflare 创建 session;SRS 分配 stream + 签发 JWT。
func (s *Service) Publish(ctx context.Context, req *gen.PublishRequest) (*gen.PublishResponse, error) {
room := req.GetRoom()
if room == "" {
return nil, fmt.Errorf("room required")
}
backend := backendName(req.GetBackend())
if backend == "" {
return nil, fmt.Errorf("unknown backend")
}
identity := req.GetIdentity()
if identity == "" {
identity = randomID()
}
switch backend {
case "cloudflare":
if !s.cf.Configured() {
return nil, fmt.Errorf("cloudflare realtime not configured (set CF_APP_ID / CF_APP_SECRET)")
}
sessionID, err := s.cf.Client().CreateSession(room)
if err != nil {
return nil, fmt.Errorf("cf create session: %w", err)
}
info, err := s.cf.Client().GetSession(sessionID)
if err != nil {
return nil, fmt.Errorf("cf get session: %w", err)
}
target := &gen.StreamTarget{
Backend: gen.BackendKind_BACKEND_KIND_CLOUDFLARE,
SessionId: sessionID,
PublishedAt: time.Now().Unix(),
}
s.hub.setTarget(room, "cloudflare", target)
return &gen.PublishResponse{
SessionId: sessionID,
IceServers: iceServersToProto(info.IceServers),
Target: target,
}, nil
case "srs":
stream := "live-" + room
token, _ := signStreamToken(s.cfg.TokenSecret, stream, identity, "publish", 2*time.Hour)
target := &gen.StreamTarget{
Backend: gen.BackendKind_BACKEND_KIND_SRS,
Stream: stream,
PublishToken: token,
Url: fmt.Sprintf("/rtc/v1/whep/?app=%s&stream=%s", s.srsP.App(), stream),
PublishedAt: time.Now().Unix(),
}
s.hub.setTarget(room, "srs", target)
return &gen.PublishResponse{
Stream: stream,
PublishToken: token,
IceServers: []*gen.IceServer{{Urls: []string{s.cf.Stun()}}},
Target: target,
}, nil
}
return nil, fmt.Errorf("unsupported backend")
}
// Subscribe 订阅某房间在某后端的分发目标,返回建立 WebRTC 所需的 session / stream。
func (s *Service) Subscribe(ctx context.Context, req *gen.SubscribeRequest) (*gen.SubscribeResponse, error) {
room := req.GetRoom()
if room == "" {
return nil, fmt.Errorf("room required")
}
backend := backendName(req.GetBackend())
if backend == "" {
return nil, fmt.Errorf("unknown backend")
}
pub := s.findTarget(room, req.GetBackend())
if pub == nil {
return nil, fmt.Errorf("room %q not live on %s", room, backend)
}
switch backend {
case "cloudflare":
viewerSession, err := s.cf.Client().CreateSession(room)
if err != nil {
return nil, fmt.Errorf("cf create viewer session: %w", err)
}
info, err := s.cf.Client().GetSession(viewerSession)
if err != nil {
return nil, fmt.Errorf("cf get viewer session: %w", err)
}
return &gen.SubscribeResponse{
SessionId: viewerSession,
PublisherSessionId: pub.GetSessionId(),
IceServers: iceServersToProto(info.IceServers),
}, nil
case "srs":
return &gen.SubscribeResponse{
Stream: pub.GetStream(),
IceServers: []*gen.IceServer{{Urls: []string{s.cf.Stun()}}},
}, nil
}
return nil, fmt.Errorf("unsupported backend")
}
func (s *Service) StopStream(ctx context.Context, req *gen.StopStreamRequest) (*gen.StopStreamResponse, error) {
room := req.GetRoom()
backend := backendName(req.GetBackend())
if room == "" || backend == "" {
return nil, fmt.Errorf("room and backend required")
}
if backend == "cloudflare" {
for _, t := range s.hub.targets(room) {
if t.GetBackend() == gen.BackendKind_BACKEND_KIND_CLOUDFLARE && t.GetSessionId() != "" {
_ = s.cf.Client().DeleteSession(t.GetSessionId())
}
}
}
s.hub.removeTarget(room, backend)
return &gen.StopStreamResponse{Ok: true}, nil
}
// WatchRoom 服务端流式推送房间分发拓扑(分流)变化。
func (s *Service) WatchRoom(req *gen.WatchRoomRequest, stream gen.LiveSFU_WatchRoomServer) error {
room := req.GetRoom()
ch, unsub := s.hub.subscribe(room)
defer unsub()
if err := stream.Send(&gen.RoomEvent{Room: room, Targets: s.hub.targets(room)}); err != nil {
return err
}
ticker := time.NewTicker(20 * time.Second)
defer ticker.Stop()
for {
select {
case <-stream.Context().Done():
return nil
case msg := <-ch:
if err := stream.Send(&gen.RoomEvent{Room: room, Targets: decodeTargets(msg)}); err != nil {
return err
}
case <-ticker.C:
if err := stream.Send(&gen.RoomEvent{Room: room, Targets: s.hub.targets(room)}); err != nil {
return err
}
}
}
}
func (s *Service) findTarget(room string, kind gen.BackendKind) *gen.StreamTarget {
for _, t := range s.hub.targets(room) {
if t.GetBackend() == kind {
return t
}
}
return nil
}
func iceServersToProto(in []cloudflare.IceServer) []*gen.IceServer {
out := make([]*gen.IceServer, 0, len(in))
for _, s := range in {
out = append(out, &gen.IceServer{Urls: s.URLs, Username: s.Username, Credential: s.Credential})
}
return out
}

View File

@ -0,0 +1,230 @@
// 直播 SFU 分流 Demo 前端逻辑。控制面走 JSON 网关(protojson),媒体面走反向代理直连 SFU。
window.LiveSFU = (function () {
const ENUM = { cloudflare: "BACKEND_KIND_CLOUDFLARE", srs: "BACKEND_KIND_SRS" };
const NAME = { BACKEND_KIND_CLOUDFLARE: "Cloudflare Realtime", BACKEND_KIND_SRS: "SRS" };
function getRoom() {
const p = new URLSearchParams(location.search).get("room");
return (p && p.trim()) || "demo";
}
function shortName(enumStr) { return NAME[enumStr] || enumStr; }
async function apiGet(path) {
const r = await fetch(path);
if (!r.ok) throw new Error(path + " -> " + r.status);
return r.json();
}
async function apiPost(path, body) {
const r = await fetch(path, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify(body),
});
const text = await r.text();
if (!r.ok) throw new Error((text || r.status));
return text ? JSON.parse(text) : {};
}
async function loadConfig(el) {
try {
const cfg = await apiGet("/api/config");
el.innerHTML = "";
(cfg.backends || []).forEach((b) => {
const span = document.createElement("span");
span.className = "badge" + (b.primary ? " primary" : "");
span.textContent = shortName(b.kind) + (b.primary ? " · 主" : "") + (b.configured ? " · 就绪" : " · 未配置");
el.appendChild(span);
});
} catch (e) {
el.innerHTML = '<span class="badge bad">控制面不可达</span>';
}
}
function log(el, msg) {
const t = new Date().toLocaleTimeString();
el.textContent += "[" + t + "] " + msg + "\n";
el.scrollTop = el.scrollHeight;
}
// ---------- WebRTC:Cloudflare Realtime ----------
async function publishCF(sessionId, localStream, iceServers) {
const pc = new RTCPeerConnection({ iceServers });
localStream.getTracks().forEach((t) => pc.addTrack(t, localStream));
const offer = await pc.createOffer();
await pc.setLocalDescription(offer);
const resp = await fetch("/api/cf/sessions/" + sessionId + "/tracks/new", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({
sessionDescription: { type: "offer", sdp: offer.sdp },
tracks: localStream.getTracks().map((t) => ({ location: "local", trackName: t.kind, kind: t.kind })),
autoDiscover: true,
}),
});
const data = await resp.json();
await pc.setRemoteDescription({ type: data.sessionDescription.type, sdp: data.sessionDescription.sdp });
return pc;
}
async function watchCF(viewerSessionId, publisherSessionId, iceServers, video) {
const pc = new RTCPeerConnection({ iceServers });
pc.ontrack = (e) => { video.srcObject = e.streams[0]; };
const offer = await pc.createOffer();
await pc.setLocalDescription(offer);
const resp = await fetch("/api/cf/sessions/" + viewerSessionId + "/tracks/new", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({
sessionDescription: { type: "offer", sdp: offer.sdp },
tracks: [
{ location: "remote", sessionId: publisherSessionId, trackName: "audio" },
{ location: "remote", sessionId: publisherSessionId, trackName: "video" },
],
autoDiscover: true,
}),
});
const data = await resp.json();
await pc.setRemoteDescription({ type: data.sessionDescription.type, sdp: data.sessionDescription.sdp });
return pc;
}
// ---------- WebRTC:SRS (WHIP/WHEP) ----------
async function publishSRS(room, token, stream, iceServers, localStream) {
const pc = new RTCPeerConnection({ iceServers });
localStream.getTracks().forEach((t) => pc.addTrack(t, localStream));
const offer = await pc.createOffer();
await pc.setLocalDescription(offer);
const url = "/rtc/v1/whip/?app=live&stream=" + encodeURIComponent(stream) + "&token=" + encodeURIComponent(token);
const resp = await fetch(url, { method: "POST", headers: { "Content-Type": "application/sdp" }, body: offer.sdp });
const answer = await resp.text();
await pc.setRemoteDescription({ type: "answer", sdp: answer });
return pc;
}
async function watchSRS(room, stream, iceServers, video) {
const pc = new RTCPeerConnection({ iceServers });
pc.ontrack = (e) => { video.srcObject = e.streams[0]; };
const offer = await pc.createOffer();
await pc.setLocalDescription(offer);
const url = "/rtc/v1/whep/?app=live&stream=" + encodeURIComponent(stream);
const resp = await fetch(url, { method: "POST", headers: { "Content-Type": "application/sdp" }, body: offer.sdp });
const answer = await resp.text();
await pc.setRemoteDescription({ type: "answer", sdp: answer });
return pc;
}
function findTarget(targets, enumStr) {
return (targets || []).find((t) => t.backend === enumStr);
}
// ---------- 发布页 ----------
function initPublish() {
const room = getRoom();
document.getElementById("roomName").textContent = room;
const localVideo = document.getElementById("local");
const logEl = document.getElementById("log");
const statusEl = document.getElementById("status");
const pcs = {};
document.getElementById("start").onclick = async () => {
try {
const backends = Array.from(document.querySelectorAll('input[name="backend"]:checked')).map((i) => i.value);
if (!backends.length) return log(logEl, "请至少选择一个分发后端");
const stream = await navigator.mediaDevices.getUserMedia({ video: true, audio: true });
localVideo.srcObject = stream;
for (const b of backends) {
try {
const resp = await apiPost("/api/publish", { room, backend: ENUM[b] });
if (b === "cloudflare") {
pcs[b] = await publishCF(resp.sessionId, stream, resp.iceServers);
} else {
pcs[b] = await publishSRS(room, resp.publishToken, resp.stream, resp.iceServers, stream);
}
log(logEl, "已在 " + shortName(resp.target.backend) + " 发布(" + (resp.stream || resp.sessionId) + ")");
} catch (e) {
log(logEl, "发布到 " + b + " 失败:" + e.message);
}
}
renderStatus(statusEl, pcs);
} catch (e) {
log(logEl, "获取摄像头失败:" + e.message);
}
};
document.getElementById("stop").onclick = async () => {
for (const b of Object.keys(pcs)) {
try { await apiPost("/api/stop", { room, backend: ENUM[b] }); } catch (e) {}
if (pcs[b]) pcs[b].close();
delete pcs[b];
}
if (localVideo.srcObject) localVideo.srcObject.getTracks().forEach((t) => t.stop());
log(logEl, "已停止推流");
renderStatus(statusEl, pcs);
};
}
function renderStatus(el, pcs) {
const keys = Object.keys(pcs);
el.innerHTML = keys.length
? keys.map((k) => '<span class="target ' + k + '">分发中:' + shortName(ENUM[k]) + "</span>").join("")
: '<span class="muted">未发布</span>';
}
// ---------- 观看页 ----------
function initWatch() {
const room = getRoom();
document.getElementById("roomName").textContent = room;
const remoteVideo = document.getElementById("remote");
const logEl = document.getElementById("log");
const statusEl = document.getElementById("status");
let pc = null;
let es = null;
let watching = false;
function backend() { return document.querySelector('input[name="watchBackend"]:checked').value; }
async function subscribe() {
if (watching) return;
const b = backend();
try {
const resp = await apiPost("/api/subscribe", { room, backend: ENUM[b] });
if (b === "cloudflare") {
pc = await watchCF(resp.sessionId, resp.publisherSessionId, resp.iceServers, remoteVideo);
} else {
pc = await watchSRS(room, resp.stream, resp.iceServers, remoteVideo);
}
watching = true;
log(logEl, "已从 " + shortName(ENUM[b]) + " 拉流");
statusEl.innerHTML = '<span class="target ' + b + '">观看中:' + shortName(ENUM[b]) + "</span>";
} catch (e) {
log(logEl, "订阅失败:" + e.message);
}
}
document.getElementById("start").onclick = () => {
es = new EventSource("/api/room/" + encodeURIComponent(room) + "/events");
es.addEventListener("room", (ev) => {
const data = JSON.parse(ev.data);
const t = findTarget(data.targets, ENUM[backend()]);
statusEl.innerHTML = t
? '<span class="target ' + backend() + '">房间在 ' + shortName(ENUM[backend()]) + " 已直播</span>"
: '<span class="muted">等待 ' + shortName(ENUM[backend()]) + " 推流…</span>";
if (t && !watching) subscribe();
});
log(logEl, "已连接房间事件流");
};
document.getElementById("stop").onclick = () => {
if (pc) pc.close();
pc = null;
if (es) es.close();
es = null;
watching = false;
remoteVideo.srcObject = null;
statusEl.innerHTML = '<span class="muted">未观看</span>';
log(logEl, "已停止观看");
};
}
return { ENUM, getRoom, loadConfig, log, initPublish, initWatch, apiGet, apiPost };
})();

View File

@ -0,0 +1,54 @@
<!doctype html>
<html lang="zh-CN">
<head>
<meta charset="utf-8" />
<meta name="viewport" content="width=device-width, initial-scale=1" />
<title>直播 SFU 分流 Demo</title>
<link rel="stylesheet" href="/static/styles.css" />
</head>
<body>
<div class="topbar">
<div class="brand">直播 SFU 分流 Demo <small>Cloudflare Realtime · SRS</small></div>
<nav>
<a href="/publish">开始直播</a>
<a href="/watch">观看直播</a>
</nav>
</div>
<div class="band">
<div class="hero">
<h1>一路推流,多路扇出</h1>
<p>基于 GOSpeak 技术栈的直播 SFU 分流演示。控制面用 protobuf/gRPC 定义,主 SFU 为
<b>Cloudflare Realtime</b>,<b>SRS</b> 作为本地对照后端。同一路发布可同时落到多条分发链路,观众任选后端拉流。</p>
</div>
<div class="grid">
<div class="card">
<h3>进入房间</h3>
<div class="row">
<input type="text" id="room" placeholder="房间名,例如 demo" />
<button onclick="goPublish()">开始直播</button>
<button class="secondary" onclick="goWatch()">观看直播</button>
</div>
<div class="muted" style="margin-top:10px">提示:推流与观看使用同一房间名即可配对。</div>
</div>
<div class="card">
<h3>后端能力</h3>
<div id="backends" class="targets"><span class="muted">加载中…</span></div>
<div class="muted" style="margin-top:10px">Cloudflare Realtime 需在 <code>.env</code> 配置 <code>CF_APP_ID</code> / <code>CF_APP_SECRET</code>;
SRS 由本地 <code>docker compose up -d srs</code> 提供。</div>
</div>
</div>
</div>
<script src="/static/app.js"></script>
<script>
LiveSFU.loadConfig(document.getElementById('backends'));
function goPublish() {
const r = document.getElementById('room').value.trim() || 'demo';
location.href = '/publish?room=' + encodeURIComponent(r);
}
function goWatch() {
const r = document.getElementById('room').value.trim() || 'demo';
location.href = '/watch?room=' + encodeURIComponent(r);
}
</script>
</body>
</html>

View File

@ -0,0 +1,43 @@
<!doctype html>
<html lang="zh-CN">
<head>
<meta charset="utf-8" />
<meta name="viewport" content="width=device-width, initial-scale=1" />
<title>开始直播 · 直播 SFU 分流 Demo</title>
<link rel="stylesheet" href="/static/styles.css" />
</head>
<body>
<div class="topbar">
<div class="brand">直播 SFU 分流 Demo <small>Cloudflare Realtime · SRS</small></div>
<nav><a href="/">首页</a><a href="/watch">观看直播</a></nav>
</div>
<div class="band">
<div class="hero"><h1>开始直播</h1><p>房间 <b id="roomName">demo</b> · 选择分发后端,一路推流将扇出到所选 SFU。</p></div>
<div class="two">
<div class="card">
<h3>本地预览</h3>
<video id="local" autoplay playsinline muted></video>
<div class="check">
<label><input type="checkbox" name="backend" value="cloudflare" checked /> Cloudflare Realtime(主 SFU)</label>
</div>
<div class="check">
<label><input type="checkbox" name="backend" value="srs" checked /> SRS(本地对照)</label>
</div>
<div class="row" style="margin-top:10px">
<button id="start">开始推流</button>
<button id="stop" class="danger">停止推流</button>
</div>
<div id="status" style="margin-top:10px"><span class="muted">未发布</span></div>
</div>
<div class="card">
<h3>分发状态</h3>
<div id="status2" class="muted">尚未开始</div>
<h3 style="margin-top:14px">日志</h3>
<div id="log" class="log"></div>
</div>
</div>
</div>
<script src="/static/app.js"></script>
<script>LiveSFU.initPublish();</script>
</body>
</html>

View File

@ -0,0 +1,79 @@
:root {
--bg: #0b0f17;
--panel: #121826;
--panel-2: #18202f;
--border: #243044;
--text: #e6edf6;
--muted: #93a1b5;
--accent: #38bdf8;
--accent-2: #f59e0b;
--ok: #34d399;
--bad: #f87171;
}
* { box-sizing: border-box; }
body {
margin: 0;
font-family: ui-sans-serif, system-ui, -apple-system, "Segoe UI", Roboto, "PingFang SC", "Microsoft YaHei", sans-serif;
background: var(--bg);
color: var(--text);
line-height: 1.5;
}
a { color: var(--accent); text-decoration: none; }
.topbar {
display: flex; align-items: center; justify-content: space-between;
padding: 14px 22px; border-bottom: 1px solid var(--border);
background: linear-gradient(180deg, #0e1422, #0b0f17);
}
.topbar .brand { font-weight: 700; letter-spacing: .3px; }
.topbar .brand small { color: var(--muted); font-weight: 500; margin-left: 8px; }
.topbar nav a { margin-left: 18px; color: var(--muted); font-size: 14px; }
.topbar nav a:hover { color: var(--text); }
.band { max-width: 1100px; margin: 0 auto; padding: 28px 22px; }
.hero h1 { font-size: 30px; margin: 0 0 8px; }
.hero p { color: var(--muted); max-width: 720px; margin: 0 0 18px; }
.grid { display: grid; gap: 16px; }
.row { display: flex; gap: 12px; flex-wrap: wrap; align-items: center; }
.card {
background: var(--panel); border: 1px solid var(--border);
border-radius: 10px; padding: 16px 18px;
}
.card h3 { margin: 0 0 10px; font-size: 15px; }
input[type="text"] {
background: var(--panel-2); border: 1px solid var(--border); color: var(--text);
padding: 10px 12px; border-radius: 8px; font-size: 15px; min-width: 240px;
}
button {
background: var(--accent); color: #04121c; border: 0; border-radius: 8px;
padding: 10px 16px; font-size: 14px; font-weight: 600; cursor: pointer;
}
button.secondary { background: var(--panel-2); color: var(--text); border: 1px solid var(--border); }
button.danger { background: var(--bad); color: #1a0606; }
button:disabled { opacity: .5; cursor: not-allowed; }
button:hover:not(:disabled) { filter: brightness(1.08); }
.check { display: flex; gap: 10px; align-items: center; padding: 8px 0; }
.check label { display: flex; gap: 8px; align-items: center; cursor: pointer; }
.badge { font-size: 12px; padding: 2px 8px; border-radius: 999px; border: 1px solid var(--border); color: var(--muted); }
.badge.primary { color: var(--accent-2); border-color: #4a3a16; }
.badge.ok { color: var(--ok); border-color: #14402f; }
.badge.bad { color: var(--bad); border-color: #4a1c1c; }
video {
width: 100%; background: #000; border-radius: 10px; border: 1px solid var(--border);
aspect-ratio: 16 / 9; object-fit: contain;
}
.log {
font-family: ui-monospace, SFMono-Regular, Menlo, monospace; font-size: 12.5px;
background: #0a0e16; border: 1px solid var(--border); border-radius: 8px;
padding: 10px 12px; height: 180px; overflow: auto; color: var(--muted); white-space: pre-wrap;
}
.kv { display: flex; justify-content: space-between; padding: 4px 0; border-bottom: 1px dashed var(--border); font-size: 13px; }
.kv:last-child { border-bottom: 0; }
.targets { display: flex; gap: 10px; flex-wrap: wrap; }
.target {
border: 1px solid var(--border); border-radius: 8px; padding: 8px 12px; font-size: 13px;
background: var(--panel-2);
}
.target.cf { border-color: #4a3a16; }
.target.srs { border-color: #1b3550; }
.muted { color: var(--muted); font-size: 13px; }
.two { display: grid; grid-template-columns: 1fr 1fr; gap: 16px; }
@media (max-width: 820px) { .two { grid-template-columns: 1fr; } }

View File

@ -0,0 +1,41 @@
<!doctype html>
<html lang="zh-CN">
<head>
<meta charset="utf-8" />
<meta name="viewport" content="width=device-width, initial-scale=1" />
<title>观看直播 · 直播 SFU 分流 Demo</title>
<link rel="stylesheet" href="/static/styles.css" />
</head>
<body>
<div class="topbar">
<div class="brand">直播 SFU 分流 Demo <small>Cloudflare Realtime · SRS</small></div>
<nav><a href="/">首页</a><a href="/publish">开始直播</a></nav>
</div>
<div class="band">
<div class="hero"><h1>观看直播</h1><p>房间 <b id="roomName">demo</b> · 选择从哪个 SFU 拉流,房间一旦直播即自动连接。</p></div>
<div class="two">
<div class="card">
<h3>远端画面</h3>
<video id="remote" autoplay playsinline></video>
<div class="check">
<label><input type="radio" name="watchBackend" value="cloudflare" checked /> Cloudflare Realtime(主 SFU)</label>
</div>
<div class="check">
<label><input type="radio" name="watchBackend" value="srs" /> SRS(本地对照)</label>
</div>
<div class="row" style="margin-top:10px">
<button id="start">开始观看</button>
<button id="stop" class="danger">停止观看</button>
</div>
<div id="status" style="margin-top:10px"><span class="muted">未观看</span></div>
</div>
<div class="card">
<h3>日志</h3>
<div id="log" class="log"></div>
</div>
</div>
</div>
<script src="/static/app.js"></script>
<script>LiveSFU.initWatch();</script>
</body>
</html>

65
internal/server/token.go Normal file
View File

@ -0,0 +1,65 @@
package server
import (
"crypto/hmac"
"crypto/rand"
"crypto/sha256"
"encoding/base64"
"encoding/json"
"errors"
"strings"
"time"
)
// SRS 推流 JWT(HS256)。服务端签发、并在 WHIP 反代层校验 role=publish,
// 形成「分发入口」的访问控制。Cloudflare 走自身 AppSecret,无需此 token。
func b64url(b []byte) string { return base64.RawURLEncoding.EncodeToString(b) }
type streamClaims struct {
Room string `json:"room"`
Identity string `json:"identity"`
Role string `json:"role"`
Exp int64 `json:"exp"`
Iat int64 `json:"iat"`
}
func signStreamToken(secret, room, identity, role string, ttl time.Duration) (string, error) {
now := time.Now()
c := streamClaims{Room: room, Identity: identity, Role: role, Iat: now.Unix(), Exp: now.Add(ttl).Unix()}
payload, _ := json.Marshal(c)
header, _ := json.Marshal(map[string]string{"alg": "HS256", "typ": "JWT"})
signing := b64url(header) + "." + b64url(payload)
mac := hmac.New(sha256.New, []byte(secret))
mac.Write([]byte(signing))
return signing + "." + b64url(mac.Sum(nil)), nil
}
func verifyStreamToken(secret, token string) (*streamClaims, error) {
parts := strings.Split(token, ".")
if len(parts) != 3 {
return nil, errors.New("bad token")
}
mac := hmac.New(sha256.New, []byte(secret))
mac.Write([]byte(parts[0] + "." + parts[1]))
if !hmac.Equal([]byte(b64url(mac.Sum(nil))), []byte(parts[2])) {
return nil, errors.New("bad signature")
}
p, err := base64.RawURLEncoding.DecodeString(parts[1])
if err != nil {
return nil, err
}
var c streamClaims
if err := json.Unmarshal(p, &c); err != nil {
return nil, err
}
if c.Exp < time.Now().Unix() {
return nil, errors.New("token expired")
}
return &c, nil
}
func randomID() string {
b := make([]byte, 9)
rand.Read(b)
return base64.RawURLEncoding.EncodeToString(b)
}

View File

@ -0,0 +1,40 @@
package server
import (
"testing"
"time"
)
func TestStreamTokenRoundtrip(t *testing.T) {
secret := "s3cr3t"
tok, err := signStreamToken(secret, "live-demo", "id1", "publish", time.Hour)
if err != nil {
t.Fatalf("sign: %v", err)
}
claims, err := verifyStreamToken(secret, tok)
if err != nil {
t.Fatalf("verify: %v", err)
}
if claims.Room != "live-demo" || claims.Role != "publish" {
t.Fatalf("claims = %+v", claims)
}
}
func TestStreamTokenRejects(t *testing.T) {
secret := "s3cr3t"
tok, _ := signStreamToken(secret, "live-demo", "id1", "publish", time.Hour)
if _, err := verifyStreamToken("other", tok); err == nil {
t.Fatal("expected reject on wrong secret")
}
if _, err := verifyStreamToken(secret, tok+"x"); err == nil {
t.Fatal("expected reject on tampered token")
}
exp, _ := signStreamToken(secret, "r", "i", "publish", -time.Hour)
if _, err := verifyStreamToken(secret, exp); err == nil {
t.Fatal("expected reject on expired token")
}
if _, err := verifyStreamToken(secret, "not.a.jwt"); err == nil {
t.Fatal("expected reject on malformed token")
}
}

View File

@ -0,0 +1,161 @@
package cloudflare
import (
"bytes"
"encoding/json"
"fmt"
"io"
"net/http"
"net/url"
"time"
)
const defaultBaseURL = "https://rtc.live.cloudflare.com/v1"
// IceServer 对应 Cloudflare Realtime session 返回的 ICE 配置(可能含 TURN 凭证)。
type IceServer struct {
URLs []string `json:"urls"`
Username string `json:"username,omitempty"`
Credential string `json:"credential,omitempty"`
}
type SessionInfo struct {
SessionID string `json:"sessionId"`
AppID string `json:"appId"`
RequesterIP string `json:"requesterIp,omitempty"`
IceServers []IceServer `json:"iceServers,omitempty"`
}
type SessionDescription struct {
Type string `json:"type"`
SDP string `json:"sdp"`
}
// TrackSpec 描述一条要加入 session 的轨道。
// location=local 表示本端发布;location=remote + SessionID 表示订阅对端 SFU 轨道。
type TrackSpec struct {
Location string `json:"location"`
SessionID string `json:"sessionId,omitempty"`
TrackName string `json:"trackName,omitempty"`
Kind string `json:"kind,omitempty"`
BidirectionalMediaStream bool `json:"bidirectionalMediaStream,omitempty"`
}
type TrackRequest struct {
SessionDescription *SessionDescription `json:"sessionDescription,omitempty"`
Tracks []TrackSpec `json:"tracks,omitempty"`
AutoDiscover bool `json:"autoDiscover,omitempty"`
}
type TrackResult struct {
TrackName string `json:"trackName,omitempty"`
MID string `json:"mid,omitempty"`
Location string `json:"location,omitempty"`
SessionID string `json:"sessionId,omitempty"`
ErrorCode string `json:"errorCode,omitempty"`
ErrorDescription string `json:"errorDescription,omitempty"`
}
type TracksResponse struct {
SessionDescription *SessionDescription `json:"sessionDescription,omitempty"`
Tracks []TrackResult `json:"tracks,omitempty"`
RequiresImmediateRenegotiation bool `json:"requiresImmediateRenegotiation,omitempty"`
ErrorCode string `json:"errorCode,omitempty"`
ErrorDescription string `json:"errorDescription,omitempty"`
}
// Client 封装 Cloudflare Realtime REST API(服务端持有 AppSecret)。
type Client struct {
appID string
appSecret string
baseURL string
httpClient *http.Client
}
func NewClient(appID, appSecret, baseURL string) *Client {
if baseURL == "" {
baseURL = defaultBaseURL
}
return &Client{
appID: appID,
appSecret: appSecret,
baseURL: baseURL,
httpClient: &http.Client{Timeout: 30 * time.Second},
}
}
func (c *Client) appsPath() string { return c.baseURL + "/apps/" + c.appID }
// CreateSession 创建一个空的 Cloudflare Realtime session(correlationId 用于房间关联)。
func (c *Client) CreateSession(correlationID string) (string, error) {
path := c.appsPath() + "/sessions/new"
if correlationID != "" {
path += "?correlationId=" + url.QueryEscape(correlationID)
}
var resp struct {
SessionID string `json:"sessionId"`
}
if err := c.doJSON(http.MethodPost, path, nil, &resp); err != nil {
return "", err
}
return resp.SessionID, nil
}
func (c *Client) GetSession(sessionID string) (*SessionInfo, error) {
var resp SessionInfo
if err := c.doJSON(http.MethodGet, c.appsPath()+"/sessions/"+sessionID, nil, &resp); err != nil {
return nil, err
}
return &resp, nil
}
// AddTracks 在 session 上发布(location=local)或订阅(location=remote)轨道,返回 SDP answer。
func (c *Client) AddTracks(sessionID string, req *TrackRequest) (*TracksResponse, error) {
var resp TracksResponse
if err := c.doJSON(http.MethodPost, c.appsPath()+"/sessions/"+sessionID+"/tracks/new", req, &resp); err != nil {
return nil, err
}
return &resp, nil
}
// DeleteSession 终止 session 并关闭其全部轨道(用于停推 / 踢人)。
func (c *Client) DeleteSession(sessionID string) error {
return c.doJSON(http.MethodDelete, c.appsPath()+"/sessions/"+sessionID, nil, nil)
}
func (c *Client) doJSON(method, path string, body, target interface{}) error {
var reader io.Reader
if body != nil {
data, err := json.Marshal(body)
if err != nil {
return fmt.Errorf("cf marshal: %w", err)
}
reader = bytes.NewReader(data)
}
req, err := http.NewRequest(method, path, reader)
if err != nil {
return fmt.Errorf("cf build request: %w", err)
}
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
if c.appSecret != "" {
req.Header.Set("Authorization", "Bearer "+c.appSecret)
}
resp, err := c.httpClient.Do(req)
if err != nil {
return fmt.Errorf("cf request: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
b, _ := io.ReadAll(resp.Body)
return fmt.Errorf("cf api error status=%d body=%s", resp.StatusCode, string(b))
}
if target == nil {
return nil
}
if err := json.NewDecoder(resp.Body).Decode(target); err != nil {
return fmt.Errorf("cf decode: %w", err)
}
return nil
}

View File

@ -0,0 +1,84 @@
package cloudflare
import (
"encoding/json"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
)
// 用 httptest 校验 Cloudflare Realtime REST 的请求契约(路径 / Bearer / body 形状),
// 无需真实凭证即可验证主 SFU 的集成代码路径。
func TestClientSessionAndTracks(t *testing.T) {
var gotAuth string
var gotBody TrackRequest
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
gotAuth = r.Header.Get("Authorization")
_ = json.Unmarshal(body, &gotBody)
switch {
case strings.HasSuffix(r.URL.Path, "/sessions/new"):
w.Write([]byte(`{"sessionId":"sess-pub"}`))
case strings.HasSuffix(r.URL.Path, "/tracks/new"):
w.Write([]byte(`{"sessionDescription":{"type":"answer","sdp":"v=0"},"tracks":[]}`))
case r.Method == http.MethodDelete:
w.WriteHeader(http.StatusOK)
default:
w.WriteHeader(http.StatusNotFound)
}
}))
defer srv.Close()
c := NewClient("app1", "secret1", srv.URL+"/v1")
if c.appID != "app1" {
t.Fatalf("appID = %q", c.appID)
}
sid, err := c.CreateSession("room1")
if err != nil {
t.Fatalf("CreateSession: %v", err)
}
if sid != "sess-pub" {
t.Fatalf("sessionId = %q", sid)
}
if gotAuth != "Bearer secret1" {
t.Fatalf("auth = %q", gotAuth)
}
// 发布:location=local
resp, err := c.AddTracks(sid, &TrackRequest{
SessionDescription: &SessionDescription{Type: "offer", SDP: "v=0"},
Tracks: []TrackSpec{{Location: "local", TrackName: "video", Kind: "video"}},
AutoDiscover: true,
})
if err != nil {
t.Fatalf("AddTracks publish: %v", err)
}
if resp.SessionDescription == nil || resp.SessionDescription.Type != "answer" {
t.Fatalf("unexpected answer: %+v", resp.SessionDescription)
}
if len(gotBody.Tracks) != 1 || gotBody.Tracks[0].Location != "local" {
t.Fatalf("publish track location = %+v", gotBody.Tracks)
}
if !gotBody.AutoDiscover {
t.Fatalf("autodiscover not set")
}
// 订阅:location=remote + 发布者 sessionId
_, err = c.AddTracks("sess-view", &TrackRequest{
Tracks: []TrackSpec{{Location: "remote", SessionID: "sess-pub", TrackName: "video"}},
})
if err != nil {
t.Fatalf("AddTracks subscribe: %v", err)
}
if len(gotBody.Tracks) != 1 || gotBody.Tracks[0].Location != "remote" || gotBody.Tracks[0].SessionID != "sess-pub" {
t.Fatalf("subscribe track = %+v", gotBody.Tracks)
}
if err := c.DeleteSession("sess-pub"); err != nil {
t.Fatalf("DeleteSession: %v", err)
}
}

View File

@ -0,0 +1,42 @@
package cloudflare
import "gospeak-live-sfu-demo/gen"
// Provider 暴露 Cloudflare Realtime 后端的能力信息,并持有底层 REST 客户端。
type Provider struct {
client *Client
appID string
stun string
configured bool
}
func NewProvider(appID, appSecret, baseURL, stun string) *Provider {
return &Provider{
client: NewClient(appID, appSecret, baseURL),
appID: appID,
stun: stun,
configured: appID != "" && appSecret != "",
}
}
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) Client() *Client { return p.client }
func (p *Provider) AppID() string { return p.appID }
func (p *Provider) Stun() string {
if p.stun == "" {
return "stun:stun.cloudflare.com:3478"
}
return p.stun
}
func (p *Provider) BackendInfo() *gen.BackendInfo {
return &gen.BackendInfo{
Kind: gen.BackendKind_BACKEND_KIND_CLOUDFLARE,
Name: "Cloudflare Realtime",
Configured: p.configured,
Primary: true,
}
}

View File

@ -0,0 +1,32 @@
package srs
import "gospeak-live-sfu-demo/gen"
// Provider 暴露 SRS 后端能力(本地 docker 默认可达)。
type Provider struct {
baseURL string
app string
secret string
candidate string
}
func NewProvider(baseURL, app, secret, candidate string) *Provider {
return &Provider{baseURL: baseURL, app: app, secret: secret, candidate: candidate}
}
func (p *Provider) Name() string { return "srs" }
func (p *Provider) Primary() bool { return false }
func (p *Provider) Configured() bool { return true }
func (p *Provider) BaseURL() string { return p.baseURL }
func (p *Provider) App() string { return p.app }
func (p *Provider) Secret() string { return p.secret }
func (p *Provider) Candidate() string { return p.candidate }
func (p *Provider) BackendInfo() *gen.BackendInfo {
return &gen.BackendInfo{
Kind: gen.BackendKind_BACKEND_KIND_SRS,
Name: "SRS",
Configured: true,
Primary: false,
}
}