From c4a50b07c6470b2fc120c339d3609fec8d614fdf Mon Sep 17 00:00:00 2001 From: noelorin Date: Mon, 17 Aug 2026 15:14:51 +0800 Subject: [PATCH] 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 --- .env.example | 22 + .gitignore | 3 + README.md | 76 ++ api/buf.yaml | 7 + api/live_sfu.proto | 107 +++ buf.gen.yaml | 13 + cmd/server/main.go | 24 + deploy/docker-compose.yml | 16 + deploy/srs.conf | 38 + gen/live_sfu.pb.go | 1141 ++++++++++++++++++++++++ gen/live_sfu_grpc.pb.go | 325 +++++++ go.mod | 17 + go.sum | 38 + internal/config/config.go | 85 ++ internal/server/gateway.go | 160 ++++ internal/server/grpc.go | 41 + internal/server/grpc_test.go | 55 ++ internal/server/proxy.go | 52 ++ internal/server/rooms.go | 112 +++ internal/server/server.go | 91 ++ internal/server/service.go | 217 +++++ internal/server/static/app.js | 230 +++++ internal/server/static/index.html | 54 ++ internal/server/static/publish.html | 43 + internal/server/static/styles.css | 79 ++ internal/server/static/watch.html | 41 + internal/server/token.go | 65 ++ internal/server/token_test.go | 40 + internal/sfu/cloudflare/client.go | 161 ++++ internal/sfu/cloudflare/client_test.go | 84 ++ internal/sfu/cloudflare/provider.go | 42 + internal/sfu/srs/provider.go | 32 + 32 files changed, 3511 insertions(+) create mode 100644 .env.example create mode 100644 .gitignore create mode 100644 README.md create mode 100644 api/buf.yaml create mode 100644 api/live_sfu.proto create mode 100644 buf.gen.yaml create mode 100644 cmd/server/main.go create mode 100644 deploy/docker-compose.yml create mode 100644 deploy/srs.conf create mode 100644 gen/live_sfu.pb.go create mode 100644 gen/live_sfu_grpc.pb.go create mode 100644 go.mod create mode 100644 go.sum create mode 100644 internal/config/config.go create mode 100644 internal/server/gateway.go create mode 100644 internal/server/grpc.go create mode 100644 internal/server/grpc_test.go create mode 100644 internal/server/proxy.go create mode 100644 internal/server/rooms.go create mode 100644 internal/server/server.go create mode 100644 internal/server/service.go create mode 100644 internal/server/static/app.js create mode 100644 internal/server/static/index.html create mode 100644 internal/server/static/publish.html create mode 100644 internal/server/static/styles.css create mode 100644 internal/server/static/watch.html create mode 100644 internal/server/token.go create mode 100644 internal/server/token_test.go create mode 100644 internal/sfu/cloudflare/client.go create mode 100644 internal/sfu/cloudflare/client_test.go create mode 100644 internal/sfu/cloudflare/provider.go create mode 100644 internal/sfu/srs/provider.go diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..eff815f --- /dev/null +++ b/.env.example @@ -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 diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..82194b9 --- /dev/null +++ b/.gitignore @@ -0,0 +1,3 @@ +.DS_Store +.env +*.local.env diff --git a/README.md b/README.md new file mode 100644 index 0000000..4cad176 --- /dev/null +++ b/README.md @@ -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 实时广播,观众任选后端拉流 +``` diff --git a/api/buf.yaml b/api/buf.yaml new file mode 100644 index 0000000..4cb6a23 --- /dev/null +++ b/api/buf.yaml @@ -0,0 +1,7 @@ +version: v2 +lint: + use: + - DEFAULT +breaking: + use: + - FILE diff --git a/api/live_sfu.proto b/api/live_sfu.proto new file mode 100644 index 0000000..6dc37a7 --- /dev/null +++ b/api/live_sfu.proto @@ -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; +} diff --git a/buf.gen.yaml b/buf.gen.yaml new file mode 100644 index 0000000..eb7346d --- /dev/null +++ b/buf.gen.yaml @@ -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 diff --git a/cmd/server/main.go b/cmd/server/main.go new file mode 100644 index 0000000..8b4844f --- /dev/null +++ b/cmd/server/main.go @@ -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) + } +} diff --git a/deploy/docker-compose.yml b/deploy/docker-compose.yml new file mode 100644 index 0000000..ed22981 --- /dev/null +++ b/deploy/docker-compose.yml @@ -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 diff --git a/deploy/srs.conf b/deploy/srs.conf new file mode 100644 index 0000000..f5ec0b1 --- /dev/null +++ b/deploy/srs.conf @@ -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; + } +} diff --git a/gen/live_sfu.pb.go b/gen/live_sfu.pb.go new file mode 100644 index 0000000..cb85d45 --- /dev/null +++ b/gen/live_sfu.pb.go @@ -0,0 +1,1141 @@ +// Code generated by protoc-gen-go. DO NOT EDIT. +// versions: +// protoc-gen-go v1.36.11 +// protoc (unknown) +// source: live_sfu.proto + +package gen + +import ( + protoreflect "google.golang.org/protobuf/reflect/protoreflect" + protoimpl "google.golang.org/protobuf/runtime/protoimpl" + reflect "reflect" + sync "sync" + unsafe "unsafe" +) + +const ( + // Verify that this generated code is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(20 - protoimpl.MinVersion) + // Verify that runtime/protoimpl is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) +) + +type BackendKind int32 + +const ( + BackendKind_BACKEND_KIND_UNSPECIFIED BackendKind = 0 + BackendKind_BACKEND_KIND_CLOUDFLARE BackendKind = 1 // Cloudflare Realtime(主 SFU) + BackendKind_BACKEND_KIND_SRS BackendKind = 2 // SRS(次级 / 本地对照) +) + +// Enum value maps for BackendKind. +var ( + BackendKind_name = map[int32]string{ + 0: "BACKEND_KIND_UNSPECIFIED", + 1: "BACKEND_KIND_CLOUDFLARE", + 2: "BACKEND_KIND_SRS", + } + BackendKind_value = map[string]int32{ + "BACKEND_KIND_UNSPECIFIED": 0, + "BACKEND_KIND_CLOUDFLARE": 1, + "BACKEND_KIND_SRS": 2, + } +) + +func (x BackendKind) Enum() *BackendKind { + p := new(BackendKind) + *p = x + return p +} + +func (x BackendKind) String() string { + return protoimpl.X.EnumStringOf(x.Descriptor(), protoreflect.EnumNumber(x)) +} + +func (BackendKind) Descriptor() protoreflect.EnumDescriptor { + return file_live_sfu_proto_enumTypes[0].Descriptor() +} + +func (BackendKind) Type() protoreflect.EnumType { + return &file_live_sfu_proto_enumTypes[0] +} + +func (x BackendKind) Number() protoreflect.EnumNumber { + return protoreflect.EnumNumber(x) +} + +// Deprecated: Use BackendKind.Descriptor instead. +func (BackendKind) EnumDescriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{0} +} + +type IceServer struct { + state protoimpl.MessageState `protogen:"open.v1"` + Urls []string `protobuf:"bytes,1,rep,name=urls,proto3" json:"urls,omitempty"` + Username string `protobuf:"bytes,2,opt,name=username,proto3" json:"username,omitempty"` + Credential string `protobuf:"bytes,3,opt,name=credential,proto3" json:"credential,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *IceServer) Reset() { + *x = IceServer{} + mi := &file_live_sfu_proto_msgTypes[0] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *IceServer) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*IceServer) ProtoMessage() {} + +func (x *IceServer) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[0] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use IceServer.ProtoReflect.Descriptor instead. +func (*IceServer) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{0} +} + +func (x *IceServer) GetUrls() []string { + if x != nil { + return x.Urls + } + return nil +} + +func (x *IceServer) GetUsername() string { + if x != nil { + return x.Username + } + return "" +} + +func (x *IceServer) GetCredential() string { + if x != nil { + return x.Credential + } + return "" +} + +type GetConfigRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *GetConfigRequest) Reset() { + *x = GetConfigRequest{} + mi := &file_live_sfu_proto_msgTypes[1] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *GetConfigRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*GetConfigRequest) ProtoMessage() {} + +func (x *GetConfigRequest) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[1] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use GetConfigRequest.ProtoReflect.Descriptor instead. +func (*GetConfigRequest) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{1} +} + +type BackendInfo struct { + state protoimpl.MessageState `protogen:"open.v1"` + Kind BackendKind `protobuf:"varint,1,opt,name=kind,proto3,enum=gospeak.livedemo.v1.BackendKind" json:"kind,omitempty"` + Name string `protobuf:"bytes,2,opt,name=name,proto3" json:"name,omitempty"` + Configured bool `protobuf:"varint,3,opt,name=configured,proto3" json:"configured,omitempty"` + Primary bool `protobuf:"varint,4,opt,name=primary,proto3" json:"primary,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *BackendInfo) Reset() { + *x = BackendInfo{} + mi := &file_live_sfu_proto_msgTypes[2] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *BackendInfo) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*BackendInfo) ProtoMessage() {} + +func (x *BackendInfo) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[2] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use BackendInfo.ProtoReflect.Descriptor instead. +func (*BackendInfo) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{2} +} + +func (x *BackendInfo) GetKind() BackendKind { + if x != nil { + return x.Kind + } + return BackendKind_BACKEND_KIND_UNSPECIFIED +} + +func (x *BackendInfo) GetName() string { + if x != nil { + return x.Name + } + return "" +} + +func (x *BackendInfo) GetConfigured() bool { + if x != nil { + return x.Configured + } + return false +} + +func (x *BackendInfo) GetPrimary() bool { + if x != nil { + return x.Primary + } + return false +} + +type GetConfigResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Backends []*BackendInfo `protobuf:"bytes,1,rep,name=backends,proto3" json:"backends,omitempty"` + Candidate string `protobuf:"bytes,2,opt,name=candidate,proto3" json:"candidate,omitempty"` + TokenRequired bool `protobuf:"varint,3,opt,name=token_required,json=tokenRequired,proto3" json:"token_required,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *GetConfigResponse) Reset() { + *x = GetConfigResponse{} + mi := &file_live_sfu_proto_msgTypes[3] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *GetConfigResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*GetConfigResponse) ProtoMessage() {} + +func (x *GetConfigResponse) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[3] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use GetConfigResponse.ProtoReflect.Descriptor instead. +func (*GetConfigResponse) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{3} +} + +func (x *GetConfigResponse) GetBackends() []*BackendInfo { + if x != nil { + return x.Backends + } + return nil +} + +func (x *GetConfigResponse) GetCandidate() string { + if x != nil { + return x.Candidate + } + return "" +} + +func (x *GetConfigResponse) GetTokenRequired() bool { + if x != nil { + return x.TokenRequired + } + return false +} + +// StreamTarget 表示一个房间在某后端的一条分发目标(SFU 扇出出口)。 +type StreamTarget struct { + state protoimpl.MessageState `protogen:"open.v1"` + Backend BackendKind `protobuf:"varint,1,opt,name=backend,proto3,enum=gospeak.livedemo.v1.BackendKind" json:"backend,omitempty"` + SessionId string `protobuf:"bytes,2,opt,name=session_id,json=sessionId,proto3" json:"session_id,omitempty"` // Cloudflare: 发布者 sessionId + Stream string `protobuf:"bytes,3,opt,name=stream,proto3" json:"stream,omitempty"` // SRS: stream 名 + PublishToken string `protobuf:"bytes,4,opt,name=publish_token,json=publishToken,proto3" json:"publish_token,omitempty"` // SRS: 推流 JWT(可选) + Url string `protobuf:"bytes,5,opt,name=url,proto3" json:"url,omitempty"` // 拉流播放地址(SRS WHEP,可选) + PublishedAt int64 `protobuf:"varint,6,opt,name=published_at,json=publishedAt,proto3" json:"published_at,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *StreamTarget) Reset() { + *x = StreamTarget{} + mi := &file_live_sfu_proto_msgTypes[4] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *StreamTarget) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*StreamTarget) ProtoMessage() {} + +func (x *StreamTarget) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[4] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use StreamTarget.ProtoReflect.Descriptor instead. +func (*StreamTarget) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{4} +} + +func (x *StreamTarget) GetBackend() BackendKind { + if x != nil { + return x.Backend + } + return BackendKind_BACKEND_KIND_UNSPECIFIED +} + +func (x *StreamTarget) GetSessionId() string { + if x != nil { + return x.SessionId + } + return "" +} + +func (x *StreamTarget) GetStream() string { + if x != nil { + return x.Stream + } + return "" +} + +func (x *StreamTarget) GetPublishToken() string { + if x != nil { + return x.PublishToken + } + return "" +} + +func (x *StreamTarget) GetUrl() string { + if x != nil { + return x.Url + } + return "" +} + +func (x *StreamTarget) GetPublishedAt() int64 { + if x != nil { + return x.PublishedAt + } + return 0 +} + +type Room struct { + state protoimpl.MessageState `protogen:"open.v1"` + Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"` + Targets []*StreamTarget `protobuf:"bytes,2,rep,name=targets,proto3" json:"targets,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *Room) Reset() { + *x = Room{} + mi := &file_live_sfu_proto_msgTypes[5] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *Room) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*Room) ProtoMessage() {} + +func (x *Room) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[5] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use Room.ProtoReflect.Descriptor instead. +func (*Room) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{5} +} + +func (x *Room) GetName() string { + if x != nil { + return x.Name + } + return "" +} + +func (x *Room) GetTargets() []*StreamTarget { + if x != nil { + return x.Targets + } + return nil +} + +type ListRoomsRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListRoomsRequest) Reset() { + *x = ListRoomsRequest{} + mi := &file_live_sfu_proto_msgTypes[6] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListRoomsRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListRoomsRequest) ProtoMessage() {} + +func (x *ListRoomsRequest) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[6] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListRoomsRequest.ProtoReflect.Descriptor instead. +func (*ListRoomsRequest) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{6} +} + +type ListRoomsResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Rooms []*Room `protobuf:"bytes,1,rep,name=rooms,proto3" json:"rooms,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListRoomsResponse) Reset() { + *x = ListRoomsResponse{} + mi := &file_live_sfu_proto_msgTypes[7] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListRoomsResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListRoomsResponse) ProtoMessage() {} + +func (x *ListRoomsResponse) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[7] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListRoomsResponse.ProtoReflect.Descriptor instead. +func (*ListRoomsResponse) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{7} +} + +func (x *ListRoomsResponse) GetRooms() []*Room { + if x != nil { + return x.Rooms + } + return nil +} + +type PublishRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Room string `protobuf:"bytes,1,opt,name=room,proto3" json:"room,omitempty"` + Backend BackendKind `protobuf:"varint,2,opt,name=backend,proto3,enum=gospeak.livedemo.v1.BackendKind" json:"backend,omitempty"` + Identity string `protobuf:"bytes,3,opt,name=identity,proto3" json:"identity,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *PublishRequest) Reset() { + *x = PublishRequest{} + mi := &file_live_sfu_proto_msgTypes[8] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *PublishRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*PublishRequest) ProtoMessage() {} + +func (x *PublishRequest) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[8] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use PublishRequest.ProtoReflect.Descriptor instead. +func (*PublishRequest) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{8} +} + +func (x *PublishRequest) GetRoom() string { + if x != nil { + return x.Room + } + return "" +} + +func (x *PublishRequest) GetBackend() BackendKind { + if x != nil { + return x.Backend + } + return BackendKind_BACKEND_KIND_UNSPECIFIED +} + +func (x *PublishRequest) GetIdentity() string { + if x != nil { + return x.Identity + } + return "" +} + +type PublishResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + SessionId string `protobuf:"bytes,1,opt,name=session_id,json=sessionId,proto3" json:"session_id,omitempty"` + Stream string `protobuf:"bytes,2,opt,name=stream,proto3" json:"stream,omitempty"` + PublishToken string `protobuf:"bytes,3,opt,name=publish_token,json=publishToken,proto3" json:"publish_token,omitempty"` + IceServers []*IceServer `protobuf:"bytes,4,rep,name=ice_servers,json=iceServers,proto3" json:"ice_servers,omitempty"` + Target *StreamTarget `protobuf:"bytes,5,opt,name=target,proto3" json:"target,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *PublishResponse) Reset() { + *x = PublishResponse{} + mi := &file_live_sfu_proto_msgTypes[9] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *PublishResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*PublishResponse) ProtoMessage() {} + +func (x *PublishResponse) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[9] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use PublishResponse.ProtoReflect.Descriptor instead. +func (*PublishResponse) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{9} +} + +func (x *PublishResponse) GetSessionId() string { + if x != nil { + return x.SessionId + } + return "" +} + +func (x *PublishResponse) GetStream() string { + if x != nil { + return x.Stream + } + return "" +} + +func (x *PublishResponse) GetPublishToken() string { + if x != nil { + return x.PublishToken + } + return "" +} + +func (x *PublishResponse) GetIceServers() []*IceServer { + if x != nil { + return x.IceServers + } + return nil +} + +func (x *PublishResponse) GetTarget() *StreamTarget { + if x != nil { + return x.Target + } + return nil +} + +type SubscribeRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Room string `protobuf:"bytes,1,opt,name=room,proto3" json:"room,omitempty"` + Backend BackendKind `protobuf:"varint,2,opt,name=backend,proto3,enum=gospeak.livedemo.v1.BackendKind" json:"backend,omitempty"` + Identity string `protobuf:"bytes,3,opt,name=identity,proto3" json:"identity,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *SubscribeRequest) Reset() { + *x = SubscribeRequest{} + mi := &file_live_sfu_proto_msgTypes[10] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *SubscribeRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*SubscribeRequest) ProtoMessage() {} + +func (x *SubscribeRequest) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[10] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use SubscribeRequest.ProtoReflect.Descriptor instead. +func (*SubscribeRequest) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{10} +} + +func (x *SubscribeRequest) GetRoom() string { + if x != nil { + return x.Room + } + return "" +} + +func (x *SubscribeRequest) GetBackend() BackendKind { + if x != nil { + return x.Backend + } + return BackendKind_BACKEND_KIND_UNSPECIFIED +} + +func (x *SubscribeRequest) GetIdentity() string { + if x != nil { + return x.Identity + } + return "" +} + +type SubscribeResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + SessionId string `protobuf:"bytes,1,opt,name=session_id,json=sessionId,proto3" json:"session_id,omitempty"` + Stream string `protobuf:"bytes,2,opt,name=stream,proto3" json:"stream,omitempty"` + PublisherSessionId string `protobuf:"bytes,3,opt,name=publisher_session_id,json=publisherSessionId,proto3" json:"publisher_session_id,omitempty"` + IceServers []*IceServer `protobuf:"bytes,4,rep,name=ice_servers,json=iceServers,proto3" json:"ice_servers,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *SubscribeResponse) Reset() { + *x = SubscribeResponse{} + mi := &file_live_sfu_proto_msgTypes[11] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *SubscribeResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*SubscribeResponse) ProtoMessage() {} + +func (x *SubscribeResponse) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[11] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use SubscribeResponse.ProtoReflect.Descriptor instead. +func (*SubscribeResponse) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{11} +} + +func (x *SubscribeResponse) GetSessionId() string { + if x != nil { + return x.SessionId + } + return "" +} + +func (x *SubscribeResponse) GetStream() string { + if x != nil { + return x.Stream + } + return "" +} + +func (x *SubscribeResponse) GetPublisherSessionId() string { + if x != nil { + return x.PublisherSessionId + } + return "" +} + +func (x *SubscribeResponse) GetIceServers() []*IceServer { + if x != nil { + return x.IceServers + } + return nil +} + +type StopStreamRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Room string `protobuf:"bytes,1,opt,name=room,proto3" json:"room,omitempty"` + Backend BackendKind `protobuf:"varint,2,opt,name=backend,proto3,enum=gospeak.livedemo.v1.BackendKind" json:"backend,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *StopStreamRequest) Reset() { + *x = StopStreamRequest{} + mi := &file_live_sfu_proto_msgTypes[12] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *StopStreamRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*StopStreamRequest) ProtoMessage() {} + +func (x *StopStreamRequest) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[12] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use StopStreamRequest.ProtoReflect.Descriptor instead. +func (*StopStreamRequest) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{12} +} + +func (x *StopStreamRequest) GetRoom() string { + if x != nil { + return x.Room + } + return "" +} + +func (x *StopStreamRequest) GetBackend() BackendKind { + if x != nil { + return x.Backend + } + return BackendKind_BACKEND_KIND_UNSPECIFIED +} + +type StopStreamResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Ok bool `protobuf:"varint,1,opt,name=ok,proto3" json:"ok,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *StopStreamResponse) Reset() { + *x = StopStreamResponse{} + mi := &file_live_sfu_proto_msgTypes[13] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *StopStreamResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*StopStreamResponse) ProtoMessage() {} + +func (x *StopStreamResponse) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[13] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use StopStreamResponse.ProtoReflect.Descriptor instead. +func (*StopStreamResponse) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{13} +} + +func (x *StopStreamResponse) GetOk() bool { + if x != nil { + return x.Ok + } + return false +} + +type WatchRoomRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Room string `protobuf:"bytes,1,opt,name=room,proto3" json:"room,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *WatchRoomRequest) Reset() { + *x = WatchRoomRequest{} + mi := &file_live_sfu_proto_msgTypes[14] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *WatchRoomRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*WatchRoomRequest) ProtoMessage() {} + +func (x *WatchRoomRequest) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[14] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use WatchRoomRequest.ProtoReflect.Descriptor instead. +func (*WatchRoomRequest) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{14} +} + +func (x *WatchRoomRequest) GetRoom() string { + if x != nil { + return x.Room + } + return "" +} + +type RoomEvent struct { + state protoimpl.MessageState `protogen:"open.v1"` + Room string `protobuf:"bytes,1,opt,name=room,proto3" json:"room,omitempty"` + Targets []*StreamTarget `protobuf:"bytes,2,rep,name=targets,proto3" json:"targets,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *RoomEvent) Reset() { + *x = RoomEvent{} + mi := &file_live_sfu_proto_msgTypes[15] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *RoomEvent) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*RoomEvent) ProtoMessage() {} + +func (x *RoomEvent) ProtoReflect() protoreflect.Message { + mi := &file_live_sfu_proto_msgTypes[15] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use RoomEvent.ProtoReflect.Descriptor instead. +func (*RoomEvent) Descriptor() ([]byte, []int) { + return file_live_sfu_proto_rawDescGZIP(), []int{15} +} + +func (x *RoomEvent) GetRoom() string { + if x != nil { + return x.Room + } + return "" +} + +func (x *RoomEvent) GetTargets() []*StreamTarget { + if x != nil { + return x.Targets + } + return nil +} + +var File_live_sfu_proto protoreflect.FileDescriptor + +const file_live_sfu_proto_rawDesc = "" + + "\n" + + "\x0elive_sfu.proto\x12\x13gospeak.livedemo.v1\"[\n" + + "\tIceServer\x12\x12\n" + + "\x04urls\x18\x01 \x03(\tR\x04urls\x12\x1a\n" + + "\busername\x18\x02 \x01(\tR\busername\x12\x1e\n" + + "\n" + + "credential\x18\x03 \x01(\tR\n" + + "credential\"\x12\n" + + "\x10GetConfigRequest\"\x91\x01\n" + + "\vBackendInfo\x124\n" + + "\x04kind\x18\x01 \x01(\x0e2 .gospeak.livedemo.v1.BackendKindR\x04kind\x12\x12\n" + + "\x04name\x18\x02 \x01(\tR\x04name\x12\x1e\n" + + "\n" + + "configured\x18\x03 \x01(\bR\n" + + "configured\x12\x18\n" + + "\aprimary\x18\x04 \x01(\bR\aprimary\"\x96\x01\n" + + "\x11GetConfigResponse\x12<\n" + + "\bbackends\x18\x01 \x03(\v2 .gospeak.livedemo.v1.BackendInfoR\bbackends\x12\x1c\n" + + "\tcandidate\x18\x02 \x01(\tR\tcandidate\x12%\n" + + "\x0etoken_required\x18\x03 \x01(\bR\rtokenRequired\"\xdb\x01\n" + + "\fStreamTarget\x12:\n" + + "\abackend\x18\x01 \x01(\x0e2 .gospeak.livedemo.v1.BackendKindR\abackend\x12\x1d\n" + + "\n" + + "session_id\x18\x02 \x01(\tR\tsessionId\x12\x16\n" + + "\x06stream\x18\x03 \x01(\tR\x06stream\x12#\n" + + "\rpublish_token\x18\x04 \x01(\tR\fpublishToken\x12\x10\n" + + "\x03url\x18\x05 \x01(\tR\x03url\x12!\n" + + "\fpublished_at\x18\x06 \x01(\x03R\vpublishedAt\"W\n" + + "\x04Room\x12\x12\n" + + "\x04name\x18\x01 \x01(\tR\x04name\x12;\n" + + "\atargets\x18\x02 \x03(\v2!.gospeak.livedemo.v1.StreamTargetR\atargets\"\x12\n" + + "\x10ListRoomsRequest\"D\n" + + "\x11ListRoomsResponse\x12/\n" + + "\x05rooms\x18\x01 \x03(\v2\x19.gospeak.livedemo.v1.RoomR\x05rooms\"|\n" + + "\x0ePublishRequest\x12\x12\n" + + "\x04room\x18\x01 \x01(\tR\x04room\x12:\n" + + "\abackend\x18\x02 \x01(\x0e2 .gospeak.livedemo.v1.BackendKindR\abackend\x12\x1a\n" + + "\bidentity\x18\x03 \x01(\tR\bidentity\"\xe9\x01\n" + + "\x0fPublishResponse\x12\x1d\n" + + "\n" + + "session_id\x18\x01 \x01(\tR\tsessionId\x12\x16\n" + + "\x06stream\x18\x02 \x01(\tR\x06stream\x12#\n" + + "\rpublish_token\x18\x03 \x01(\tR\fpublishToken\x12?\n" + + "\vice_servers\x18\x04 \x03(\v2\x1e.gospeak.livedemo.v1.IceServerR\n" + + "iceServers\x129\n" + + "\x06target\x18\x05 \x01(\v2!.gospeak.livedemo.v1.StreamTargetR\x06target\"~\n" + + "\x10SubscribeRequest\x12\x12\n" + + "\x04room\x18\x01 \x01(\tR\x04room\x12:\n" + + "\abackend\x18\x02 \x01(\x0e2 .gospeak.livedemo.v1.BackendKindR\abackend\x12\x1a\n" + + "\bidentity\x18\x03 \x01(\tR\bidentity\"\xbd\x01\n" + + "\x11SubscribeResponse\x12\x1d\n" + + "\n" + + "session_id\x18\x01 \x01(\tR\tsessionId\x12\x16\n" + + "\x06stream\x18\x02 \x01(\tR\x06stream\x120\n" + + "\x14publisher_session_id\x18\x03 \x01(\tR\x12publisherSessionId\x12?\n" + + "\vice_servers\x18\x04 \x03(\v2\x1e.gospeak.livedemo.v1.IceServerR\n" + + "iceServers\"c\n" + + "\x11StopStreamRequest\x12\x12\n" + + "\x04room\x18\x01 \x01(\tR\x04room\x12:\n" + + "\abackend\x18\x02 \x01(\x0e2 .gospeak.livedemo.v1.BackendKindR\abackend\"$\n" + + "\x12StopStreamResponse\x12\x0e\n" + + "\x02ok\x18\x01 \x01(\bR\x02ok\"&\n" + + "\x10WatchRoomRequest\x12\x12\n" + + "\x04room\x18\x01 \x01(\tR\x04room\"\\\n" + + "\tRoomEvent\x12\x12\n" + + "\x04room\x18\x01 \x01(\tR\x04room\x12;\n" + + "\atargets\x18\x02 \x03(\v2!.gospeak.livedemo.v1.StreamTargetR\atargets*^\n" + + "\vBackendKind\x12\x1c\n" + + "\x18BACKEND_KIND_UNSPECIFIED\x10\x00\x12\x1b\n" + + "\x17BACKEND_KIND_CLOUDFLARE\x10\x01\x12\x14\n" + + "\x10BACKEND_KIND_SRS\x10\x022\xa8\x04\n" + + "\aLiveSFU\x12Z\n" + + "\tGetConfig\x12%.gospeak.livedemo.v1.GetConfigRequest\x1a&.gospeak.livedemo.v1.GetConfigResponse\x12Z\n" + + "\tListRooms\x12%.gospeak.livedemo.v1.ListRoomsRequest\x1a&.gospeak.livedemo.v1.ListRoomsResponse\x12T\n" + + "\aPublish\x12#.gospeak.livedemo.v1.PublishRequest\x1a$.gospeak.livedemo.v1.PublishResponse\x12Z\n" + + "\tSubscribe\x12%.gospeak.livedemo.v1.SubscribeRequest\x1a&.gospeak.livedemo.v1.SubscribeResponse\x12]\n" + + "\n" + + "StopStream\x12&.gospeak.livedemo.v1.StopStreamRequest\x1a'.gospeak.livedemo.v1.StopStreamResponse\x12T\n" + + "\tWatchRoom\x12%.gospeak.livedemo.v1.WatchRoomRequest\x1a\x1e.gospeak.livedemo.v1.RoomEvent0\x01B\x1bZ\x19gospeak-live-sfu-demo/genb\x06proto3" + +var ( + file_live_sfu_proto_rawDescOnce sync.Once + file_live_sfu_proto_rawDescData []byte +) + +func file_live_sfu_proto_rawDescGZIP() []byte { + file_live_sfu_proto_rawDescOnce.Do(func() { + file_live_sfu_proto_rawDescData = protoimpl.X.CompressGZIP(unsafe.Slice(unsafe.StringData(file_live_sfu_proto_rawDesc), len(file_live_sfu_proto_rawDesc))) + }) + return file_live_sfu_proto_rawDescData +} + +var file_live_sfu_proto_enumTypes = make([]protoimpl.EnumInfo, 1) +var file_live_sfu_proto_msgTypes = make([]protoimpl.MessageInfo, 16) +var file_live_sfu_proto_goTypes = []any{ + (BackendKind)(0), // 0: gospeak.livedemo.v1.BackendKind + (*IceServer)(nil), // 1: gospeak.livedemo.v1.IceServer + (*GetConfigRequest)(nil), // 2: gospeak.livedemo.v1.GetConfigRequest + (*BackendInfo)(nil), // 3: gospeak.livedemo.v1.BackendInfo + (*GetConfigResponse)(nil), // 4: gospeak.livedemo.v1.GetConfigResponse + (*StreamTarget)(nil), // 5: gospeak.livedemo.v1.StreamTarget + (*Room)(nil), // 6: gospeak.livedemo.v1.Room + (*ListRoomsRequest)(nil), // 7: gospeak.livedemo.v1.ListRoomsRequest + (*ListRoomsResponse)(nil), // 8: gospeak.livedemo.v1.ListRoomsResponse + (*PublishRequest)(nil), // 9: gospeak.livedemo.v1.PublishRequest + (*PublishResponse)(nil), // 10: gospeak.livedemo.v1.PublishResponse + (*SubscribeRequest)(nil), // 11: gospeak.livedemo.v1.SubscribeRequest + (*SubscribeResponse)(nil), // 12: gospeak.livedemo.v1.SubscribeResponse + (*StopStreamRequest)(nil), // 13: gospeak.livedemo.v1.StopStreamRequest + (*StopStreamResponse)(nil), // 14: gospeak.livedemo.v1.StopStreamResponse + (*WatchRoomRequest)(nil), // 15: gospeak.livedemo.v1.WatchRoomRequest + (*RoomEvent)(nil), // 16: gospeak.livedemo.v1.RoomEvent +} +var file_live_sfu_proto_depIdxs = []int32{ + 0, // 0: gospeak.livedemo.v1.BackendInfo.kind:type_name -> gospeak.livedemo.v1.BackendKind + 3, // 1: gospeak.livedemo.v1.GetConfigResponse.backends:type_name -> gospeak.livedemo.v1.BackendInfo + 0, // 2: gospeak.livedemo.v1.StreamTarget.backend:type_name -> gospeak.livedemo.v1.BackendKind + 5, // 3: gospeak.livedemo.v1.Room.targets:type_name -> gospeak.livedemo.v1.StreamTarget + 6, // 4: gospeak.livedemo.v1.ListRoomsResponse.rooms:type_name -> gospeak.livedemo.v1.Room + 0, // 5: gospeak.livedemo.v1.PublishRequest.backend:type_name -> gospeak.livedemo.v1.BackendKind + 1, // 6: gospeak.livedemo.v1.PublishResponse.ice_servers:type_name -> gospeak.livedemo.v1.IceServer + 5, // 7: gospeak.livedemo.v1.PublishResponse.target:type_name -> gospeak.livedemo.v1.StreamTarget + 0, // 8: gospeak.livedemo.v1.SubscribeRequest.backend:type_name -> gospeak.livedemo.v1.BackendKind + 1, // 9: gospeak.livedemo.v1.SubscribeResponse.ice_servers:type_name -> gospeak.livedemo.v1.IceServer + 0, // 10: gospeak.livedemo.v1.StopStreamRequest.backend:type_name -> gospeak.livedemo.v1.BackendKind + 5, // 11: gospeak.livedemo.v1.RoomEvent.targets:type_name -> gospeak.livedemo.v1.StreamTarget + 2, // 12: gospeak.livedemo.v1.LiveSFU.GetConfig:input_type -> gospeak.livedemo.v1.GetConfigRequest + 7, // 13: gospeak.livedemo.v1.LiveSFU.ListRooms:input_type -> gospeak.livedemo.v1.ListRoomsRequest + 9, // 14: gospeak.livedemo.v1.LiveSFU.Publish:input_type -> gospeak.livedemo.v1.PublishRequest + 11, // 15: gospeak.livedemo.v1.LiveSFU.Subscribe:input_type -> gospeak.livedemo.v1.SubscribeRequest + 13, // 16: gospeak.livedemo.v1.LiveSFU.StopStream:input_type -> gospeak.livedemo.v1.StopStreamRequest + 15, // 17: gospeak.livedemo.v1.LiveSFU.WatchRoom:input_type -> gospeak.livedemo.v1.WatchRoomRequest + 4, // 18: gospeak.livedemo.v1.LiveSFU.GetConfig:output_type -> gospeak.livedemo.v1.GetConfigResponse + 8, // 19: gospeak.livedemo.v1.LiveSFU.ListRooms:output_type -> gospeak.livedemo.v1.ListRoomsResponse + 10, // 20: gospeak.livedemo.v1.LiveSFU.Publish:output_type -> gospeak.livedemo.v1.PublishResponse + 12, // 21: gospeak.livedemo.v1.LiveSFU.Subscribe:output_type -> gospeak.livedemo.v1.SubscribeResponse + 14, // 22: gospeak.livedemo.v1.LiveSFU.StopStream:output_type -> gospeak.livedemo.v1.StopStreamResponse + 16, // 23: gospeak.livedemo.v1.LiveSFU.WatchRoom:output_type -> gospeak.livedemo.v1.RoomEvent + 18, // [18:24] is the sub-list for method output_type + 12, // [12:18] is the sub-list for method input_type + 12, // [12:12] is the sub-list for extension type_name + 12, // [12:12] is the sub-list for extension extendee + 0, // [0:12] is the sub-list for field type_name +} + +func init() { file_live_sfu_proto_init() } +func file_live_sfu_proto_init() { + if File_live_sfu_proto != nil { + return + } + type x struct{} + out := protoimpl.TypeBuilder{ + File: protoimpl.DescBuilder{ + GoPackagePath: reflect.TypeOf(x{}).PkgPath(), + RawDescriptor: unsafe.Slice(unsafe.StringData(file_live_sfu_proto_rawDesc), len(file_live_sfu_proto_rawDesc)), + NumEnums: 1, + NumMessages: 16, + NumExtensions: 0, + NumServices: 1, + }, + GoTypes: file_live_sfu_proto_goTypes, + DependencyIndexes: file_live_sfu_proto_depIdxs, + EnumInfos: file_live_sfu_proto_enumTypes, + MessageInfos: file_live_sfu_proto_msgTypes, + }.Build() + File_live_sfu_proto = out.File + file_live_sfu_proto_goTypes = nil + file_live_sfu_proto_depIdxs = nil +} diff --git a/gen/live_sfu_grpc.pb.go b/gen/live_sfu_grpc.pb.go new file mode 100644 index 0000000..38c2644 --- /dev/null +++ b/gen/live_sfu_grpc.pb.go @@ -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", +} diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..42f8fe5 --- /dev/null +++ b/go.mod @@ -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 +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..98ff7bd --- /dev/null +++ b/go.sum @@ -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= diff --git a/internal/config/config.go b/internal/config/config.go new file mode 100644 index 0000000..0bfea40 --- /dev/null +++ b/internal/config/config.go @@ -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 +} diff --git a/internal/server/gateway.go b/internal/server/gateway.go new file mode 100644 index 0000000..23d8c04 --- /dev/null +++ b/internal/server/gateway.go @@ -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 +} diff --git a/internal/server/grpc.go b/internal/server/grpc.go new file mode 100644 index 0000000..cff6814 --- /dev/null +++ b/internal/server/grpc.go @@ -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) +} diff --git a/internal/server/grpc_test.go b/internal/server/grpc_test.go new file mode 100644 index 0000000..98fc726 --- /dev/null +++ b/internal/server/grpc_test.go @@ -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) + } +} diff --git a/internal/server/proxy.go b/internal/server/proxy.go new file mode 100644 index 0000000..a666c9a --- /dev/null +++ b/internal/server/proxy.go @@ -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 +} diff --git a/internal/server/rooms.go b/internal/server/rooms.go new file mode 100644 index 0000000..f276b30 --- /dev/null +++ b/internal/server/rooms.go @@ -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: + } + } +} diff --git a/internal/server/server.go b/internal/server/server.go new file mode 100644 index 0000000..e04165d --- /dev/null +++ b/internal/server/server.go @@ -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)) + } +} diff --git a/internal/server/service.go b/internal/server/service.go new file mode 100644 index 0000000..e188946 --- /dev/null +++ b/internal/server/service.go @@ -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 +} diff --git a/internal/server/static/app.js b/internal/server/static/app.js new file mode 100644 index 0000000..a2713ee --- /dev/null +++ b/internal/server/static/app.js @@ -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 = '控制面不可达'; + } + } + + 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) => '分发中:' + shortName(ENUM[k]) + "").join("") + : '未发布'; + } + + // ---------- 观看页 ---------- + 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 = '观看中:' + shortName(ENUM[b]) + ""; + } 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 + ? '房间在 ' + shortName(ENUM[backend()]) + " 已直播" + : '等待 ' + shortName(ENUM[backend()]) + " 推流…"; + 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 = '未观看'; + log(logEl, "已停止观看"); + }; + } + + return { ENUM, getRoom, loadConfig, log, initPublish, initWatch, apiGet, apiPost }; +})(); diff --git a/internal/server/static/index.html b/internal/server/static/index.html new file mode 100644 index 0000000..ef942df --- /dev/null +++ b/internal/server/static/index.html @@ -0,0 +1,54 @@ + + + + + + 直播 SFU 分流 Demo + + + +
+
直播 SFU 分流 Demo Cloudflare Realtime · SRS
+ +
+
+
+

一路推流,多路扇出

+

基于 GOSpeak 技术栈的直播 SFU 分流演示。控制面用 protobuf/gRPC 定义,主 SFU 为 + Cloudflare Realtime,SRS 作为本地对照后端。同一路发布可同时落到多条分发链路,观众任选后端拉流。

+
+
+
+

进入房间

+
+ + + +
+
提示:推流与观看使用同一房间名即可配对。
+
+
+

后端能力

+
加载中…
+
Cloudflare Realtime 需在 .env 配置 CF_APP_ID / CF_APP_SECRET; + SRS 由本地 docker compose up -d srs 提供。
+
+
+
+ + + + diff --git a/internal/server/static/publish.html b/internal/server/static/publish.html new file mode 100644 index 0000000..30edde9 --- /dev/null +++ b/internal/server/static/publish.html @@ -0,0 +1,43 @@ + + + + + + 开始直播 · 直播 SFU 分流 Demo + + + +
+
直播 SFU 分流 Demo Cloudflare Realtime · SRS
+ +
+
+

开始直播

房间 demo · 选择分发后端,一路推流将扇出到所选 SFU。

+
+
+

本地预览

+ +
+ +
+
+ +
+
+ + +
+
未发布
+
+
+

分发状态

+
尚未开始
+

日志

+
+
+
+
+ + + + diff --git a/internal/server/static/styles.css b/internal/server/static/styles.css new file mode 100644 index 0000000..4f68064 --- /dev/null +++ b/internal/server/static/styles.css @@ -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; } } diff --git a/internal/server/static/watch.html b/internal/server/static/watch.html new file mode 100644 index 0000000..e5a1e12 --- /dev/null +++ b/internal/server/static/watch.html @@ -0,0 +1,41 @@ + + + + + + 观看直播 · 直播 SFU 分流 Demo + + + +
+
直播 SFU 分流 Demo Cloudflare Realtime · SRS
+ +
+
+

观看直播

房间 demo · 选择从哪个 SFU 拉流,房间一旦直播即自动连接。

+
+
+

远端画面

+ +
+ +
+
+ +
+
+ + +
+
未观看
+
+
+

日志

+
+
+
+
+ + + + diff --git a/internal/server/token.go b/internal/server/token.go new file mode 100644 index 0000000..6116c04 --- /dev/null +++ b/internal/server/token.go @@ -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) +} diff --git a/internal/server/token_test.go b/internal/server/token_test.go new file mode 100644 index 0000000..d963c7a --- /dev/null +++ b/internal/server/token_test.go @@ -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") + } +} diff --git a/internal/sfu/cloudflare/client.go b/internal/sfu/cloudflare/client.go new file mode 100644 index 0000000..fbc204d --- /dev/null +++ b/internal/sfu/cloudflare/client.go @@ -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 +} diff --git a/internal/sfu/cloudflare/client_test.go b/internal/sfu/cloudflare/client_test.go new file mode 100644 index 0000000..ddbc83c --- /dev/null +++ b/internal/sfu/cloudflare/client_test.go @@ -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) + } +} diff --git a/internal/sfu/cloudflare/provider.go b/internal/sfu/cloudflare/provider.go new file mode 100644 index 0000000..4f9a740 --- /dev/null +++ b/internal/sfu/cloudflare/provider.go @@ -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, + } +} diff --git a/internal/sfu/srs/provider.go b/internal/sfu/srs/provider.go new file mode 100644 index 0000000..8b6edb5 --- /dev/null +++ b/internal/sfu/srs/provider.go @@ -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, + } +}