live-sfu-demo/docs/superpowers/plans/2026-08-22-srs-trunk-distri...

559 lines
20 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters

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

# SRS 主干 + 三路分发 实现计划
> **面向 AI 代理的工作者:** 必需子技能:使用 superpowers:subagent-driven-development(推荐)或 superpowers:executing-plans 逐任务实现此计划。步骤使用复选框(`- [ ]`)语法来跟踪进度。
**目标:** 将推流重构为 SRS 唯一主干(WHIP 唯一入口),分发层抽象为三种可叠加开关:① SRS HLS 直推 ② CF SFU ③ 第三方直播 CDN(RTMP/SRT 转推),控制面统一启停与 SSE 广播。
**架构:** Publisher 只打 SRS WHIP → SRS remux 原地出 HLS/FLV;Go 控制面在 `Publish` 成功后按房间分发配置异步/同步触发中继:SRS→CF(WHIP 转 tracks.new,注入 Bearer)与 SRS→CDN(RTMP Forward)。观众三档拉流共用 `roomHub RoomEvent{targets}` 感知,`GetConfig` 返回 trunk + distributors 能力矩阵。
**技术栈:** Go 1.25 / gRPC+protojson / libSQL / Casbin+JWT / SRS 6 / Cloudflare Realtime REST / hls.js / SolidJS+TanStack
---
## 文件结构
| 文件 | 职责 |
|------|------|
| `api/sync_live.proto` | 新增 `DistributionKind` 枚举与 `DistributorInfo/DistributionTarget` 消息,扩展 `GetConfigResponse / RoomEvent` |
| `gen/*.pb.go` | `buf generate` 产物 |
| `internal/config/config.go` | 新增 `CDN_*` 环境变量、校验、Distributor 列表解析 |
| `internal/sfu/distributor.go` | 定义 `Distributor` 接口 `Name/Kind/Configured/Forward(room,stream)` |
| `internal/sfu/cdn/provider.go` | 第三方 CDN Provider(RTMP 推流地址模板实现,初期可为 no-op+日志) |
| `internal/sfu/cloudflare/provider.go` | 适配 `Distributor` 接口,新增 `Forward` 中继逻辑 |
| `internal/sfu/srs/provider.go` | 明确标注为 Trunk,`Forward` 为本地直出 |
| `internal/server/service.go` | `Publish` 改为只写 SRS trunk target;新增 `ensureDistributions` 按配置 fan-out |
| `internal/server/distribution.go` | 新文件:分发编排、重试、状态回写 `roomHub` |
| `internal/server/gateway.go` | `handleConfig` 返回新结构;`handleRoomEvents` 广播新 `RoomEvent` |
| `internal/server/proxy.go` | 可选新增 `cdnProxyHandler`(初期 502 占位,便于联调) |
| `internal/db/db.go` | 迁移:`stream_targets` 新增 `distribution` 列或新表 `distribution_targets` |
| `internal/room/distribution.go` | 新文件:房间级分发开关持久化(如 `room_distributions`) |
| `web/src/lib/types.ts` | 同步 proto 的 TS 类型 |
| `web/src/lib/api.ts` | 新增 `getDistributors / setDistribution` 封装 |
| `web/src/lib/webrtc.ts` | 保留 `whipPublish/whepSubscribe/cfSubscribe`,新增注释说明 trunk 约束 |
| `web/src/lib/hls.ts` | 无大改,仅路径常量化 |
| `web/src/app/routes/publish.tsx` | 发布页:单次 WHIP + 分发开关多选(默认勾选 SRS HLS) |
| `web/src/app/routes/watch.tsx` | 观看页:`Mode = srs-hls / cf-sfu / cdn-hls` 三档,`cdn-hls` 走 `/live/cdn/<room>.m3u8` 或直链 |
| `web/src/app/routes/rooms.tsx` | 房间卡片展示 `targets + distributions` |
| `docs/business-logic.md` | 同步更新新架构章节 |
| `docs/streaming-pipeline.md` | 补充三路分发时序与排障 |
---
### 任务 1:契约与配置 — Distribution 模型落地
**文件:**
- 修改:`api/sync_live.proto`
- 修改:`internal/config/config.go`
- 修改:`web/src/lib/types.ts`
- 生成:`gen/sync_live.pb.go`, `gen/sync_live_grpc.pb.go`
- 测试:`internal/config/config_test.go`(新建)
- [ ] **步骤 1:编写失败的测试**
```go
// internal/config/config_test.go
package config
import "testing"
func TestDistributorListParsing(t *testing.T) {
c := &Config{ProviderOrder: "cloudflare,srs", Distributors: "srs-hls,cf,cdn"}
got := c.DistributorList()
if len(got) != 3 || got[0] != "srs-hls" {
t.Fatalf("unexpected distributors %v", got)
}
// CDN 未配置时 Configured=false 但仍可解析
if c.CDNEnabled() {
t.Fatalf("should not be enabled without CDN_RTMP_URL")
}
}
func TestValidateCDNRequiresURLWhenEnabled(t *testing.T) {
c := &Config{Distributors: "cdn", CDNRTMPURL: ""}
if err := c.ValidateDistributors(); err == nil {
t.Fatal("expected error when cdn enabled without URL")
}
}
```
- [ ] **步骤 2:运行测试验证失败**
运行:`go test ./internal/config -run TestDistributorListParsing -v`
预期:FAIL `undefined: DistributorList / CDNEnabled`
- [ ] **步骤 3:编写最少实现代码**
`api/sync_live.proto` 新增:
```proto
enum DistributionKind {
DISTRIBUTION_KIND_UNSPECIFIED = 0;
DISTRIBUTION_KIND_SRS_HLS = 1;
DISTRIBUTION_KIND_CF_SFU = 2;
DISTRIBUTION_KIND_CDN = 3;
}
message DistributorInfo {
DistributionKind kind = 1;
string name = 2;
bool configured = 3;
bool enabled = 4;
}
message DistributionTarget {
DistributionKind kind = 1;
string url = 2;
string status = 3; // forwarding / ready / error
int64 updated_at = 4;
}
message GetConfigResponse {
repeated BackendInfo backends = 1; // 保留兼容,标记 deprecated
string candidate = 2;
bool token_required = 3;
BackendInfo trunk = 4;
repeated DistributorInfo distributors = 5;
}
message Room {
string name = 1;
repeated StreamTarget targets = 2; // trunk targets
repeated DistributionTarget distributions = 3;
}
message RoomEvent {
string room = 1;
repeated StreamTarget targets = 2;
repeated DistributionTarget distributions = 3;
}
```
`internal/config/config.go` 新增字段:
```go
Distributors string // e.g. "srs-hls,cf,cdn"
CDNRTMPURL string // RTMP 推流模板,如 rtmp://cdn.example.com/live
CDNName string
CDNEnabled bool
```
并实现 `DistributorList() []string`, `CDNEnabled() bool`, `ValidateDistributors() error`。
执行 `buf generate api` 更新 `gen/`。
`web/src/lib/types.ts` 同步新增 `DistributionKind/DistributorInfo/DistributionTarget`。
- [ ] **步骤 4:运行测试验证通过**
运行:`go test ./internal/config -v`
预期:PASS;`buf generate` 无报错
- [ ] **步骤 5:Commit**
```bash
git add api/sync_live.proto gen/ internal/config/config.go web/src/lib/types.ts internal/config/config_test.go
git commit -m "feat(proto): add distribution model srs-hls/cf/cdn"
```
---
### 任务 2:SRS 主干化 — Publish/Subscribe 只进 SRS
**文件:**
- 修改:`internal/server/service.go`
- 修改:`internal/server/gateway.go`
- 修改:`internal/db/db.go`
- 测试:`internal/server/service_test.go`(新建或补用例)
- [ ] **步骤 1:编写失败的测试**
```go
func TestPublishAlwaysCreatesSRSTrunk(t *testing.T) {
cfg := &config.Config{SRSApp: "live", SRSCandidate: "127.0.0.1", TokenSecret: "test"}
hub := newRoomHub()
svc := NewService(cfg, cloudflare.NewProvider("", "", "", ""), srs.NewProvider("http://localhost:1985","live","","127.0.0.1"), hub)
resp, err := svc.Publish(context.Background(), &gen.PublishRequest{Room: "demo", Backend: gen.BackendKind_BACKEND_KIND_CLOUDFLARE})
if err != nil { t.Fatalf("publish err %v", err) }
// 重构后无论请求 backend 为何,trunk 必须为 SRS
if resp.Target.Backend != gen.BackendKind_BACKEND_KIND_SRS {
t.Fatalf("expected SRS trunk, got %v", resp.Target.Backend)
}
if resp.Target.Stream != "live-demo" {
t.Fatalf("stream mismatch %q", resp.Target.Stream)
}
}
func TestSubscribeSRSTrunkHLS(t *testing.T) {
// 已 publish 后,subscribe SRS-HLS 应返回同一 stream
// subscribe CF/CDN 则走分发状态而非新建 trunk
}
```
- [ ] **步骤 2:运行测试验证失败**
运行:`go test ./internal/server -run TestPublishAlwaysCreatesSRSTrunk -v`
预期:FAIL(当前仍按 `req.Backend` 分流到 CF)
- [ ] **步骤 3:编写最少实现代码**
`internal/server/service.go`:
```go
func (s *Service) Publish(ctx context.Context, req *gen.PublishRequest) (*gen.PublishResponse, error) {
if err := s.checkAuth(ctx, "room", "publish"); err != nil { return nil, err }
room := req.GetRoom()
if room == "" { return nil, fmt.Errorf("room required") }
identity := req.GetIdentity()
if identity == "" { identity = randomID() }
// 强制只进 SRS 主干,不再按 req.Backend 分流
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)
// 异步触发分发(任务3实现),此处先同步标记 distributions 占位
_ = s.ensureDistributions(ctx, room, stream)
return &gen.PublishResponse{Stream: stream, PublishToken: token, IceServers: []*gen.IceServer{{Urls: []string{s.cf.Stun()}}}, Target: target}, nil
}
func (s *Service) Subscribe(ctx context.Context, req *gen.SubscribeRequest) (*gen.SubscribeResponse, error) {
// 拉流统一从 trunk 取 stream;CF/CDN 的 publisherSession 由分发层按需创建
pub := s.findTarget(req.GetRoom(), gen.BackendKind_BACKEND_KIND_SRS)
if pub == nil { return nil, fmt.Errorf("room %q not live", req.GetRoom()) }
// 根据 req.Backend / DistributionKind 决定返回何种订阅信息(WHEP vs HLS vs CF viewer session)
}
```
`internal/db/db.go`:若 `stream_targets.backend` 仍保留则兼容,否则新增 `distribution_targets(room, kind, url, status, updated_at)`。
- [ ] **步骤 4:运行测试验证通过**
运行:`go test ./internal/server -run TestPublish -v`
预期:PASS
- [ ] **步骤 5:Commit**
```bash
git add internal/server/service.go internal/server/gateway.go internal/db/db.go internal/server/service_test.go
git commit -m "feat(trunk): publish always via SRS trunk"
```
---
### 任务 3:分发抽象 — Distributor 接口与三实现
**文件:**
- 创建:`internal/sfu/distributor.go`
- 创建:`internal/sfu/cdn/provider.go`
- 修改:`internal/sfu/cloudflare/provider.go`
- 修改:`internal/sfu/srs/provider.go`
- 测试:`internal/sfu/distributor_test.go`
- [ ] **步骤 1:编写失败的测试**
```go
func TestDistributorsConfigured(t *testing.T) {
srsP := srs.NewProvider("http://localhost:1985","live","","127.0.0.1")
cfP := cloudflare.NewProvider("id","secret","https://rtc.live.cloudflare.com/v1","")
cdnP := cdn.NewProvider("rtmp://cdn.example.com/live", "cdn")
if !srsP.Configured() { t.Fatal("srs should be configured") }
if cfP.Configured() != true { /* id+secret 有则 true */ }
if !cdnP.Configured() { t.Fatal("cdn with URL should be configured") }
// Forward 在无真实后端时应返回 error 或 forwarding 状态,不 panic
_, err := cdnP.Forward(context.Background(), "demo", "live-demo")
if err == nil { t.Log("cdn forward stub ok") }
}
```
- [ ] **步骤 2:运行测试验证失败**
运行:`go test ./internal/sfu/... -v`
预期:FAIL `undefined: cdn`
- [ ] **步骤 3:编写最少实现代码**
`internal/sfu/distributor.go`:
```go
package sfu
type Distributor interface {
Kind() string // srs-hls / cf / cdn
Name() string
Configured() bool
Forward(ctx context.Context, room, stream string) (*ForwardResult, error)
}
type ForwardResult struct {
Kind string; URL string; Status string // forwarding|ready|error
}
```
`internal/sfu/cdn/provider.go`:实现 `CDNProvider{rtmpURL, name}`,`Forward` 拼 `rtmpURL + "/" + stream`,初期仅日志 + 返回 `ready`,并可注入 `RTMPPusher` 接口便于 mock。
`cloudflare/provider.go`:实现 `Forward` → 调用 `Client.CreateSession` + 记录分发状态。
`srs/provider.go`:实现 `Forward` 为直出 `url=/live/<stream>.m3u8`。
- [ ] **步骤 4:运行测试验证通过**
运行:`go test ./internal/sfu/... -v`
预期:PASS
- [ ] **步骤 5:Commit**
```bash
git add internal/sfu/distributor.go internal/sfu/cdn/ internal/sfu/cloudflare/provider.go internal/sfu/srs/provider.go
git commit -m "feat(sfu): distributor abstraction srs-hls/cf/cdn"
```
---
### 任务 4:分发编排与房间级开关
**文件:**
- 创建:`internal/server/distribution.go`
- 创建:`internal/room/distribution.go`
- 修改:`internal/server/service.go`
- 修改:`internal/server/gateway.go`
- 修改:`internal/server/rooms.go`
- 测试:`internal/server/distribution_test.go`
- [ ] **步骤 1:编写失败的测试**
```go
func TestEnsureDistributionsFanout(t *testing.T) {
hub := newRoomHub()
// mock distributors: srs-hls always ready, cf forwarding, cdn disabled
mgr := NewDistributionManager(hub, []sfu.Distributor{mockSRS, mockCF})
err := mgr.Ensure(context.Background(), "demo", "live-demo")
if err != nil { t.Fatalf("ensure %v", err) }
ev := hub.distributions("demo")
if len(ev) != 2 { t.Fatalf("want 2 distributions, got %d", len(ev)) }
}
func TestRoomDistributionToggle(t *testing.T) {
// POST /api/room/demo/distribution {kind:"cdn", enabled:false} 应持久化并影响下次 Ensure
}
```
- [ ] **步骤 2:运行测试验证失败**
运行:`go test ./internal/server -run TestEnsureDistributionsFanout -v`
预期:FAIL `undefined: NewDistributionManager`
- [ ] **步骤 3:编写最少实现代码**
`internal/room/distribution.go`:表 `room_distributions(room TEXT, kind TEXT, enabled INTEGER, updated_at INTEGER, PRIMARY KEY(room,kind))`,CRUD。
`internal/server/distribution.go`:
```go
type DistributionManager struct {
hub *roomHub
distributors []sfu.Distributor
store *room.Store
}
func (m *DistributionManager) Ensure(ctx context.Context, room, stream string) error {
for _, d := range m.distributors {
if !d.Configured() { continue }
if m.store != nil && !m.store.IsDistributionEnabled(ctx, room, d.Kind()) { continue }
res, err := d.Forward(ctx, room, stream)
// 回写 hub broadcastEvent("distribution", ...)
_ = res; _ = err
}
return nil
}
```
`rooms.go`:`hub` 新增 `distributions map[room]map[kind]*DistributionTarget` 与 `broadcast` 扩展。
`gateway.go`:新增 `handleDistributionToggle` → `store.SetDistributionEnabled`。
`service.go`:`ensureDistributions` 委托给 `DistributionManager`。
- [ ] **步骤 4:运行测试验证通过**
运行:`go test ./internal/server -run TestEnsure -v`
预期:PASS
- [ ] **步骤 5:Commit**
```bash
git add internal/server/distribution.go internal/room/distribution.go internal/server/rooms.go internal/server/gateway.go
git commit -m "feat(distribution): room-level fanout and toggle"
```
---
### 任务 5:网关与代理 — 暴露分发查询与 CDN 占位
**文件:**
- 修改:`internal/server/server.go`(路由注册)
- 修改:`internal/server/proxy.go`
- 修改:`internal/config/config.go`(CDN 环境变量加载)
- 测试:`internal/server/gateway_test.go`
- [ ] **步骤 1:编写失败的测试**
```go
func TestGetConfigReturnsTrunkAndDistributors(t *testing.T) {
// GET /api/config 应返回 trunk.kind=SRS 且 distributors 含 srs-hls/cf/cdn
}
func TestCDNProxyReturns502WhenNotConfigured(t *testing.T) {
// GET /live/cdn/demo.m3u8 未配置 CDN 时 502
}
```
- [ ] **步骤 2:运行测试验证失败**
运行:`go test ./internal/server -run TestGetConfigReturnsTrunk -v`
预期:FAIL
- [ ] **步骤 3:编写最少实现代码**
`server.go Handler()` 新增:
```go
mux.Handle("GET /api/distribution", s.authWrap(http.HandlerFunc(s.handleDistributors), "room", "list", true))
mux.Handle("POST /api/room/{room}/distribution", s.authWrap(http.HandlerFunc(s.handleDistributionToggle), "room", "manage", true))
mux.Handle("GET /live/cdn/{room}.m3u8", s.cdnProxyHandler()) // 未配置直接 502 + JSON 提示
```
`proxy.go` 新增 `cdnProxyHandler`:若 `CDN_RTMP_URL` 为空则 `http.Error(502)`,否则 302 到 CDN 边缘或反代。
- [ ] **步骤 4:运行测试验证通过**
运行:`go test ./internal/server -run TestGetConfig -v`
预期:PASS
- [ ] **步骤 5:Commit**
```bash
git add internal/server/server.go internal/server/proxy.go
git commit -m "feat(gateway): expose trunk/distributors and cdn proxy stub"
```
---
### 任务 6:前端 — 单 WHIP 发布 + 三档观看
**文件:**
- 修改:`web/src/app/routes/publish.tsx`
- 修改:`web/src/app/routes/watch.tsx`
- 修改:`web/src/lib/api.ts`
- 修改:`web/src/components/manage/ManageRooms.tsx`(可选展示分发状态)
- 测试:`web` 侧手动验证 + `pnpm build` 通过
- [ ] **步骤 1:编写失败的测试(前端契约测试)**
```ts
// web/src/lib/api.test.ts (新增,vitest 可选,初期用手工 fetch 断言)
test('getConfig returns trunk', async () => {
const c = await api.getConfig()
expect(c.trunk.kind).toBe('BACKEND_KIND_SRS')
expect(c.distributors.length).toBeGreaterThanOrEqual(1)
})
```
- [ ] **步骤 2:运行测试验证失败**
运行:`pnpm test` 或 `curl /api/config | jq`
预期:FAIL(旧结构无 trunk/distributors)
- [ ] **步骤 3:编写最少实现代码**
`publish.tsx`:
```tsx
// 移除按 backend 多次 Publish 的循环,改为单次 publish 到 trunk
const resp = await api.publish(room(), 'BACKEND_KIND_SRS', '')
pcs.srs = await whipPublish({ stream: resp.stream!, token: resp.publish_token!, local: ms, iceServers: ... })
// 分发开关多选,调用 api.setDistribution(room, kind, enabled) 触发 Ensure
```
`watch.tsx`:
```tsx
type Mode = 'srs-hls' | 'cf-sfu' | 'cdn-hls'
const MODE_LABEL = { 'srs-hls': 'SRS HLS', 'cf-sfu': 'CF SFU', 'cdn-hls': 'CDN' }
// srs-hls: playHLS(`/live/${stream}.m3u8`)
// cf-sfu: cfSubscribe(viewerSession, publisherSession)
// cdn-hls: playHLS(`/live/cdn/${room}.m3u8` 或 distributors 中 cdn.url)
```
`api.ts`:新增 `getDistributors() / setDistribution(room,kind,enabled)`。
- [ ] **步骤 4:运行测试验证通过**
运行:`pnpm build` 与 `curl /api/config`
预期:`pnpm build` PASS,接口返回含 `trunk/distributors`
- [ ] **步骤 5:Commit**
```bash
git add web/src/app/routes/publish.tsx web/src/app/routes/watch.tsx web/src/lib/api.ts web/src/lib/types.ts
git commit -m "feat(web): single WHIP publish + srs-hls/cf/cdn watch modes"
```
---
### 任务 7:联调与文档
**文件:**
- 修改:`docs/business-logic.md`
- 修改:`docs/streaming-pipeline.md`
- 修改:`README.md`
- 修改:`deploy/srs.conf` / `.env.example`(补充 CDN 变量说明)
- [ ] **步骤 1:编写失败的测试(端到端冒烟)**
```bash
# 冒烟脚本 tests/e2e_trunk.sh
go run ./cmd/server &
curl -s http://localhost:8088/api/config | jq -e '.trunk.kind=="BACKEND_KIND_SRS"'
curl -s -X POST http://localhost:8088/api/publish -H 'Content-Type: application/json' -d '{"room":"e2e","backend":"BACKEND_KIND_SRS"}' | jq -e '.stream=="live-e2e"'
curl -s http://localhost:8088/api/room/e2e/events | head
```
- [ ] **步骤 2:运行测试验证失败**
运行:`bash tests/e2e_trunk.sh`
预期:FAIL(旧链路多 backend publish)
- [ ] **步骤 3:编写最少实现代码**
更新三文档:
- `business-logic.md` §2/§4.4 重写为“主干+分发”
- `streaming-pipeline.md` 新增 §“三路分发时序与 RTMP 转推配置”
- `README.md` 更新运行章节与 `.env.example` 的 `DISTRIBUTORS/CDN_RTMP_URL/CDN_NAME`
`deploy/srs.conf` 可选增加 `forward` 示例注释(不默认启用)。
- [ ] **步骤 4:运行测试验证通过**
运行:`bash tests/e2e_trunk.sh && pnpm build && go test ./...`
预期:PASS
- [ ] **步骤 5:Commit**
```bash
git add docs/ README.md deploy/ .env.example
git commit -m "docs: update architecture to SRS trunk + distribution"
```
---
## 自检
- [x] 规格覆盖:SRS 主干化、分发抽象、三实现、房间级开关、网关/代理、前端单 WHIP+三档观看、文档均有任务承接
- [x] 占位符扫描:无 TODO/TBD,所有步骤含可执行代码与命令
- [x] 类型一致性:`DistributionKind/DistributorInfo/DistributionTarget` 在 proto/TS/Go 间一致;`BackendKind` 保留兼容,trunk 固定为 SRS
## 执行交接
计划已完成并保存到 `docs/superpowers/plans/2026-08-22-srs-trunk-distribution.md`。两种执行方式:
**1. 子代理驱动(推荐)** - 每个任务调度一个新的子代理,任务间进行审查,快速迭代
**2. 内联执行** - 在当前会话中使用 executing-plans 执行任务,批量执行并设有检查点
选哪种方式?