mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-11 07:46:37 +08:00
修复 115 云转码播放中途断流后直接掉到本地转码
115 云 HLS 会话只存在服务端内存里:服务重新部署、会话过期(30 分钟无活动)、 上游分片地址失效,都会让播放中的客户端在后续分片上拿到 404/410。原实现此时 直接退回本地转码,而低配宿主机本地转码追不上播放速度(实测 2 vCPU 只有 0.62x),用户看到的就是播放报废。 - 云 HLS 中途 fatal 时先重建会话续播(10 分钟内最多 3 次),仍失败才退回本地 - manifestLoadingMaxRetry 1 -> 4:master.m3u8 每次都要现调 115 接口,一次抖动 不该直接判死刑 - 播放列表重写改用上游地址哈希做 key:原来用自增序号,每次重写 entries 就膨胀 一轮、分片编号漂移(实测 762 段列表重写几次后 key 涨到 e2288..e3049) - 会话条目数加上限,异常播放列表不会无限占用内存 - 客户端断开(seek / 切清晰度取消在途请求)不再记成 502,改记 499 与 nginx 对齐;会话失效返回带 code 的 404/410,便于前端判断需要重建会话 - 修正进度/时长上报:云 HLS 时间轴起点恒为 0,不应再叠加 hlsStartSec(历史上 把 7624s 的片子记成 8346s 并标成「已看完」) - 切回直连 / 原画时保住播放位置:HLS 卸载会把 currentTime 归零,必须在切换前 先记下目标位置
This commit is contained in:
@@ -2,6 +2,8 @@ package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -12,6 +14,8 @@ import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
|
||||
"github.com/truewhile/MeBox/internal/service/cloud115"
|
||||
)
|
||||
|
||||
@@ -24,6 +28,10 @@ const (
|
||||
cloud115HLSSessionTTL = 30 * time.Minute
|
||||
cloud115HLSMaxSessions = 256
|
||||
cloud115HLSMaxManifest = 8 << 20
|
||||
// cloud115HLSMaxEntries 单个会话最多记住多少个上游地址。key 由上游地址
|
||||
// 派生(见 cloud115HLSKeyForURL),所以条目数只随「实际出现过的地址」增长,
|
||||
// 正常一部片子就是几千条;这个上限只是防御异常/恶意播放列表的兜底。
|
||||
cloud115HLSMaxEntries = 20000
|
||||
)
|
||||
|
||||
type cloud115HLSSession struct {
|
||||
@@ -35,7 +43,17 @@ type cloud115HLSSession struct {
|
||||
|
||||
mu sync.Mutex
|
||||
entries map[string]string
|
||||
next int
|
||||
}
|
||||
|
||||
// cloud115HLSKeyForURL 由上游地址派生稳定的分片 key。
|
||||
//
|
||||
// 早期实现用自增序号(e1、e2…)在每次重写播放列表时重新编号:同一分片每次
|
||||
// 重写都会拿到新 key,entries 随重写次数线性膨胀(实测一个 762 段的列表反复
|
||||
// 重写后 key 已经涨到 e2288…e3049),客户端拿到的分片编号也会无谓漂移。
|
||||
// 改成地址哈希后,同一分片在任何一次重写里都是同一个 key。
|
||||
func cloud115HLSKeyForURL(upstream string) string {
|
||||
sum := sha256.Sum256([]byte(upstream))
|
||||
return hex.EncodeToString(sum[:12])
|
||||
}
|
||||
|
||||
// Cloud115HLSProxy 把 115 云端 HLS 转成 MeBox 同源 HLS。
|
||||
@@ -302,12 +320,19 @@ func (p *Cloud115HLSProxy) rewriteManifest(session *cloud115HLSSession, text, ba
|
||||
}
|
||||
|
||||
func (p *Cloud115HLSProxy) proxyURL(session *cloud115HLSSession, upstream, rawQuery string) string {
|
||||
key := cloud115HLSKeyForURL(upstream)
|
||||
session.mu.Lock()
|
||||
session.next++
|
||||
key := "e" + strconv.Itoa(session.next)
|
||||
session.entries[key] = upstream
|
||||
_, known := session.entries[key]
|
||||
if !known && len(session.entries) < cloud115HLSMaxEntries {
|
||||
session.entries[key] = upstream
|
||||
known = true
|
||||
}
|
||||
mediaID := session.MediaID
|
||||
session.mu.Unlock()
|
||||
if !known && p != nil && p.service != nil && p.service.log != nil {
|
||||
p.service.log.Warn("cloud115 hls: session entry cap reached, segment will not resolve",
|
||||
zap.String("session", session.ID), zap.Int("cap", cloud115HLSMaxEntries))
|
||||
}
|
||||
|
||||
query := childProxyQuery(rawQuery, mediaID)
|
||||
return "/api/cloud115/hls/" + url.PathEscape(session.ID) + "/" + url.PathEscape(key) + "?" + query
|
||||
|
||||
@@ -25,7 +25,8 @@ func TestCloud115HLSProxyRewriteManifest(t *testing.T) {
|
||||
if strings.Contains(rewritten, "cpats01.115.com") {
|
||||
t.Fatalf("upstream URL leaked into rewritten manifest: %s", rewritten)
|
||||
}
|
||||
if !strings.Contains(rewritten, "/api/cloud115/hls/sess/e1?") {
|
||||
wantKey := cloud115HLSKeyForURL("https://cpats01.115.com/a.m3u8?u=1&se=2")
|
||||
if !strings.Contains(rewritten, "/api/cloud115/hls/sess/"+wantKey+"?") {
|
||||
t.Fatalf("proxy URL missing: %s", rewritten)
|
||||
}
|
||||
if !strings.Contains(rewritten, "media_id=media-1") || !strings.Contains(rewritten, "token=t") {
|
||||
@@ -36,6 +37,32 @@ func TestCloud115HLSProxyRewriteManifest(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestCloud115HLSProxyRewriteIsStableAcrossRewrites(t *testing.T) {
|
||||
proxy := &Cloud115HLSProxy{sessions: map[string]*cloud115HLSSession{}}
|
||||
session := &cloud115HLSSession{
|
||||
ID: "sess",
|
||||
MediaID: "media-1",
|
||||
entries: map[string]string{},
|
||||
}
|
||||
manifest := "#EXTM3U\n#EXT-X-TARGETDURATION:10\n" +
|
||||
"https://cpats01.115.com/seg0.ts?x=0\nhttps://cpats01.115.com/seg1.ts?x=1\n"
|
||||
|
||||
first := proxy.rewriteManifest(session, manifest, "https://cpats01.115.com/v.m3u8", "media_id=media-1")
|
||||
entriesAfterFirst := len(session.entries)
|
||||
second := proxy.rewriteManifest(session, manifest, "https://cpats01.115.com/v.m3u8", "media_id=media-1")
|
||||
|
||||
// 同一分片在每次重写里都必须是同一个 key,客户端拿到的地址才不会无谓漂移。
|
||||
if first != second {
|
||||
t.Fatalf("rewrite is not stable:\nfirst = %s\nsecond = %s", first, second)
|
||||
}
|
||||
if len(session.entries) != entriesAfterFirst {
|
||||
t.Fatalf("entries grew on rewrite: %d -> %d", entriesAfterFirst, len(session.entries))
|
||||
}
|
||||
if entriesAfterFirst != 2 {
|
||||
t.Fatalf("session entries = %d, want 2", entriesAfterFirst)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCloud115HLSProxyRewriteKeyURI(t *testing.T) {
|
||||
proxy := &Cloud115HLSProxy{sessions: map[string]*cloud115HLSSession{}}
|
||||
session := &cloud115HLSSession{
|
||||
@@ -151,7 +178,6 @@ func TestCloud115HLSProxyServeChildRewritesVariant(t *testing.T) {
|
||||
MediaID: "media-1",
|
||||
ExpiresAt: time.Now().Add(time.Hour),
|
||||
entries: map[string]string{"e1": "https://cpats01.115.com/v.m3u8"},
|
||||
next: 1,
|
||||
}
|
||||
proxy.sessions[session.ID] = session
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/cloud115/hls/sess/e1?media_id=media-1&token=t", nil)
|
||||
@@ -162,7 +188,8 @@ func TestCloud115HLSProxyServeChildRewritesVariant(t *testing.T) {
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want 200", rec.Code)
|
||||
}
|
||||
if !strings.Contains(rec.Body.String(), "/api/cloud115/hls/sess/e2?") {
|
||||
wantKey := cloud115HLSKeyForURL("https://cpats01.115.com/seg.ts?x=1")
|
||||
if !strings.Contains(rec.Body.String(), "/api/cloud115/hls/sess/"+wantKey+"?") {
|
||||
t.Fatalf("rewritten variant missing proxy segment: %s", rec.Body.String())
|
||||
}
|
||||
if strings.Contains(rec.Body.String(), "cpats01.115.com") {
|
||||
@@ -192,7 +219,6 @@ func TestCloud115HLSProxyServeChildStreamsRange(t *testing.T) {
|
||||
MediaID: "media-1",
|
||||
ExpiresAt: time.Now().Add(time.Hour),
|
||||
entries: map[string]string{"e1": "https://cpats01.115.com/seg.ts?x=1"},
|
||||
next: 1,
|
||||
}
|
||||
proxy.sessions[session.ID] = session
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/cloud115/hls/sess/e1?media_id=media-1&token=t", nil)
|
||||
|
||||
Reference in New Issue
Block a user