Files
MeBox/internal/service/cloud115_hls_proxy.go
T
truewhile 5f4b613bde 修复 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 归零,必须在切换前
  先记下目标位置
2026-10-08 22:52:31 +08:00

401 lines
12 KiB
Go
Raw 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.
package service
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"strconv"
"strings"
"sync"
"time"
"go.uber.org/zap"
"github.com/truewhile/MeBox/internal/service/cloud115"
)
var (
ErrCloud115HLSSessionNotFound = errors.New("115 hls session not found")
ErrCloud115HLSUpstreamExpired = errors.New("115 hls upstream expired")
)
const (
cloud115HLSSessionTTL = 30 * time.Minute
cloud115HLSMaxSessions = 256
cloud115HLSMaxManifest = 8 << 20
// cloud115HLSMaxEntries 单个会话最多记住多少个上游地址。key 由上游地址
// 派生(见 cloud115HLSKeyForURL),所以条目数只随「实际出现过的地址」增长,
// 正常一部片子就是几千条;这个上限只是防御异常/恶意播放列表的兜底。
cloud115HLSMaxEntries = 20000
)
type cloud115HLSSession struct {
ID string
MediaID string
Definition int
CreatedAt time.Time
ExpiresAt time.Time
mu sync.Mutex
entries map[string]string
}
// 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。
//
// 浏览器不能直接请求 115 的 m3u8:master/variant/分片的 CORS 只允许
// https://115.com,且 master 还是 HTTP。代理在服务端拉取并重写播放列表,
// 分片按 Range 流式转发。
type Cloud115HLSProxy struct {
service *Cloud115PlaybackService
client *http.Client
mu sync.Mutex
sessions map[string]*cloud115HLSSession
}
func newCloud115HLSProxy(service *Cloud115PlaybackService) *Cloud115HLSProxy {
return &Cloud115HLSProxy{
service: service,
client: &http.Client{
Timeout: 0,
CheckRedirect: func(req *http.Request, via []*http.Request) error {
if len(via) >= 6 {
return errors.New("stopped after 6 redirects")
}
return nil
},
},
sessions: make(map[string]*cloud115HLSSession),
}
}
// ServeMaster 解析指定清晰度并返回重写后的 master.m3u8。
func (p *Cloud115HLSProxy) ServeMaster(ctx context.Context, w http.ResponseWriter, r *http.Request, mediaID string, definition int) error {
if p == nil || p.service == nil {
return ErrCloud115NotApplicable
}
upstream, _, err := p.service.ResolveCloud115URL(ctx, mediaID, definition)
if err != nil {
return err
}
session := p.newSession(mediaID, definition)
fetchCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
resp, err := p.fetchUpstream(fetchCtx, upstream, "")
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("115 云端播放列表返回 HTTP %d", resp.StatusCode)
}
body, err := io.ReadAll(io.LimitReader(resp.Body, cloud115HLSMaxManifest))
if err != nil {
return err
}
baseURL := upstream
if resp.Request != nil && resp.Request.URL != nil {
baseURL = resp.Request.URL.String()
}
rewritten := p.rewriteManifest(session, string(body), baseURL, r.URL.RawQuery)
p.storeSession(session)
w.Header().Set("Content-Type", "application/vnd.apple.mpegurl")
w.Header().Set("Cache-Control", "no-store, no-cache, must-revalidate, max-age=0")
w.Header().Set("Content-Length", strconv.Itoa(len(rewritten)))
if r.Method != http.MethodHead {
_, err = io.WriteString(w, rewritten)
}
return err
}
// ServeChild 代理 variant/分片;variant 播放列表会继续重写为同源地址。
func (p *Cloud115HLSProxy) ServeChild(ctx context.Context, w http.ResponseWriter, r *http.Request, sessionID, key string) error {
session, upstream, ok := p.lookup(sessionID, key)
if !ok {
return ErrCloud115HLSSessionNotFound
}
resp, err := p.fetchUpstream(ctx, upstream, r.Header.Get("Range"))
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusForbidden || resp.StatusCode == http.StatusGone {
p.deleteSession(sessionID)
return ErrCloud115HLSUpstreamExpired
}
contentType := strings.ToLower(strings.TrimSpace(resp.Header.Get("Content-Type")))
isPlaylist := strings.Contains(contentType, "mpegurl") ||
strings.Contains(contentType, "application/vnd.apple.mpegurl")
if !isPlaylist {
// Content-Type 缺失时用前 4KB 内容嗅探是否为播放列表;无论结果如何
// 都要把已读前缀接回 body,避免分片/播放列表丢开头字节。
body, readErr := io.ReadAll(io.LimitReader(resp.Body, 4096))
if readErr != nil {
return readErr
}
if strings.HasPrefix(strings.TrimSpace(string(body)), "#EXTM3U") {
isPlaylist = true
}
resp.Body = io.NopCloser(io.MultiReader(strings.NewReader(string(body)), resp.Body))
}
copyUpstreamHeaders(w, resp, isPlaylist)
if isPlaylist {
body, err := io.ReadAll(io.LimitReader(resp.Body, cloud115HLSMaxManifest))
if err != nil {
return err
}
baseURL := upstream
if resp.Request != nil && resp.Request.URL != nil {
baseURL = resp.Request.URL.String()
}
rewritten := p.rewriteManifest(session, string(body), baseURL, r.URL.RawQuery)
w.Header().Set("Content-Type", "application/vnd.apple.mpegurl")
w.Header().Set("Cache-Control", "no-store, no-cache, must-revalidate, max-age=0")
w.Header().Set("Content-Length", strconv.Itoa(len(rewritten)))
w.WriteHeader(resp.StatusCode)
if r.Method != http.MethodHead {
_, err = io.WriteString(w, rewritten)
}
return err
}
w.WriteHeader(resp.StatusCode)
if r.Method != http.MethodHead {
_, err = io.Copy(w, resp.Body)
}
return err
}
func (p *Cloud115HLSProxy) newSession(mediaID string, definition int) *cloud115HLSSession {
now := time.Now()
return &cloud115HLSSession{
ID: cloud115.RandomString(24),
MediaID: mediaID,
Definition: definition,
CreatedAt: now,
ExpiresAt: now.Add(cloud115HLSSessionTTL),
entries: make(map[string]string),
}
}
func (p *Cloud115HLSProxy) storeSession(session *cloud115HLSSession) {
if p == nil || session == nil {
return
}
p.mu.Lock()
defer p.mu.Unlock()
now := time.Now()
if len(p.sessions) >= cloud115HLSMaxSessions {
for id, existing := range p.sessions {
if now.After(existing.ExpiresAt) {
delete(p.sessions, id)
}
}
}
if len(p.sessions) >= cloud115HLSMaxSessions {
// ExpiresAt 恒为「最近一次活跃时间 + TTL」(lookup 会滑动续期),
// 因此最小者即最久未活跃的会话,淘汰它对正在播放的影响最小。
evictID := ""
var evictAt time.Time
for id, existing := range p.sessions {
if evictID == "" || existing.ExpiresAt.Before(evictAt) {
evictID, evictAt = id, existing.ExpiresAt
}
}
if evictID != "" {
delete(p.sessions, evictID)
}
}
session.ExpiresAt = now.Add(cloud115HLSSessionTTL)
p.sessions[session.ID] = session
}
func (p *Cloud115HLSProxy) lookup(sessionID, key string) (*cloud115HLSSession, string, bool) {
if p == nil {
return nil, "", false
}
p.mu.Lock()
session := p.sessions[sessionID]
if session != nil {
if time.Now().After(session.ExpiresAt) {
delete(p.sessions, sessionID)
session = nil
} else {
// 滑动续期:VOD 点播整个播放期间只有分片请求,master 只在起播时取
// 一次,不续期的话会话会在 TTL 后被删,长视频播放到一半必然断流。
session.ExpiresAt = time.Now().Add(cloud115HLSSessionTTL)
}
}
p.mu.Unlock()
if session == nil {
return nil, "", false
}
session.mu.Lock()
upstream, ok := session.entries[key]
session.mu.Unlock()
return session, upstream, ok
}
func (p *Cloud115HLSProxy) deleteSession(sessionID string) {
if p == nil {
return
}
p.mu.Lock()
delete(p.sessions, sessionID)
p.mu.Unlock()
}
func (p *Cloud115HLSProxy) fetchUpstream(ctx context.Context, rawURL, rangeHeader string) (*http.Response, error) {
parsed, err := url.Parse(strings.TrimSpace(rawURL))
if err != nil || parsed == nil || !isAllowed115UpstreamHost(parsed.Hostname()) {
return nil, fmt.Errorf("115 hls upstream host not allowed: %s", rawURL)
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, rawURL, nil)
if err != nil {
return nil, err
}
req.Header.Set("User-Agent", cloud115.DefaultUA)
if strings.TrimSpace(rangeHeader) != "" {
req.Header.Set("Range", rangeHeader)
}
return p.client.Do(req)
}
func isAllowed115UpstreamHost(host string) bool {
host = strings.ToLower(strings.TrimSpace(host))
if host == "" {
return false
}
for _, suffix := range []string{".115.com", ".115cdn.com", ".115cdn.net"} {
if strings.HasSuffix(host, suffix) {
return true
}
}
return host == "115.com" || host == "115cdn.com" || host == "115cdn.net"
}
func (p *Cloud115HLSProxy) rewriteManifest(session *cloud115HLSSession, text, baseURL, rawQuery string) string {
if session == nil {
return text
}
lines := strings.SplitAfter(text, "\n")
for i, line := range lines {
trimmed := strings.TrimSpace(line)
if trimmed == "" {
continue
}
if strings.HasPrefix(trimmed, "#") {
if strings.HasPrefix(trimmed, "#EXT-X-MEDIA:") ||
strings.HasPrefix(trimmed, "#EXT-X-KEY:") ||
strings.HasPrefix(trimmed, "#EXT-X-MAP:") {
lines[i] = replaceManifestURI(line, func(uri string) string {
return p.proxyURL(session, resolveManifestURL(baseURL, uri), rawQuery)
})
}
continue
}
lines[i] = p.proxyURL(session, resolveManifestURL(baseURL, trimmed), rawQuery) + lineEnding(line)
}
return strings.Join(lines, "")
}
func (p *Cloud115HLSProxy) proxyURL(session *cloud115HLSSession, upstream, rawQuery string) string {
key := cloud115HLSKeyForURL(upstream)
session.mu.Lock()
_, 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
}
func childProxyQuery(rawQuery, mediaID string) string {
values, _ := url.ParseQuery(rawQuery)
keep := url.Values{}
for _, key := range []string{"token", "api_key", "apiKey", "ApiKey", "profile_id", "profile_pin_token"} {
if value := strings.TrimSpace(values.Get(key)); value != "" {
keep.Set(key, value)
}
}
if strings.TrimSpace(mediaID) != "" {
keep.Set("media_id", mediaID)
}
return keep.Encode()
}
func replaceManifestURI(line string, replace func(string) string) string {
const marker = `URI="`
idx := strings.Index(line, marker)
if idx < 0 {
return line
}
start := idx + len(marker)
end := strings.Index(line[start:], `"`)
if end < 0 {
return line
}
end += start
return line[:start] + replace(line[start:end]) + line[end:]
}
func resolveManifestURL(baseURL, raw string) string {
base, baseErr := url.Parse(strings.TrimSpace(baseURL))
ref, refErr := url.Parse(strings.TrimSpace(raw))
if baseErr != nil || refErr != nil || base == nil || ref == nil {
return raw
}
return base.ResolveReference(ref).String()
}
func lineEnding(line string) string {
if strings.HasSuffix(line, "\r\n") {
return "\r\n"
}
if strings.HasSuffix(line, "\n") {
return "\n"
}
return ""
}
func copyUpstreamHeaders(w http.ResponseWriter, resp *http.Response, playlist bool) {
for _, key := range []string{"Content-Type", "Content-Length", "Content-Range", "Accept-Ranges", "ETag", "Last-Modified"} {
if value := resp.Header.Get(key); value != "" {
w.Header().Set(key, value)
}
}
if playlist {
w.Header().Set("Cache-Control", "no-store, no-cache, must-revalidate, max-age=0")
} else if w.Header().Get("Cache-Control") == "" {
w.Header().Set("Cache-Control", "public, max-age=3600")
}
}