chore: commit remaining worktree changes (room/playlist/docs/deps)
This commit is contained in:
parent
a3a6a3a77a
commit
3256839d5d
|
|
@ -0,0 +1,464 @@
|
|||
# SyncLive 业务逻辑文档
|
||||
|
||||
> 版本:2026-08-22 · 基于 `main` 分支代码整理 · 对标 SyncTV 的直播分流姐妹项目
|
||||
|
||||
---
|
||||
|
||||
## 1. 项目定位
|
||||
|
||||
**SyncLive = Sync + Live(一起播·一起看直播)**,对标 `SyncTV(Sync + TV 一起看)`。
|
||||
|
||||
- **一句话**:一路推流,多端同步扇出。控制面 `protobuf/gRPC`,媒体面 `WebRTC + HLS/FLV`。
|
||||
- **双 SFU 分流**:同一路发布可同时落到 `Cloudflare Realtime(主 SFU,托管扇出)` 与 `SRS(本地对照,WHIP/WHEP/HLS/FLV 三协议同出)`。
|
||||
- **设计约束**:浏览器只与服务端交换 SDP/信令,SFU 凭证不出服务端;房间拓扑持久化到嵌入式 libSQL,无外部依赖即可开箱。
|
||||
|
||||
```
|
||||
publisher ──WHIP/tracks.new──▶ Cloudflare Realtime ──▶ viewer(s)
|
||||
└────WHIP──────────▶ SRS ─┬─▶ WHEP (0.2-0.5s 低延时)
|
||||
├─▶ HLS (5-10s,切片 /live/*.m3u8)
|
||||
└─▶ FLV (http-flv)
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 2. 总体架构
|
||||
|
||||
```
|
||||
浏览器 (SolidJS SPA, TanStack Router/Query, Ark UI, Tailwind v4)
|
||||
│ HTTP / SSE / WebRTC(信令经服务端反代,媒体直连 SFU)
|
||||
▼
|
||||
控制面 (Go, :8088 HTTP + :9090 gRPC)
|
||||
├─ gRPC + JSON 网关 (protojson,见 api/sync_live.proto)
|
||||
├─ 鉴权 (Casbin RBAC + JWT + UserStore)
|
||||
├─ 房间拓扑 (roomHub 内存 + libSQL 持久化 stream_targets)
|
||||
├─ 房间成员/权限/播放列表 (room.Service + room.Store)
|
||||
├─ 聊天弹幕 (chatHub 内存广播 + 滑动窗口限流)
|
||||
└─ 媒体面反代 (httputil: /rtc/v1/* → SRS, /api/cf/* → Cloudflare, /live/* → SRS HLS)
|
||||
▼
|
||||
媒体面 (SFU Provider 抽象)
|
||||
├─ Cloudflare Realtime Provider (REST Client + BackendInfo)
|
||||
└─ SRS Provider (WHIP/WHEP 信令 + HLS 静态资源)
|
||||
▼
|
||||
数据层
|
||||
├─ libSQL 嵌入式 file:./data/sync-live.db (WAL)
|
||||
└─ data/users.json (用户库,bcrypt 哈希)
|
||||
```
|
||||
|
||||
目录映射:
|
||||
|
||||
| 目录 | 职责 |
|
||||
|------|------|
|
||||
| `api/sync_live.proto` | 控制面唯一契约,经 `buf generate` 到 `gen/` |
|
||||
| `internal/config` | 环境变量加载 + `Validate()` 密钥校验 + `ProviderList()` |
|
||||
| `internal/db` | 嵌入式 libSQL 打开/迁移/CRUD (stream_targets, rooms_v2 等) |
|
||||
| `internal/auth` | JWT 双令牌、用户文件存储、Casbin Enforcer、中间件 |
|
||||
| `internal/room` | 房间/成员/角色/权限/播放列表的领域模型与存储 |
|
||||
| `internal/sfu/cloudflare` | Cloudflare Realtime REST 客户端 |
|
||||
| `internal/sfu/srs` | SRS Provider 元信息 |
|
||||
| `internal/server` | HTTP/gRPC 装配、网关、反代、SSE、聊天、Manage 面 |
|
||||
| `cmd/server` | 进程入口与优雅关闭 |
|
||||
| `web/src` | SolidJS 前端,构建产物 `internal/server/static` 被 Go embed |
|
||||
| `deploy` | `docker-compose.yml` + `srs.conf` |
|
||||
|
||||
---
|
||||
|
||||
## 3. 技术栈
|
||||
|
||||
- **后端**:Go 1.25, `google.golang.org/grpc + protobuf`, `casbin/casbin/v2`, `golang-jwt/jwt/v5`, `tursodatabase/go-libsql`, `golang.org/x/crypto/bcrypt`
|
||||
- **前端**:SolidJS + TanStack Router/Query/Virtual + Ark UI + Tailwind v4 + Vite + `hls.js` + `lucide-solid` + `playwright` (截图)
|
||||
- **SFU**:Cloudflare Realtime (`rtc.live.cloudflare.com/v1`), SRS 6 (`:1935/:1985/:8080/:8000`)
|
||||
- **持久化**:libSQL file 模式,`MaxOpenConns=1`,WAL
|
||||
|
||||
---
|
||||
|
||||
## 4. 核心业务模块
|
||||
|
||||
### 4.1 房间 (Room)
|
||||
|
||||
**模型** `internal/room/room.go`:
|
||||
|
||||
| 字段 | 含义 |
|
||||
|------|------|
|
||||
| `name` | 主键,URL `?room=` 指定的房间名 |
|
||||
| `display_name/description` | 展示用 |
|
||||
| `creator` | 创建者 username |
|
||||
| `status` | `active / closed / banned`,状态机 `CanTransitionTo` |
|
||||
| `is_public` | 是否公开可发现 |
|
||||
| `password_hash` | 可选房间密码 (bcrypt) |
|
||||
| `require_approval` | 入房需审批 → 成员 `pending` |
|
||||
| `max_members` | 上限默认 100 |
|
||||
| `settings_json` | `RoomSettings` JSON,含细粒度权限覆盖 |
|
||||
| `auto_play` | `enabled/mode/delay`,`sequential/repeat_one/repeat_all/shuffle` |
|
||||
| `category_id/cover_file_id/label_ids` | 分类/标签/封面 (兼容 SyncTV) |
|
||||
|
||||
**房间服务** `internal/room/service.go`:
|
||||
|
||||
- `CreateRoom`:建 `rooms_v2` + 写入 `creator` 为 `RoleCreator` 的成员记录
|
||||
- `JoinRoom`:校验 `banned/closed`、密码、人数上限、重复加入、封禁;`require_approval=true` 时置 `MemberPending`
|
||||
- `EnsurePermission/CheckPermission`:按 `RoomSettings` 与成员 `PermissionSet` 判定单项权限;`closed/banned` 房间仅允许 `view_members/view_chat_history`
|
||||
- `UpdateMemberRole/Kick/Ban/Leave/ApproveJoin/RejectJoin/TransferOwnership`:均校验 `PermManageMembers/PermRemoveMembers` 与 `Role.CanManage` 等级
|
||||
|
||||
**成员模型** `internal/room/member.go` + `role.go`:
|
||||
|
||||
- 角色 `RoleCreator(1) > RoleAdmin(2) > RoleMember(3) > RoleGuest(4)`,`Rank()` 分 100/80/50/10
|
||||
- 状态 `MemberActive / Banned / Pending / Kicked`
|
||||
- 权限位 `PermissionSet uint64`,在 `added/removed/admin_added/admin_removed` 四列上做增量覆盖
|
||||
|
||||
### 4.2 鉴权与 RBAC (Auth)
|
||||
|
||||
**用户存储** `internal/auth/store.go` + `data/users.json`:
|
||||
|
||||
- 文件锁 + 内存 map,`Create` 时 bcrypt 哈希、`Verify` 比对
|
||||
- `BootstrapAdmin`:用户库为空且 `BOOTSTRAP_ADMIN_USER/PASS` 已设时创建首个 `admin`(可选 `BOOTSTRAP_ADMIN_ROLE=root` 提升为 super admin)
|
||||
- 字段:`username / hash / role / status(active|banned) / created_at`
|
||||
|
||||
**JWT** `internal/auth/jwt.go`:
|
||||
|
||||
- `JWTManager{secret, ttl=2h, refreshTTL=7d, issuer="sync-live"}`,HS256
|
||||
- `Sign(username, role)` → `access token (TokenType=access)`;`SignRefresh` → `refresh token`
|
||||
- `Verify / VerifyRefresh / IssuePair / Refresh`;`Validate()` 要求 `JWT_SECRET` 非空,否则拒绝启动
|
||||
|
||||
**Casbin** `internal/auth/model.conf + policy.csv + enforcer.go`:
|
||||
|
||||
- 模型:`r=sub,obj,act | p=sub,obj,act | g=_,_ | e=some(p.eft==allow) | m=g(r.sub,p.sub) && r.obj==p.obj && r.act==p.act`
|
||||
- 策略表(节选):
|
||||
- `root/admin/publisher/viewer/guest` 五档,`guest` 仅 `config:read + srs:streams + room:chat`
|
||||
- `viewer` 可 `room:subscribe/watch/chat`
|
||||
- `publisher` 可 `room:publish/stop`
|
||||
- `admin/root` 可 `room:manage + user:list/manage + system:manage`
|
||||
- `Enforcer.AddUserRole` 在登录/改角色时同步 `g` 关系
|
||||
|
||||
**中间件** `internal/auth/middleware.go`:
|
||||
|
||||
- 从 `Authorization: Bearer / Cookie(token|access_token|refresh_token) / ?token=` 提取 token
|
||||
- 注入 `context.Context` 的 `AuthedUser{Username, Role, Token}`
|
||||
- `AuthorizeMiddleware(obj, act, needAuth)`:`guest` 放行 `config:read + srs:streams + room:chat`,其余走 `Enforce`
|
||||
|
||||
### 4.3 房间分发拓扑 (roomHub + stream_targets)
|
||||
|
||||
**内存 Hub** `internal/server/rooms.go`:
|
||||
|
||||
- `roomHub{rooms map[name]*roomEntry, db *sql.DB}`,`roomEntry{targets map[backend]*StreamTarget, subs map[chan]struct{}}`
|
||||
- `setTarget(room, backend, target)`:写入内存 → `broadcast(room)` → 异步 `db.SaveTarget`
|
||||
- `removeTarget`:删 backend,空房间且无订阅者则删 entry + `db.DeleteTarget`
|
||||
- `subscribe(room)`:每房间 8 缓冲 channel,`broadcast` 非阻塞 fan-out
|
||||
- `broadcastEvent(room, type, payload)`:复用同一订阅通道,以 `{"type": "...", "payload": ...}` 信封推送 `playlist_update / playback` 等通用事件
|
||||
- `newRoomHubWithDB` 启动时 `db.LoadAll` 恢复全量房间
|
||||
|
||||
**持久化** `internal/db/db.go`:
|
||||
|
||||
- `Open(dsn)` 仅允许 `file:` / `:memory:`,自动 `MkdirAll` + `Migrate`
|
||||
- 表:
|
||||
- `stream_targets(room, backend PK, session_id, stream, publish_token, url, published_at)` + `idx_room`
|
||||
- `rooms(name PK, created_at, updated_at)` (轻量索引)
|
||||
- `rooms_v2 / room_members / room_bans / room_join_requests / room_categories / room_labels / user_bans / user_registration_requests`
|
||||
- `SaveTarget = INSERT OR REPLACE` + 维护 `rooms` 占位;`DeleteTarget` 空房间时删 `rooms`
|
||||
|
||||
**控制面 API** `internal/server/service.go + gateway.go`:
|
||||
|
||||
- `GetConfig`:按 `SFU_PROVIDER` 顺序组 `BackendInfo[] + candidate + token_required`
|
||||
- `Publish(room, backend, identity)`:
|
||||
- `cloudflare`:`CreateSession(room) → GetSession → target{session_id}`;未配置时直接报错
|
||||
- `srs`:`stream=live-<room>` + `signStreamToken(secret, stream, identity, "publish", 2h)` → `target{stream, publish_token, url=/rtc/v1/whep/?app=live&stream=...}`
|
||||
- 统一 `hub.setTarget` 并广播
|
||||
- `Subscribe(room, backend)`:
|
||||
- `cloudflare`:新建 viewer session,返回 `session_id + publisher_session_id + ice_servers`
|
||||
- `srs`:返回 `stream + ice_servers(stun)`,观看侧 WHEP 与 HLS 共用同一 `stream`
|
||||
- `StopStream(room, backend)`:`cloudflare` 删发布 session + `hub.removeTarget`
|
||||
- `WatchRoom (gRPC stream) / GET /api/room/{room}/events (SSE)`:首包即时推送 `targets`,后续 fan-out + 20s 心跳全量重推
|
||||
- `handleRoomsQuery GET /api/rooms/query`:支持 `page/page_size/search/status/creator/is_public/sort_by/sort_direction`,优先走 `room.Store.QueryRooms` (DB),否则回退内存 hub
|
||||
|
||||
### 4.4 媒体面反代与 SFU
|
||||
|
||||
**SRS** `internal/server/proxy.go + deploy/srs.conf`:
|
||||
|
||||
- `srsProxyHandler`:`SRS_API_BASE (默认 :1985)` 的 WHIP/WHEP 信令反代;`SFU_TOKEN_REQUIRED=1` 时校验 `?token=` 为 `role=publish` 且 `room` 一致的 HS256 JWT
|
||||
- `srsHlsProxyHandler`:`SRS_HTTP_BASE (默认 :8080)` 的 HLS/FLV 静态切片反代,挂 `GET/HEAD /live/`,自动补 CORS `Allow-Origin: *`,`OPTIONS` 204
|
||||
- SRS 配置:`http_server :8080` 切片落盘 `objs/nginx/html`,`http_api :1985` 信令,`rtc_server :8000` 媒体,`vhost __defaultVhost__ { rtc + http_remux(flv) + hls(fragment 10s, window 60s) }`
|
||||
|
||||
**Cloudflare** `internal/sfu/cloudflare/client.go + provider.go`:
|
||||
|
||||
- `Provider{client, appID, stun, configured}`,`BackendInfo{kind=CLOUDFLARE, name="Cloudflare Realtime", configured, primary=true}`
|
||||
- `cfProxyHandler`:将 `/api/cf/*` 转 `CF_BASE_URL/apps/<CF_APP_ID>/*`,注入 `Authorization: Bearer CF_APP_SECRET`,浏览器只做 `tracks/new` SDP 交换
|
||||
|
||||
**前端媒体链路**:
|
||||
|
||||
- 推流:`RTCPeerConnection → createOffer → POST /rtc/v1/whip/?app=live&stream=live-<room>&token=... (信令经网关)` → `setRemoteDescription(answer)`,成功后展示 `hlsUrl=/live/<stream>.m3u8 + flvUrl`
|
||||
- 观看:
|
||||
- `WHEP`:`POST /rtc/v1/whep/?app=live&stream=...` 同流程,`ontrack → <video>`
|
||||
- `HLS`:`video.canPlayType('application/vnd.apple.mpegurl')` 走原生,否则 `hls.js loadSource/attachMedia`;观看页三档切 `cloudflare / srs(whep) / srs-hls`
|
||||
|
||||
### 4.5 聊天弹幕
|
||||
|
||||
**hub** `internal/server/chat.go`:
|
||||
|
||||
- `chatMessage{id, room, user, role, message, color, host, ts}`;`chatHub{rooms map[room]map[chan]struct{}, recentMsgs map[room][]*chatMessage}`,`recentLimit=50`
|
||||
- `subscribe(room)` 16 缓冲;`broadcast` 进 `recentMsgs` 并裁剪;`recent(room)` 快照回放
|
||||
- 空房间无订阅者时清 `recentMsgs`,防内存无界
|
||||
- `GET /api/room/{room}/chat (SSE)`:进场回放 recent + 25s ping;`POST /api/room/{room}/chat` 与 `POST /api/room/{room}/broadcast` (主播专用,`host=true`) 均经 `chatHub.broadcast`
|
||||
- 限流 `chatRateLimiter(limit=5, window=3s)`,key=`room/user/ip` (broadcast 另缀 `/broadcast/`),滑动窗口 + 分钟级 `cleanup`
|
||||
- `clientIP` 取 `X-Forwarded-For[0] / X-Real-IP / RemoteAddr`
|
||||
|
||||
**权限**:`POST /chat` 需 `room:chat` (guest 亦可);`POST /broadcast` 需 `room:publish` (publisher/admin),走 `authWrap(..., needAuth=true)`
|
||||
|
||||
### 4.6 播放列表与同步播放
|
||||
|
||||
`internal/server/playlist.go + internal/room/playlist.go`:
|
||||
|
||||
- `PlaylistItem{room, url, title, added_by}`,`PlaybackState{room, position, playing, rate, current_item, updated_by}`
|
||||
- `GET /api/room/{room}/playlist` 列件;`POST` (需 `PermControlStream`) 增件并 `broadcastEvent("playlist_update", items)`;`DELETE ?id=` 删件同广播
|
||||
- `GET /api/room/{room}/playback` 读状态(无则 `rate=1`);`POST/PUT` 需 `PermControlStream`,`rate<=0` 归一,落库后 `broadcastEvent("playback", ps)`,观看端经同一 SSE 连接按 `event: playlist_update / playback` 分发
|
||||
|
||||
### 4.7 房间内权限位
|
||||
|
||||
`internal/room/permission.go`:
|
||||
|
||||
- 16 项:`send_chat / view_chat_history / view_members / subscribe / watch_room / publish / stop_stream / broadcast / manage_members / add_members / remove_members / manage_permissions / manage_room / delete_room / delete_chat_messages / control_stream`
|
||||
- `PermissionSet uint64` 位集,`Grant/Revoke/Has/Names`
|
||||
- 默认集:
|
||||
- `PermDefaultCreator = All`
|
||||
- `PermDefaultAdmin` = 除 `delete_room` 外大部分
|
||||
- `PermDefaultMember` = `chat/view/subscribe/watch/publish/stop`
|
||||
- `PermDefaultGuest` = `view_chat/view_members/subscribe/watch`;`PermGuestAssignable = Guest + send_chat`
|
||||
- `RoomSettings.Admin/Member/GuestPermissions()` 在 `Default*` 上叠加 `added/removed` 增量
|
||||
|
||||
### 4.8 管理后台 (Manage)
|
||||
|
||||
`internal/server/manage.go` 挂于 `/api/manage/*`,均 `authWrap(needAuth=true)`:
|
||||
|
||||
- `GET /api/manage/overview`:汇总计数
|
||||
- `GET/POST /api/manage/rooms`:列/建房间(含 `is_public/require_approval/max_members/password/display_name/description`)
|
||||
- `GET/PUT/PATCH/DELETE /api/manage/rooms/{room}`:详情/改设置/删房
|
||||
- `GET/POST /api/manage/rooms/{room}/members`、`PUT/DELETE .../members/{user}`、`POST .../perms`、`POST .../approve|reject`、`POST .../transfer`:成员增删改、权限覆盖、审批、转让房主 (`RoleCreator` 互转)
|
||||
- `GET /api/manage/permissions[/{room}]`:权限矩阵
|
||||
- `GET/PUT /api/manage/rooms/{room}/settings`、`POST .../join|leave`:房间设置与自助进出
|
||||
|
||||
---
|
||||
|
||||
## 5. 关键流程时序
|
||||
|
||||
### 5.1 登录 / 注册 / 鉴权
|
||||
|
||||
1. `POST /api/auth/login|register {username,password,role?}` → `UserStore.Verify/Create` → `JWT.Sign` → `Set-Cookie: token=... HttpOnly + {token,user,expires_at}`
|
||||
2. 后续请求带 `Cookie` 或 `Authorization: Bearer`,`AuthorizeMiddleware` → `JWT.Verify` → 注入 `AuthedUser` → `Enforcer.Enforce(sub,obj,act)`
|
||||
3. `POST /api/auth/refresh` 用 `refresh_token` 换新 `access`;`GET /api/auth/me` 返回当前用户与 `expires_at/issued_at`
|
||||
4. 未登录 `guest` 默认仅 `config:read` 与 `srs:streams` 与 `room:chat` 放行,其余 401/403
|
||||
|
||||
### 5.2 房间生命周期
|
||||
|
||||
1. `POST /api/rooms {name}` 或 `POST /api/manage/rooms` → `hub.createRoom` (内存占位);若走 `room.Service.CreateRoom` 则同时落 `rooms_v2 + room_members(creator)`
|
||||
2. `GET /api/rooms | GET /api/rooms/query?...` 分页查;`GET /api/manage/rooms/{room}` 看成员与 `targets(live)`;`GET /live/*.m3u8` 经网关拉 HLS 无需订阅
|
||||
3. 转让:`POST /api/manage/rooms/{room}/transfer {target}` 需当前 `RoleCreator`,事务内改 `rooms_v2.creator` 并互换 `creator↔admin` 角色
|
||||
4. 关闭/封禁:`RoomStatus` 流转 `active↔closed/banned`,`EnsurePermission` 据此限权
|
||||
|
||||
### 5.3 发布 (Publish)
|
||||
|
||||
```
|
||||
前端 → POST /api/publish {room, backend} (需 room:publish)
|
||||
→ Service.Publish 按 backend 分流
|
||||
├─ cloudflare: CreateSession(room) → target{session_id}
|
||||
└─ srs: stream=live-<room>, token=HS256(room,identity,publish,2h) → target{stream,token,url}
|
||||
→ hub.setTarget(room, backend, target) → SSE 广播 RoomEvent{targets}
|
||||
→ 返回 {session_id/stream/publish_token/ice_servers/target}
|
||||
前端 → RTCPeerConnection offer → POST /rtc/v1/whip/?app=live&stream=...&token=... (SRS 反代校验) 或 POST /api/cf/.../tracks/new (CF 反代注入 Bearer)
|
||||
→ setRemoteDescription(answer) → 推流成功,展示 HLS/FLV 链接
|
||||
```
|
||||
|
||||
### 5.4 订阅 / 观看
|
||||
|
||||
```
|
||||
前端 → POST /api/subscribe {room, backend} (需 room:subscribe)
|
||||
→ hub.findTarget(room, backend) 取发布侧 target
|
||||
├─ cloudflare: CreateSession(room) viewer → {session_id, publisher_session_id, ice_servers}
|
||||
└─ srs: {stream, ice_servers}
|
||||
前端 → 按 backends[].configured 展示线路,未配置标「未配置」
|
||||
├─ cloudflare: tracks/new location=remote
|
||||
├─ srs-whep: POST /rtc/v1/whep/?app=live&stream=live-<room> → ontrack
|
||||
└─ srs-hls: GET /live/live-<room>.m3u8 → hls.js / 原生
|
||||
→ 同时订阅 GET /api/room/{room}/events (SSE) + GET /api/room/{room}/chat (SSE),拓扑/弹幕/播放状态同链路按 event type 分发
|
||||
```
|
||||
|
||||
### 5.5 弹幕与主播广播
|
||||
|
||||
```
|
||||
观众 POST /api/room/{room}/chat {message,color} (room:chat)
|
||||
→ 取身份 user/role/id (guest 可发)
|
||||
→ chatRL.allow(room/user/ip, 5/3s) 限流
|
||||
→ chatHub.broadcast → SSE fan-out (event: chat) → 前端 danmaku 叠加
|
||||
|
||||
主播 POST /api/room/{room}/broadcast {message,color} (room:publish, 需 Bearer)
|
||||
→ 同限流 key=room/broadcast/user/ip
|
||||
→ chatMessage{host:true} → 同一 hub 广播,观众/主播端同可见
|
||||
```
|
||||
|
||||
### 5.6 停止
|
||||
|
||||
`POST /api/stop {room, backend} (room:stop)` → `cloudflare` 删 `session_id` + `hub.removeTarget(room, backend)` → 广播空 targets → 前端自动切离或提示下播
|
||||
|
||||
---
|
||||
|
||||
## 6. 数据模型与契约
|
||||
|
||||
### 6.1 Protobuf 契约 `api/sync_live.proto`
|
||||
|
||||
- 服务 `SyncLive`:`GetConfig / ListRooms / Publish / Subscribe / StopStream / WatchRoom(stream) / Login / Register / GetMe / ListUsers / UpdateUserRole`
|
||||
- 枚举 `BackendKind`: `UNSPECIFIED(0) / CLOUDFLARE(1) / SRS(2)`
|
||||
- 消息:`BackendInfo{kind,name,configured,primary}`, `IceServer{urls,username,credential}`, `StreamTarget{backend,session_id,stream,publish_token,url,published_at}`, `Room{name,targets[]}`, `RoomEvent{room,targets[]}`, `UserInfo{username,role,created_at}`, `Login/Register/GetMe/ListUsers/UpdateUserRole` 族
|
||||
|
||||
### 6.2 libSQL 表
|
||||
|
||||
- `stream_targets`:房间分发的事实表,`PRIMARY KEY(room,backend)`
|
||||
- `rooms / rooms_v2`:`rooms` 为轻量拓扑索引,`rooms_v2` 为完整房间档案(含 `settings_json`)
|
||||
- `room_members(room,user_id PK, username, role, status, added/removed/admin_added/admin_removed, joined_at, version)` + 两索引
|
||||
- `room_bans / room_join_requests / room_categories / room_labels / user_bans / user_registration_requests`
|
||||
- 增量兼容:启动时 `ALTER TABLE rooms_v2 ADD COLUMN category_id/cover_file_id/label_ids` 忽略已存在错误
|
||||
|
||||
### 6.3 用户库 `data/users.json`
|
||||
|
||||
- `map[username]*User{username, hash(bcrypt), role, status, created_at}`,文件锁读写,`BOOTSTRAP_ADMIN_*` 冷启动种子
|
||||
|
||||
---
|
||||
|
||||
## 7. API 清单
|
||||
|
||||
### 7.1 JSON 网关 (同 HTTP 端口,protojson)
|
||||
|
||||
| 方法 | 路径 | 鉴权 | 说明 |
|
||||
|------|------|------|------|
|
||||
| GET | `/api/config` | `config:read` guest可 | 后端能力与 `candidate/token_required` |
|
||||
| GET | `/api/rooms` | `room:list` | 全量房间 (hub) |
|
||||
| GET | `/api/rooms/query?...` | `room:list` | 分页查询 (DB 优先) |
|
||||
| POST | `/api/rooms` | `room:list` | 创建占位 |
|
||||
| POST | `/api/publish` | `room:publish` | 发布 (多后端) |
|
||||
| POST | `/api/subscribe` | `room:subscribe` | 订阅 |
|
||||
| POST | `/api/stop` | `room:stop` | 下播 |
|
||||
| GET | `/api/srs/streams` | `srs:streams` guest可 | 代理 `SRS :1985 /api/v1/streams/` |
|
||||
| GET | `/api/room/{room}/events` | `room:watch` SSE | 房间拓扑 |
|
||||
| GET | `/api/room/{room}/playlist` | `room:watch` | 列播放列表 |
|
||||
| POST | `/api/room/{room}/playlist` | `room:watch` + `control_stream` | 增件 |
|
||||
| DELETE | `/api/room/{room}/playlist?id=` | 同上 | 删件 |
|
||||
| GET | `/api/room/{room}/playback` | `room:watch` | 读同步状态 |
|
||||
| POST/PUT | `/api/room/{room}/playback` | 同上 | 写同步状态 |
|
||||
| GET | `/api/room/{room}/chat` | `room:chat` SSE guest可 | 弹幕订阅 + recent 回放 |
|
||||
| POST | `/api/room/{room}/chat` | `room:chat` guest可 | 发弹幕 |
|
||||
| POST | `/api/room/{room}/broadcast` | `room:publish` | 主播广播 (`host:true`) |
|
||||
|
||||
### 7.2 认证
|
||||
|
||||
| 方法 | 路径 | 说明 |
|
||||
|------|------|------|
|
||||
| POST | `/api/auth/login` | 登录,签 access+refresh,种 HttpOnly Cookie |
|
||||
| POST | `/api/auth/register` | 注册,`ALLOW_REGISTER=1` 时放行,`admin/publisher` 需 admin 授权否则降为 viewer |
|
||||
| POST | `/api/auth/logout` | 清 Cookie |
|
||||
| POST | `/api/auth/refresh` | 用 refresh 换 access |
|
||||
| GET | `/api/auth/me` | 当前用户 |
|
||||
| GET | `/api/auth/users` | 列用户 (admin) |
|
||||
| POST | `/api/auth/users/role` | 改角色 (`only root can grant root`, admin 不能提权他人为 admin/root) |
|
||||
| POST | `/api/auth/users/ban|unban` | 封/解封 (admin/root) |
|
||||
| GET | `/api/auth/check` | 前端权限探针 |
|
||||
|
||||
### 7.3 管理面 `/api/manage/*`
|
||||
|
||||
见 4.8 节,覆盖 `overview/rooms/rooms/{room}/members/permissions/settings/join/leave/transfer` 全链路。
|
||||
|
||||
### 7.4 媒体反代
|
||||
|
||||
| 方法 | 路径 | 目标 | 说明 |
|
||||
|------|------|------|------|
|
||||
| ANY | `/rtc/v1/*` | `SRS_API_BASE (:1985)` | WHIP/WHEP 信令,可选 `?token` 校验 |
|
||||
| ANY | `/api/cf/*` | `CF_BASE_URL/apps/<CF_APP_ID>/*` | Cloudflare Realtime,注 Bearer |
|
||||
| GET/HEAD | `/live/*` | `SRS_HTTP_BASE (:8080)` | HLS/FLV 切片,补 CORS |
|
||||
|
||||
### 7.5 gRPC `:9090`
|
||||
|
||||
同 `SyncLive` 服务,`UnaryAuthInterceptor + StreamAuthInterceptor` 校验 `room/user/config/srs/system` 等 `obj:act`。
|
||||
|
||||
---
|
||||
|
||||
## 8. 权限矩阵
|
||||
|
||||
### 8.1 全局 Casbin (policy.csv)
|
||||
|
||||
| 角色 | 能力 |
|
||||
|------|------|
|
||||
| `root` | 全量:`config:read, room:list/publish/subscribe/stop/watch/chat/manage/delete, user:list/manage, srs:streams, system:manage` |
|
||||
| `admin` | 除 `system:manage/room:delete` 外全量 |
|
||||
| `publisher` | `config:read + room:list/publish/subscribe/stop/watch/chat + srs:streams` |
|
||||
| `viewer` | `config:read + room:list/subscribe/watch/chat + srs:streams` |
|
||||
| `guest` | `config:read + srs:streams + room:chat` (弹幕可发,拉流只看 HLS) |
|
||||
|
||||
### 8.2 房间内位权限 (PermissionSet)
|
||||
|
||||
见 4.7 表,房间设置可在默认集上按 `admin/member/guest` 三档增删位,且受 `Role.CanManage` 等级约束。关键校验入口 `room.Service.EnsurePermission`。
|
||||
|
||||
---
|
||||
|
||||
## 9. 配置与部署
|
||||
|
||||
### 9.1 环境变量 (env > 默认)
|
||||
|
||||
| 变量 | 默认 | 说明 |
|
||||
|------|------|------|
|
||||
| `HTTP_PORT/GRPC_PORT` | `8088/9090` | 控制面端口 |
|
||||
| `JWT_SECRET` | (必填) | 正式环境必须显式配置,否则 `Validate` 拒绝启动 |
|
||||
| `SFU_TOKEN_SECRET` | 同 JWT | 推流 JWT 密钥,空则复用 JWT |
|
||||
| `SFU_TOKEN_REQUIRED` | `0` | `1` 时 WHIP 强制 `?token` 校验 |
|
||||
| `SFU_PROVIDER` | `cloudflare,srs` | 后端顺序 |
|
||||
| `SRS_API_BASE/SRS_HTTP_BASE/SRS_APP/SRS_CANDIDATE` | `http://localhost:1985 / :8080 / live / 127.0.0.1` | SRS 链路 |
|
||||
| `CF_APP_ID/CF_APP_SECRET/CF_BASE_URL/CF_STUN_URL` | (空) / `https://rtc.live.cloudflare.com/v1` / `stun:stun.cloudflare.com:3478` | Cloudflare 链路 |
|
||||
| `TURSO_DATABASE_URL/DATABASE_URL` | `file:./data/sync-live.db?cache=shared&_journal_mode=WAL` | 仅 file/`:memory:`,拒 `libsql://` |
|
||||
| `AUTH_USER_FILE/CASBIN_MODEL/CASBIN_POLICY` | `data/users.json / internal/auth/model.conf / policy.csv` | 认证文件 |
|
||||
| `ALLOW_REGISTER` | `1` | 是否开放注册 |
|
||||
| `BOOTSTRAP_ADMIN_USER/PASS/ROLE` | (空) | 冷启动种子,仅空库生效 |
|
||||
| `SMTP_*/EMAIL_VERIFY_*/OAUTH_*/AUTHBOSS_ENABLED` | (多为 `0`/空) | 邮件/验证/OAuth 预留开关 |
|
||||
|
||||
### 9.2 本地启动
|
||||
|
||||
```bash
|
||||
cp .env.example .env # 填 JWT_SECRET 等
|
||||
cd deploy && SRS_CANDIDATE=127.0.0.1 docker compose up -d srs
|
||||
go run ./cmd/server # http://localhost:8088 /publish?room=xxx /watch?room=xxx /login /manage
|
||||
# 前端热更新
|
||||
cd web && pnpm i && pnpm dev # :5173
|
||||
pnpm build # 产到 internal/server/static (Go embed)
|
||||
```
|
||||
|
||||
### 9.3 SRS 部署要点
|
||||
|
||||
- `http_server :8080` 切片目录与 `hls_path` 一致;`hls_fragment 10 / window 60`;公网部署 `SRS_CANDIDATE=公网IP` 并重建容器
|
||||
- 生产可不暴露 `8080`,仅保留 `1985/8000`,前端统一经 `8088 /live/*` 拉流
|
||||
- HLS 首屏 10-20s,重试 `curl /live/<stream>.m3u8` 与 `docker exec ls objs/nginx/html/live` 排障
|
||||
|
||||
---
|
||||
|
||||
## 10. 前端
|
||||
|
||||
- 路由 `web/src/app/router.tsx`:`TanStack Router + RootRoute(Layout) + lazy(code-split)`,路由 `/ /publish?room= /watch?room= /room?room= /rooms /users /manage /login`,均为 SPA 回退 `index.html`
|
||||
- 状态:`TanStack Query` 封装 `api.ts` 的 `fetch /api/* (protojson)`;SSE 封装 `sse.ts/chat.ts`,`GET /api/room/{room}/events` 与 `/chat` 并行订阅,按 `event: room/chat/playlist_update/playback` 分发
|
||||
- 媒体:`webrtc.ts` 封装 WHIP/WHEP 的 `RTCPeerConnection` 流程,`hls.ts` 封装 `hls.js` 回退
|
||||
- 组件:`components/danmaku.tsx` 叠加层,`components/manage/*` 管理后台分 `Overview/Rooms/Permissions/System/UsersInline`
|
||||
- 构建:`vite.config.ts` 输出 `../internal/server/static`,`internal/server/server.go:embed static` 内联,`GET / /publish /watch /login /manage` 均回 `static/index.html`
|
||||
|
||||
---
|
||||
|
||||
## 11. 安全与限流
|
||||
|
||||
- 正式环境密钥强校验:缺 `JWT_SECRET` 或 `SFU_TOKEN_REQUIRED` 却无 `SFU_TOKEN_SECRET` 直接退错
|
||||
- 密码:用户 `bcrypt`,房间密码 `bcrypt` (`HashRoomPassword`)
|
||||
- 推流 token:`token.go` 自实现 `HS256(b64url header + payload).hmacSHA256(secret)`,`{room,identity,role,iat,exp}`,TTL 2h
|
||||
- 弹幕限流:每 `room/user/ip` 5 条 / 3s 滑动窗口,超限 `429`,后台分钟级清过期 key;单条 500 字符截断
|
||||
- 鉴权覆盖:HTTP `authWrap` 与 gRPC `Unary/StreamInterceptor` 双轨,`guest` 白名单最小化
|
||||
|
||||
---
|
||||
|
||||
## 12. 关联文档与验证
|
||||
|
||||
- 推流链路详述:`docs/streaming-pipeline.md`(含 SRS `srs.conf` 全量、反代代码、前端 `publishSRS/watchSRS/watchHLS` 选型表与 `curl/ffplay` 验证步骤)
|
||||
- 架构分层:`design.md`
|
||||
- 协议契约:`api/sync_live.proto` + `buf.gen.yaml`
|
||||
- 快速验证:发布页推流后 `curl -i http://localhost:8088/live/live-demo.m3u8` / `curl -I .../live-demo.flv` / `ffplay ...`,SRS 日志与切片目录联合排障
|
||||
|
||||
---
|
||||
|
||||
*维护提示:新增后端或房间字段时,同步更新 `api/*.proto → gen/`、`internal/db.Migrate / room.Store.InitSchema`、`policy.csv` 与本文件对应小节;房间拓扑的持久化与 SSE 广播是分流可见性的核心链路,改动前优先补 `grpc_test.go / auth_test.go` 用例。*
|
||||
|
||||
|
|
@ -0,0 +1,558 @@
|
|||
# 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 执行任务,批量执行并设有检查点
|
||||
|
||||
选哪种方式?
|
||||
25
go.mod
25
go.mod
|
|
@ -12,39 +12,14 @@ require (
|
|||
)
|
||||
|
||||
require (
|
||||
cloud.google.com/go/compute/metadata v0.9.0 // indirect
|
||||
github.com/aarondl/authboss/v3 v3.5.3 // indirect
|
||||
github.com/antlr4-go/antlr/v4 v4.13.0 // indirect
|
||||
github.com/bmatcuk/doublestar/v4 v4.6.1 // indirect
|
||||
github.com/casbin/govaluate v1.3.0 // indirect
|
||||
github.com/dghubble/oauth1 v0.7.3 // indirect
|
||||
github.com/friendsofgo/errors v0.9.2 // indirect
|
||||
github.com/go-oauth2/oauth2/v4 v4.5.4 // indirect
|
||||
github.com/go-pkgz/auth v1.26.0 // indirect
|
||||
github.com/go-pkgz/repeater/v2 v2.2.0 // indirect
|
||||
github.com/go-pkgz/rest v1.24.0 // indirect
|
||||
github.com/golang-jwt/jwt v3.2.2+incompatible // indirect
|
||||
github.com/golang/snappy v1.0.0 // indirect
|
||||
github.com/google/uuid v1.6.0 // indirect
|
||||
github.com/klauspost/compress v1.19.2 // indirect
|
||||
github.com/libsql/sqlite-antlr4-parser v0.0.0-20240327125255-dbf53b6cbf06 // indirect
|
||||
github.com/markbates/goth v1.82.0 // indirect
|
||||
github.com/montanaflynn/stats v0.12.4 // indirect
|
||||
github.com/rrivera/identicon v0.0.0-20240116195454-d5ba35832c0d // indirect
|
||||
github.com/wneessen/go-mail v0.8.1 // indirect
|
||||
github.com/xdg-go/pbkdf2 v1.0.0 // indirect
|
||||
github.com/xdg-go/scram v1.2.0 // indirect
|
||||
github.com/xdg-go/stringprep v1.0.4 // indirect
|
||||
github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78 // indirect
|
||||
go.etcd.io/bbolt v1.5.0 // indirect
|
||||
go.mongodb.org/mongo-driver v1.17.9 // indirect
|
||||
golang.org/x/exp v0.0.0-20230515195305-f3d0a9c9a5cc // indirect
|
||||
golang.org/x/image v0.45.0 // indirect
|
||||
golang.org/x/net v0.57.0 // indirect
|
||||
golang.org/x/oauth2 v0.36.0 // indirect
|
||||
golang.org/x/sync v0.22.0 // indirect
|
||||
golang.org/x/sys v0.47.0 // indirect
|
||||
golang.org/x/text v0.41.0 // indirect
|
||||
golang.org/x/xerrors v0.0.0-20220907171357-04be3eba64a2 // indirect
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20260120221211-b8f7ae30c516 // indirect
|
||||
)
|
||||
|
|
|
|||
70
go.sum
70
go.sum
|
|
@ -1,7 +1,3 @@
|
|||
cloud.google.com/go/compute/metadata v0.9.0 h1:pDUj4QMoPejqq20dK0Pg2N4yG9zIkYGdBtwLoEkH9Zs=
|
||||
cloud.google.com/go/compute/metadata v0.9.0/go.mod h1:E0bWwX5wTnLPedCKqk3pJmVgCBSM6qQI1yTBdEb3C10=
|
||||
github.com/aarondl/authboss/v3 v3.5.3 h1:eF2G85mGjADsRQnBekMSsgqvW5Naec6kew5B0PRT1nE=
|
||||
github.com/aarondl/authboss/v3 v3.5.3/go.mod h1:Q/t4SEfVVI+i/JPVuB6gB5LBSkqBnsf4JezLekg5X/g=
|
||||
github.com/antlr4-go/antlr/v4 v4.13.0 h1:lxCg3LAv+EUK6t1i0y1V6/SLeUi0eKEKdhQAlS8TVTI=
|
||||
github.com/antlr4-go/antlr/v4 v4.13.0/go.mod h1:pfChB/xh/Unjila75QW7+VU4TSnWnnk9UTnmpPaOR2g=
|
||||
github.com/bmatcuk/doublestar/v4 v4.6.1 h1:FH9SifrbvJhnlQpztAx++wlkk70QBf0iBWDwNy7PA4I=
|
||||
|
|
@ -12,65 +8,26 @@ github.com/casbin/govaluate v1.3.0 h1:VA0eSY0M2lA86dYd5kPPuNZMUD9QkWnOCnavGrw9my
|
|||
github.com/casbin/govaluate v1.3.0/go.mod h1:G/UnbIjZk/0uMNaLwZZmFQrR72tYRZWQkO70si/iR7A=
|
||||
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/dghubble/oauth1 v0.7.3 h1:EkEM/zMDMp3zOsX2DC/ZQ2vnEX3ELK0/l9kb+vs4ptE=
|
||||
github.com/dghubble/oauth1 v0.7.3/go.mod h1:oxTe+az9NSMIucDPDCCtzJGsPhciJV33xocHfcR2sVY=
|
||||
github.com/friendsofgo/errors v0.9.2 h1:X6NYxef4efCBdwI7BgS820zFaN7Cphrmb+Pljdzjtgk=
|
||||
github.com/friendsofgo/errors v0.9.2/go.mod h1:yCvFW5AkDIL9qn7suHVLiI/gH228n7PC4Pn44IGoTOI=
|
||||
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/go-oauth2/oauth2/v4 v4.5.4 h1:YjI0tmGW8oxVhn9QSBIxlr641QugWrJY5UWa6XmLcW0=
|
||||
github.com/go-oauth2/oauth2/v4 v4.5.4/go.mod h1:BXiOY+QZtZy2ewbsGk2B5P8TWmtz/Rf7ES5ZttQFxfQ=
|
||||
github.com/go-pkgz/auth v1.26.0 h1:2CL91EO8ntGWcHkXj2AXDx2NUfZ//qgo+DTX6ZLyffM=
|
||||
github.com/go-pkgz/auth v1.26.0/go.mod h1:z4stIOtrnxMumDchQ3BP7sDOP36Qyq7+xgT9+E/E+/c=
|
||||
github.com/go-pkgz/repeater/v2 v2.2.0 h1:8nZR/NaknmLfx2YMHbr78u9OL4Xj+8+romm9dz4FpMg=
|
||||
github.com/go-pkgz/repeater/v2 v2.2.0/go.mod h1:RgX5vUbLKq7PV82QUDP5pFbQS1os4Z+U9XzKymK23A8=
|
||||
github.com/go-pkgz/rest v1.24.0 h1:GAUCgx7U8xCOC2OynLjhCRMhtnMQH4d1mTdKpQyX2yI=
|
||||
github.com/go-pkgz/rest v1.24.0/go.mod h1:dl3EWiuFB4hRTo2Sknj6UrQGFRAYvANK6/NyW8qQPxc=
|
||||
github.com/golang-jwt/jwt v3.2.2+incompatible h1:IfV12K8xAKAnZqdXVzCZ+TOjboZ2keLg81eXfW3O+oY=
|
||||
github.com/golang-jwt/jwt v3.2.2+incompatible/go.mod h1:8pz2t5EyA70fFQQSrl6XZXzqecmYZeUEB8OUGHkxJ+I=
|
||||
github.com/golang-jwt/jwt/v5 v5.3.1 h1:kYf81DTWFe7t+1VvL7eS+jKFVWaUnK9cB1qbwn63YCY=
|
||||
github.com/golang-jwt/jwt/v5 v5.3.1/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE=
|
||||
github.com/golang/mock v1.4.4 h1:l75CXGRSwbaYNpl/Z2X1XIIAMSCquvXgpVZDhwEIJsc=
|
||||
github.com/golang/mock v1.4.4/go.mod h1:l3mdAwkq5BuhzHwde/uurv3sEJeZMXNpwsxVWU71h+4=
|
||||
github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek=
|
||||
github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps=
|
||||
github.com/golang/snappy v1.0.0 h1:Oy607GVXHs7RtbggtPBnr2RmDArIsAefDwvrdWvRhGs=
|
||||
github.com/golang/snappy v1.0.0/go.mod h1:/XxbfmMg8lxefKM7IXC3fBNl/7bRcc72aCRzEWrmP2Q=
|
||||
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=
|
||||
github.com/klauspost/compress v1.19.2 h1:hMRETovs/pu/dVWN7zIT1PGG8t509MwT6bO7XSi26R8=
|
||||
github.com/klauspost/compress v1.19.2/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
|
||||
github.com/libsql/sqlite-antlr4-parser v0.0.0-20240327125255-dbf53b6cbf06 h1:JLvn7D+wXjH9g4Jsjo+VqmzTUpl/LX7vfr6VOfSWTdM=
|
||||
github.com/libsql/sqlite-antlr4-parser v0.0.0-20240327125255-dbf53b6cbf06/go.mod h1:FUkZ5OHjlGPjnM2UyGJz9TypXQFgYqw6AFNO1UiROTM=
|
||||
github.com/markbates/goth v1.82.0 h1:8j/c34AjBSTNzO7zTsOyP5IYCQCMBTRBHAbBt/PI0bQ=
|
||||
github.com/markbates/goth v1.82.0/go.mod h1:/DRlcq0pyqkKToyZjsL2KgiA1zbF1HIjE7u2uC79rUk=
|
||||
github.com/montanaflynn/stats v0.12.4 h1:amtNRsti20yIhcrkfUJGwoYqBR82jKQFE8SNNYVgGn0=
|
||||
github.com/montanaflynn/stats v0.12.4/go.mod h1:etXPPgVO6n31NxCd9KQUMvCM+ve0ruNzt6R8Bnaayow=
|
||||
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
|
||||
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
|
||||
github.com/rrivera/identicon v0.0.0-20240116195454-d5ba35832c0d h1:l3+2LWCbVxn5itfvXAfH9n4YL9jh8l1g5zcncbIc1cs=
|
||||
github.com/rrivera/identicon v0.0.0-20240116195454-d5ba35832c0d/go.mod h1:TbpErkob6SY7cyozRVSGoB3OlO2qOAgVN8O3KAJ4fMI=
|
||||
github.com/tursodatabase/go-libsql v0.0.0-20260424063416-3051e37e6e04 h1:9nlqEMruvXDPynGbZ0RE67kKnkkg3NdnjGccvRABefc=
|
||||
github.com/tursodatabase/go-libsql v0.0.0-20260424063416-3051e37e6e04/go.mod h1:TjsB2miB8RW2Sse8sdxzVTdeGlx74GloD5zJYUC38d8=
|
||||
github.com/wneessen/go-mail v0.8.1 h1:tVcncj02/QySVFw3zr/kXOzZcuFQqBNT6K+Rbgm/pcM=
|
||||
github.com/wneessen/go-mail v0.8.1/go.mod h1:dWZ61zadzCIyvB4y1/YzC5O7MrbbzBfPkARmbosdf8w=
|
||||
github.com/xdg-go/pbkdf2 v1.0.0 h1:Su7DPu48wXMwC3bs7MCNG+z4FhcyEuz5dlvchbq0B0c=
|
||||
github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI=
|
||||
github.com/xdg-go/scram v1.2.0 h1:bYKF2AEwG5rqd1BumT4gAnvwU/M9nBp2pTSxeZw7Wvs=
|
||||
github.com/xdg-go/scram v1.2.0/go.mod h1:3dlrS0iBaWKYVt2ZfA4cj48umJZ+cAEbR6/SjLA88I8=
|
||||
github.com/xdg-go/stringprep v1.0.4 h1:XLI/Ng3O1Atzq0oBs3TWm+5ZVgkq2aqdlvP9JtoZ6c8=
|
||||
github.com/xdg-go/stringprep v1.0.4/go.mod h1:mPGuuIYwz7CmR2bT9j4GbQqutWS1zV24gijq1dTyGkM=
|
||||
github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78 h1:ilQV1hzziu+LLM3zUTJ0trRztfwgjqKnBWNtSRkbmwM=
|
||||
github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78/go.mod h1:aL8wCCfTfSfmXjznFBSZNN13rSJjlIOI1fUNAtF7rmI=
|
||||
github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY=
|
||||
go.etcd.io/bbolt v1.5.0 h1:S7GAl7Fxv12yohbwFfIbQCGDWbQbtDGPET4P/bD4lxU=
|
||||
go.etcd.io/bbolt v1.5.0/go.mod h1:mkltfYE5aUHQxUct9N9V+Kp7aSjFqjgrhcXIS70Lrdk=
|
||||
go.mongodb.org/mongo-driver v1.17.9 h1:IexDdCuuNJ3BHrELgBlyaH9p60JXAvdzWR128q+U5tU=
|
||||
go.mongodb.org/mongo-driver v1.17.9/go.mod h1:LlOhpH5NUEfhxcAwG0UEkMqwYcc4JU18gtCdGudk/tQ=
|
||||
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=
|
||||
|
|
@ -84,50 +41,23 @@ go.opentelemetry.io/otel/sdk/metric v1.39.0/go.mod h1:xq9HEVH7qeX69/JnwEfp6fVq5w
|
|||
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/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
|
||||
golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc=
|
||||
golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M=
|
||||
golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis=
|
||||
golang.org/x/exp v0.0.0-20230515195305-f3d0a9c9a5cc h1:mCRnTeVUjcrhlRmO0VK8a6k6Rrf6TF9htwo2pJVSjIU=
|
||||
golang.org/x/exp v0.0.0-20230515195305-f3d0a9c9a5cc/go.mod h1:V1LtkGg67GoY2N1AnLN78QLrzxkLyJw7RJb1gzOOz9w=
|
||||
golang.org/x/image v0.45.0 h1:FMb1nTbH5H9vF55SriQHgFw5GnNL9Jg6L25BwXKzhB0=
|
||||
golang.org/x/image v0.45.0/go.mod h1:n62x/7RqlwXDvGsSU4u6IUTUf6KghUZ9Bt7cG/T9Fx4=
|
||||
golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4=
|
||||
golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
|
||||
golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
|
||||
golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg=
|
||||
golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c=
|
||||
golang.org/x/net v0.57.0 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE=
|
||||
golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU=
|
||||
golang.org/x/oauth2 v0.34.0 h1:hqK/t4AKgbqWkdkcAeI8XLmbK+4m4G5YeQRrmiotGlw=
|
||||
golang.org/x/oauth2 v0.34.0/go.mod h1:lzm5WQJQwKZ3nwavOZ3IS5Aulzxi68dUSgRHujetwEA=
|
||||
golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs=
|
||||
golang.org/x/oauth2 v0.36.0/go.mod h1:YDBUJMTkDnJS+A4BP4eZBjCqtokkg1hODuPjwiGPO7Q=
|
||||
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek=
|
||||
golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
|
||||
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
|
||||
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
|
||||
golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8=
|
||||
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
|
||||
golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
|
||||
golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ=
|
||||
golang.org/x/text v0.3.8/go.mod h1:E6s5w1FMmriuDzIBO73fBruAKo1PCIq6d2Q6DHfQ8WQ=
|
||||
golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8=
|
||||
golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M=
|
||||
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
|
||||
golang.org/x/tools v0.0.0-20190425150028-36563e24a262/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q=
|
||||
golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
|
||||
golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc=
|
||||
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||
golang.org/x/xerrors v0.0.0-20220907171357-04be3eba64a2 h1:H2TDz8ibqkAF6YGhCdN3jS9O0/s90v0rJh3X/OLHEUk=
|
||||
golang.org/x/xerrors v0.0.0-20220907171357-04be3eba64a2/go.mod h1:K8+ghG5WaK9qNqU5K3HdILfMLy1f3aNYFI/wnl100a8=
|
||||
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=
|
||||
|
|
|
|||
|
|
@ -0,0 +1,157 @@
|
|||
package room
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
// PlaylistItem 表示房间共享播放列表中的一个媒体项。
|
||||
type PlaylistItem struct {
|
||||
ID int64 `json:"id"`
|
||||
RoomName string `json:"room_name"`
|
||||
URL string `json:"url"`
|
||||
Title string `json:"title"`
|
||||
AddedBy string `json:"added_by"`
|
||||
SortOrder int32 `json:"sort_order"`
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
}
|
||||
|
||||
// PlaybackState 表示房间当前的同步播放状态。
|
||||
type PlaybackState struct {
|
||||
Room string `json:"room"`
|
||||
Position float64 `json:"position"`
|
||||
Playing bool `json:"playing"`
|
||||
Rate float64 `json:"rate"`
|
||||
CurrentItem int64 `json:"current_item_id"`
|
||||
UpdatedAt int64 `json:"updated_at"`
|
||||
UpdatedBy string `json:"updated_by"`
|
||||
}
|
||||
|
||||
func InitPlaylistSchema(db *sql.DB) error {
|
||||
stmts := []string{
|
||||
`CREATE TABLE IF NOT EXISTS playlist_items (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
room TEXT NOT NULL,
|
||||
url TEXT NOT NULL,
|
||||
title TEXT NOT NULL DEFAULT '',
|
||||
added_by TEXT NOT NULL,
|
||||
sort_order INTEGER NOT NULL DEFAULT 0,
|
||||
created_at INTEGER NOT NULL,
|
||||
FOREIGN KEY (room) REFERENCES rooms_v2(name) ON DELETE CASCADE
|
||||
)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_playlist_room ON playlist_items(room, sort_order)`,
|
||||
`CREATE TABLE IF NOT EXISTS playback_state (
|
||||
room TEXT PRIMARY KEY,
|
||||
position REAL NOT NULL DEFAULT 0,
|
||||
playing INTEGER NOT NULL DEFAULT 0,
|
||||
rate REAL NOT NULL DEFAULT 1,
|
||||
current_item_id INTEGER NOT NULL DEFAULT 0,
|
||||
updated_at INTEGER NOT NULL,
|
||||
updated_by TEXT NOT NULL DEFAULT '',
|
||||
FOREIGN KEY (room) REFERENCES rooms_v2(name) ON DELETE CASCADE
|
||||
)`,
|
||||
}
|
||||
for _, stmt := range stmts {
|
||||
if _, err := db.Exec(stmt); err != nil {
|
||||
return fmt.Errorf("init playlist schema: %w", err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// AddPlaylistItem 向房间播放列表添加一个媒体项。
|
||||
func (s *Store) AddPlaylistItem(ctx context.Context, item *PlaylistItem) (*PlaylistItem, error) {
|
||||
if s.db == nil {
|
||||
return nil, fmt.Errorf("no db")
|
||||
}
|
||||
item.CreatedAt = time.Now()
|
||||
res, err := s.db.ExecContext(ctx,
|
||||
`INSERT INTO playlist_items (room, url, title, added_by, sort_order, created_at) VALUES (?, ?, ?, ?, ?, ?)`,
|
||||
item.RoomName, item.URL, item.Title, item.AddedBy, item.SortOrder, item.CreatedAt.Unix())
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
id, _ := res.LastInsertId()
|
||||
item.ID = id
|
||||
return item, nil
|
||||
}
|
||||
|
||||
// ListPlaylistItems 返回房间的播放列表(按排序字段排列)。
|
||||
func (s *Store) ListPlaylistItems(ctx context.Context, roomName string) ([]*PlaylistItem, error) {
|
||||
if s.db == nil {
|
||||
return nil, nil
|
||||
}
|
||||
rows, err := s.db.QueryContext(ctx,
|
||||
`SELECT id, room, url, title, added_by, sort_order, created_at FROM playlist_items WHERE room = ? ORDER BY sort_order ASC, id ASC`, roomName)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
var out []*PlaylistItem
|
||||
for rows.Next() {
|
||||
var it PlaylistItem
|
||||
var createdAt int64
|
||||
if err := rows.Scan(&it.ID, &it.RoomName, &it.URL, &it.Title, &it.AddedBy, &it.SortOrder, &createdAt); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
it.CreatedAt = time.Unix(createdAt, 0)
|
||||
out = append(out, &it)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// RemovePlaylistItem 删除播放列表中的一个媒体项。
|
||||
func (s *Store) RemovePlaylistItem(ctx context.Context, roomName string, itemID int64) error {
|
||||
if s.db == nil {
|
||||
return fmt.Errorf("no db")
|
||||
}
|
||||
_, err := s.db.ExecContext(ctx, `DELETE FROM playlist_items WHERE room = ? AND id = ?`, roomName, itemID)
|
||||
return err
|
||||
}
|
||||
|
||||
// ReorderPlaylistItems 批量更新播放列表项的排序。
|
||||
func (s *Store) ReorderPlaylistItems(ctx context.Context, roomName string, ids []int64) error {
|
||||
if s.db == nil {
|
||||
return fmt.Errorf("no db")
|
||||
}
|
||||
for i, id := range ids {
|
||||
if _, err := s.db.ExecContext(ctx,
|
||||
`UPDATE playlist_items SET sort_order = ? WHERE room = ? AND id = ?`, i, roomName, id); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// SavePlaybackState 保存或更新房间的同步播放状态。
|
||||
func (s *Store) SavePlaybackState(ctx context.Context, ps *PlaybackState) error {
|
||||
if s.db == nil {
|
||||
return fmt.Errorf("no db")
|
||||
}
|
||||
ps.UpdatedAt = time.Now().Unix()
|
||||
_, err := s.db.ExecContext(ctx,
|
||||
`INSERT INTO playback_state (room, position, playing, rate, current_item_id, updated_at, updated_by)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT(room) DO UPDATE SET position=excluded.position, playing=excluded.playing, rate=excluded.rate, current_item_id=excluded.current_item_id, updated_at=excluded.updated_at, updated_by=excluded.updated_by`,
|
||||
ps.Room, ps.Position, boolToInt(ps.Playing), ps.Rate, ps.CurrentItem, ps.UpdatedAt, ps.UpdatedBy)
|
||||
return err
|
||||
}
|
||||
|
||||
// GetPlaybackState 获取房间的同步播放状态。
|
||||
func (s *Store) GetPlaybackState(ctx context.Context, roomName string) (*PlaybackState, error) {
|
||||
if s.db == nil {
|
||||
return nil, sql.ErrNoRows
|
||||
}
|
||||
var ps PlaybackState
|
||||
var playing int
|
||||
err := s.db.QueryRowContext(ctx,
|
||||
`SELECT room, position, playing, rate, current_item_id, updated_at, updated_by FROM playback_state WHERE room = ?`, roomName).
|
||||
Scan(&ps.Room, &ps.Position, &playing, &ps.Rate, &ps.CurrentItem, &ps.UpdatedAt, &ps.UpdatedBy)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
ps.Playing = playing != 0
|
||||
return &ps, nil
|
||||
}
|
||||
|
|
@ -4,6 +4,8 @@ import (
|
|||
"context"
|
||||
"database/sql"
|
||||
"errors"
|
||||
|
||||
"golang.org/x/crypto/bcrypt"
|
||||
)
|
||||
|
||||
type Service struct {
|
||||
|
|
@ -59,8 +61,8 @@ func (s *Service) JoinRoom(ctx context.Context, roomName, username string, passw
|
|||
return nil, errors.New("room is closed")
|
||||
}
|
||||
if r.PasswordHash != "" {
|
||||
if !verifyPassword(password, r.PasswordHash) {
|
||||
return nil, errors.New("invalid room password")
|
||||
if err := verifyRoomPassword(password, r.PasswordHash); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
if cnt, _ := s.store.CountMembers(ctx, roomName); cnt >= r.MaxMembers && r.MaxMembers > 0 {
|
||||
|
|
@ -73,7 +75,13 @@ func (s *Service) JoinRoom(ctx context.Context, roomName, username string, passw
|
|||
if err == nil && member.Status == MemberBanned {
|
||||
return nil, errors.New("you are banned from this room")
|
||||
}
|
||||
if err == nil && member.Status == MemberPending {
|
||||
return nil, errors.New("join request pending approval")
|
||||
}
|
||||
newMember := NewMember(roomName, username, RoleMember)
|
||||
if r.RequireApproval {
|
||||
newMember.Status = MemberPending
|
||||
}
|
||||
if s.store != nil {
|
||||
_ = s.store.UpsertMember(ctx, newMember)
|
||||
}
|
||||
|
|
@ -255,9 +263,109 @@ func (s *Service) UpdateRoomSettings(ctx context.Context, room, actorID string,
|
|||
return s.store.UpdateRoomSettings(ctx, room, newSettings)
|
||||
}
|
||||
|
||||
func verifyPassword(input, hash string) bool {
|
||||
// HashRoomPassword 使用 bcrypt 哈希房间密码。
|
||||
func HashRoomPassword(password string) (string, error) {
|
||||
if password == "" {
|
||||
return "", nil
|
||||
}
|
||||
b, err := bcrypt.GenerateFromPassword([]byte(password), bcrypt.DefaultCost)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return string(b), nil
|
||||
}
|
||||
|
||||
func verifyRoomPassword(input, hash string) error {
|
||||
if hash == "" {
|
||||
return true
|
||||
return nil
|
||||
}
|
||||
return input == hash
|
||||
if err := bcrypt.CompareHashAndPassword([]byte(hash), []byte(input)); err != nil {
|
||||
return errors.New("invalid room password")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ApproveJoin 将 pending 成员升级为 active member,需 PermManageMembers 权限。
|
||||
func (s *Service) ApproveJoin(ctx context.Context, room, actorID, targetID string) error {
|
||||
roomObj, err := s.store.GetRoom(ctx, room)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
actor, err := s.store.GetMember(ctx, room, actorID)
|
||||
if err != nil {
|
||||
return errors.New("actor not in room")
|
||||
}
|
||||
if !actor.HasPermission(PermManageMembers, roomObj.Settings) {
|
||||
return errors.New("permission denied: manage members")
|
||||
}
|
||||
target, err := s.store.GetMember(ctx, room, targetID)
|
||||
if err != nil {
|
||||
return errors.New("target not found")
|
||||
}
|
||||
if target.Status != MemberPending {
|
||||
return errors.New("member is not pending approval")
|
||||
}
|
||||
target.Status = MemberActive
|
||||
return s.store.UpsertMember(ctx, target)
|
||||
}
|
||||
|
||||
// RejectJoin 删除一个 pending 成员的加入请求。
|
||||
func (s *Service) RejectJoin(ctx context.Context, room, actorID, targetID string) error {
|
||||
roomObj, err := s.store.GetRoom(ctx, room)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
actor, err := s.store.GetMember(ctx, room, actorID)
|
||||
if err != nil {
|
||||
return errors.New("actor not in room")
|
||||
}
|
||||
if !actor.HasPermission(PermManageMembers, roomObj.Settings) {
|
||||
return errors.New("permission denied: manage members")
|
||||
}
|
||||
target, err := s.store.GetMember(ctx, room, targetID)
|
||||
if err != nil {
|
||||
return errors.New("target not found")
|
||||
}
|
||||
if target.Status != MemberPending {
|
||||
return errors.New("member is not pending approval")
|
||||
}
|
||||
return s.store.RemoveMember(ctx, room, targetID)
|
||||
}
|
||||
|
||||
// TransferOwnership 将房间所有权从当前 creator 转移给目标 active 成员。
|
||||
func (s *Service) TransferOwnership(ctx context.Context, roomName, actorID, targetID string) error {
|
||||
r, err := s.store.GetRoom(ctx, roomName)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
_ = r
|
||||
actor, err := s.store.GetMember(ctx, roomName, actorID)
|
||||
if err != nil {
|
||||
return errors.New("actor not in room")
|
||||
}
|
||||
if actor.Role != RoleCreator {
|
||||
return errors.New("only creator can transfer ownership")
|
||||
}
|
||||
target, err := s.store.GetMember(ctx, roomName, targetID)
|
||||
if err != nil {
|
||||
return errors.New("target not in room")
|
||||
}
|
||||
if target.Status != MemberActive {
|
||||
return errors.New("target must be an active member")
|
||||
}
|
||||
if target.Role == RoleCreator {
|
||||
return errors.New("target is already the creator")
|
||||
}
|
||||
// 更新房间表的 creator 字段
|
||||
if err := s.store.UpdateRoomCreator(ctx, roomName, targetID); err != nil {
|
||||
return err
|
||||
}
|
||||
// 升级 target 为 creator
|
||||
target.Role = RoleCreator
|
||||
if err := s.store.UpsertMember(ctx, target); err != nil {
|
||||
return err
|
||||
}
|
||||
// 降级 actor 为 admin
|
||||
actor.Role = RoleAdmin
|
||||
return s.store.UpsertMember(ctx, actor)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -198,6 +198,14 @@ func (s *Store) DeleteRoom(ctx context.Context, name string) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
func (s *Store) UpdateRoomCreator(ctx context.Context, name, newCreator string) error {
|
||||
if s.db == nil {
|
||||
return fmt.Errorf("no db")
|
||||
}
|
||||
_, err := s.db.ExecContext(ctx, `UPDATE rooms_v2 SET creator = ?, updated_at = ?, version = version + 1 WHERE name = ?`, newCreator, time.Now().Unix(), name)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *Store) UpsertMember(ctx context.Context, m *RoomMember) error {
|
||||
if s.db == nil {
|
||||
return fmt.Errorf("no db")
|
||||
|
|
|
|||
|
|
@ -84,7 +84,12 @@ func (s *Server) handleManageRooms(w http.ResponseWriter, r *http.Request) {
|
|||
newRoom.RequireApproval = settings.RequireApproval
|
||||
newRoom.MaxMembers = settings.MaxMembers
|
||||
if body.Password != "" {
|
||||
newRoom.PasswordHash = body.Password
|
||||
hash, err := room.HashRoomPassword(body.Password)
|
||||
if err != nil {
|
||||
http.Error(w, "hash password: "+err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
newRoom.PasswordHash = hash
|
||||
}
|
||||
created, err := s.roomSvc.CreateRoom(r.Context(), name, creator, settings)
|
||||
if err != nil {
|
||||
|
|
@ -94,7 +99,8 @@ func (s *Server) handleManageRooms(w http.ResponseWriter, r *http.Request) {
|
|||
created.Description = body.Description
|
||||
created.DisplayName = newRoom.DisplayName
|
||||
if body.Password != "" {
|
||||
created.PasswordHash = body.Password
|
||||
hash, _ := room.HashRoomPassword(body.Password)
|
||||
created.PasswordHash = hash
|
||||
}
|
||||
s.hub.createRoom(name)
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
|
|
@ -192,7 +198,12 @@ func (s *Server) handleManageRoomDetail(w http.ResponseWriter, r *http.Request)
|
|||
rm.RequireApproval = settings.RequireApproval
|
||||
rm.MaxMembers = settings.MaxMembers
|
||||
if body.Password != nil {
|
||||
rm.PasswordHash = *body.Password
|
||||
hash, err := room.HashRoomPassword(*body.Password)
|
||||
if err != nil {
|
||||
http.Error(w, "hash password: "+err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
rm.PasswordHash = hash
|
||||
}
|
||||
rm.Settings = settings
|
||||
rm.UpdatedAt = time.Now()
|
||||
|
|
@ -477,6 +488,57 @@ func actorName(r *http.Request) string {
|
|||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func (s *Server) handleManageMemberApprove(w http.ResponseWriter, r *http.Request) {
|
||||
roomName := r.PathValue("room")
|
||||
target := r.PathValue("user")
|
||||
if roomName == "" || target == "" {
|
||||
http.Error(w, "room and user required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
if err := s.roomSvc.ApproveJoin(r.Context(), roomName, actorName(r), target); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
json.NewEncoder(w).Encode(map[string]interface{}{"ok": true})
|
||||
}
|
||||
|
||||
func (s *Server) handleManageMemberReject(w http.ResponseWriter, r *http.Request) {
|
||||
roomName := r.PathValue("room")
|
||||
target := r.PathValue("user")
|
||||
if roomName == "" || target == "" {
|
||||
http.Error(w, "room and user required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
if err := s.roomSvc.RejectJoin(r.Context(), roomName, actorName(r), target); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
json.NewEncoder(w).Encode(map[string]interface{}{"ok": true})
|
||||
}
|
||||
|
||||
func (s *Server) handleManageTransferOwnership(w http.ResponseWriter, r *http.Request) {
|
||||
roomName := r.PathValue("room")
|
||||
if roomName == "" {
|
||||
http.Error(w, "room required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
var body struct {
|
||||
TargetUser string `json:"target_user"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&body); err != nil || body.TargetUser == "" {
|
||||
http.Error(w, "target_user required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
if err := s.roomSvc.TransferOwnership(r.Context(), roomName, actorName(r), body.TargetUser); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
json.NewEncoder(w).Encode(map[string]interface{}{"ok": true})
|
||||
}
|
||||
func stringOr(cur string, override *string) string {
|
||||
if override != nil {
|
||||
return *override
|
||||
|
|
|
|||
|
|
@ -0,0 +1,171 @@
|
|||
package server
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"context"
|
||||
|
||||
"sync-live/internal/room"
|
||||
)
|
||||
|
||||
// handleRoomPlaylist 处理房间播放列表的 CRUD。
|
||||
func (s *Server) handleRoomPlaylist(w http.ResponseWriter, r *http.Request) {
|
||||
roomName := r.PathValue("room")
|
||||
if roomName == "" {
|
||||
http.Error(w, "room required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
if s.roomSvc == nil || s.roomSvc.Store().DB() == nil {
|
||||
http.Error(w, "room service not configured", http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
|
||||
switch r.Method {
|
||||
case http.MethodGet:
|
||||
items, err := s.roomSvc.Store().ListPlaylistItems(r.Context(), roomName)
|
||||
if err != nil {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
if items == nil {
|
||||
items = []*room.PlaylistItem{}
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
json.NewEncoder(w).Encode(map[string]interface{}{"items": items})
|
||||
|
||||
case http.MethodPost:
|
||||
user := actorName(r)
|
||||
if user == "" {
|
||||
http.Error(w, "unauthorized", http.StatusUnauthorized)
|
||||
return
|
||||
}
|
||||
if err := s.roomSvc.EnsurePermission(r.Context(), roomName, user, room.PermControlStream); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
var body struct {
|
||||
URL string `json:"url"`
|
||||
Title string `json:"title"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&body); err != nil || body.URL == "" {
|
||||
http.Error(w, "url required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
item := &room.PlaylistItem{
|
||||
RoomName: roomName,
|
||||
URL: body.URL,
|
||||
Title: body.Title,
|
||||
AddedBy: user,
|
||||
}
|
||||
created, err := s.roomSvc.Store().AddPlaylistItem(r.Context(), item)
|
||||
if err != nil {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
// 通过 WebSocket 广播播放列表变更
|
||||
s.broadcastPlaylistUpdate(r.Context(), roomName)
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
json.NewEncoder(w).Encode(created)
|
||||
|
||||
case http.MethodDelete:
|
||||
user := actorName(r)
|
||||
if user == "" {
|
||||
http.Error(w, "unauthorized", http.StatusUnauthorized)
|
||||
return
|
||||
}
|
||||
if err := s.roomSvc.EnsurePermission(r.Context(), roomName, user, room.PermControlStream); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
itemIDStr := r.URL.Query().Get("id")
|
||||
itemID, err := strconv.ParseInt(itemIDStr, 10, 64)
|
||||
if err != nil || itemID <= 0 {
|
||||
http.Error(w, "valid item id required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
if err := s.roomSvc.Store().RemovePlaylistItem(r.Context(), roomName, itemID); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
s.broadcastPlaylistUpdate(r.Context(), roomName)
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
json.NewEncoder(w).Encode(map[string]interface{}{"ok": true})
|
||||
|
||||
default:
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
}
|
||||
}
|
||||
|
||||
// handleRoomPlayback 处理同步播放状态的读取和更新。
|
||||
func (s *Server) handleRoomPlayback(w http.ResponseWriter, r *http.Request) {
|
||||
roomName := r.PathValue("room")
|
||||
if roomName == "" {
|
||||
http.Error(w, "room required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
if s.roomSvc == nil || s.roomSvc.Store().DB() == nil {
|
||||
http.Error(w, "room service not configured", http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
|
||||
switch r.Method {
|
||||
case http.MethodGet:
|
||||
ps, err := s.roomSvc.Store().GetPlaybackState(r.Context(), roomName)
|
||||
if err != nil {
|
||||
// 无状态时返回默认值
|
||||
ps = &room.PlaybackState{Room: roomName, Rate: 1}
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
json.NewEncoder(w).Encode(ps)
|
||||
|
||||
case http.MethodPost, http.MethodPut:
|
||||
user := actorName(r)
|
||||
if user == "" {
|
||||
http.Error(w, "unauthorized", http.StatusUnauthorized)
|
||||
return
|
||||
}
|
||||
if err := s.roomSvc.EnsurePermission(r.Context(), roomName, user, room.PermControlStream); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
var body struct {
|
||||
Position float64 `json:"position"`
|
||||
Playing bool `json:"playing"`
|
||||
Rate float64 `json:"rate"`
|
||||
CurrentItem int64 `json:"current_item_id"`
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
|
||||
http.Error(w, "bad json", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
if body.Rate <= 0 {
|
||||
body.Rate = 1
|
||||
}
|
||||
ps := &room.PlaybackState{
|
||||
Room: roomName,
|
||||
Position: body.Position,
|
||||
Playing: body.Playing,
|
||||
Rate: body.Rate,
|
||||
CurrentItem: body.CurrentItem,
|
||||
UpdatedBy: user,
|
||||
}
|
||||
if err := s.roomSvc.Store().SavePlaybackState(r.Context(), ps); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
// 通过 WebSocket 广播播放状态变更
|
||||
s.hub.broadcastEvent(roomName, "playback", ps)
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
json.NewEncoder(w).Encode(ps)
|
||||
|
||||
default:
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
}
|
||||
}
|
||||
|
||||
// broadcastPlaylistUpdate 向房间广播播放列表变更事件。
|
||||
func (s *Server) broadcastPlaylistUpdate(ctx context.Context, roomName string) {
|
||||
items, _ := s.roomSvc.Store().ListPlaylistItems(ctx, roomName)
|
||||
s.hub.broadcastEvent(roomName, "playlist_update", items)
|
||||
}
|
||||
Loading…
Reference in New Issue