Compare commits

...

1 Commits

Author SHA1 Message Date
truewhile 34ccf14cda 优化视频头尾跳过功能 2026-09-23 16:07:28 +08:00
26 changed files with 1992 additions and 47 deletions
+23 -3
View File
@@ -22,7 +22,9 @@ import (
func playbackSegmentsHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
autoSkip := resolveAutoSkipFlag(c, svc)
source := resolveSegmentSource(c, svc)
segments := []service.SegmentView{}
pending := false
m, err := findMediaForPlaybackEndpoint(c, svc, c.Param("id"))
if err != nil || m == nil || !mediaVisibleForRequest(c, svc, m) {
@@ -31,14 +33,22 @@ func playbackSegmentsHandler(svc *service.Container) gin.HandlerFunc {
}
// 远程 Emby 挂载的条目是上游库的投影,本地没有可查询的外部 ID 关联。
if svc.Segments != nil && !service.IsEmbyRemoteID(m.ID) {
rows, listErr := svc.Segments.ListForPlayback(c.Request.Context(), m)
result, listErr := svc.Segments.SegmentsForPlayback(c.Request.Context(), m, source)
if listErr != nil && svc.Log != nil {
svc.Log.Debug("list media segments failed",
zap.String("media_id", m.ID), zap.Error(listErr))
}
segments = service.ToSegmentViews(rows)
segments = service.ToSegmentViews(result.Segments)
pending = result.Pending
}
c.JSON(http.StatusOK, gin.H{"segments": segments, "auto_skip": autoSkip})
// pending 告诉客户端「章节提取还在后台跑,过几秒再拉一次」;提取完成时
// 如果片头还没播完,跳过按钮就会自己出现,已经过了片头则不会提示。
c.JSON(http.StatusOK, gin.H{
"segments": segments,
"auto_skip": autoSkip,
"pending": pending,
"source": source,
})
}
}
@@ -53,3 +63,13 @@ func resolveAutoSkipFlag(c *gin.Context, svc *service.Container) bool {
}
return profile.SkipIntro
}
// resolveSegmentSource reads the「片头片尾数据来源」choice off the active profile.
// PIN-locked profiles fall back to auto rather than leaking the profile's setting.
func resolveSegmentSource(c *gin.Context, svc *service.Container) string {
profile, locked := selectedPlayProfile(c, svc)
if locked || profile == nil {
return service.SegmentSourceAuto
}
return service.NormalizeSegmentSource(profile.SegmentSource)
}
+42
View File
@@ -0,0 +1,42 @@
package model
import "time"
// MediaProbe 是一次 ffprobe 全量探测(容器 / 轨道 / 内嵌章节)的结果缓存。
//
// 它存在的理由:探测一次要 3~4 秒——远端直链更慢,因为要跨洋跑三次 HTTP
// 事务。播放链路绝不能等它,所以第一次播放只起后台任务,结果落库后由后续请求
// 与「跳过片头」的章节数据直接读库。
//
// Payload 刻意只保存裁剪后的字段:ffprobe 原始输出里的 format.filename 是解析
// 后的播放直链(带网盘签名与 pickcode),原样落库等于把可直接下载的链接留在
// 数据库里,所以只保留与技术信息有关的字段。
type MediaProbe struct {
Base
MediaID string `gorm:"uniqueIndex;size:128;not null" json:"media_id"`
// Signature 是「探的是哪个文件」的指纹(哈希):本地文件取路径 + 大小 +
// 修改时间,STRM / 云盘取固化的播放目标。文件换了就说明缓存不再对应当前
// 内容,需要重探。
Signature string `gorm:"size:64" json:"signature,omitempty"`
// Source 记录输入形态:local | strm。
Source string `gorm:"size:16" json:"source,omitempty"`
Container string `gorm:"size:64" json:"container,omitempty"`
DurationSec int `json:"duration_sec"`
BitRate int64 `json:"bit_rate,omitempty"`
Width int `json:"width,omitempty"`
Height int `json:"height,omitempty"`
VideoCodec string `gorm:"size:64" json:"video_codec,omitempty"`
AudioCodec string `gorm:"size:64" json:"audio_codec,omitempty"`
VideoStreams int `json:"video_streams"`
AudioStreams int `json:"audio_streams"`
SubtitleStreams int `json:"subtitle_streams"`
ChapterCount int `json:"chapter_count"`
// Payload 是供详情页展示的裁剪后 JSON(容器 + 每路轨道 + 章节)。
Payload string `gorm:"type:text" json:"payload,omitempty"`
// ProbedAt 是最近一次探测的时刻。LastError 非空表示这次探测失败;失败只更新
// 这两个字段,不会覆盖此前成功的 Payload 与已经落库的章节片段。
ProbedAt time.Time `json:"probed_at"`
LastError string `gorm:"size:512" json:"last_error,omitempty"`
}
+1
View File
@@ -40,6 +40,7 @@ func AllModels() []interface{} {
&PlaybackHistory{},
&MediaSegment{},
&MediaSegmentFetch{},
&MediaProbe{},
&Favorite{},
&Playlist{},
&PlaylistItem{},
+17 -14
View File
@@ -53,18 +53,21 @@ type PlaylistItem struct {
// AllowedLibraryIDs is a JSON array of library UUIDs (empty = all).
type PlayProfile struct {
Base
UserID string `gorm:"index;size:36;not null" json:"user_id"`
Name string `gorm:"size:64;not null" json:"name"`
IsDefault bool `gorm:"default:false" json:"is_default"`
ContentRatingLimit string `gorm:"size:16" json:"content_rating_limit,omitempty"`
AllowAdult bool `gorm:"default:false" json:"allow_adult"`
RequirePIN bool `gorm:"default:false" json:"require_pin"`
PINHash string `gorm:"size:128" json:"-"`
PreferredSubtitleLang string `gorm:"size:16" json:"preferred_subtitle_lang,omitempty"`
PreferredAudioLang string `gorm:"size:16" json:"preferred_audio_lang,omitempty"`
AutoplayNext bool `gorm:"default:true" json:"autoplay_next"`
SkipIntro bool `gorm:"default:false" json:"skip_intro"`
AllowedLibraryIDs string `gorm:"type:text;default:'[]'" json:"allowed_library_ids"`
TotalWatchTime int64 `gorm:"default:0" json:"total_watch_time"`
LastActiveAt *time.Time `json:"last_active_at,omitempty"`
UserID string `gorm:"index;size:36;not null" json:"user_id"`
Name string `gorm:"size:64;not null" json:"name"`
IsDefault bool `gorm:"default:false" json:"is_default"`
ContentRatingLimit string `gorm:"size:16" json:"content_rating_limit,omitempty"`
AllowAdult bool `gorm:"default:false" json:"allow_adult"`
RequirePIN bool `gorm:"default:false" json:"require_pin"`
PINHash string `gorm:"size:128" json:"-"`
PreferredSubtitleLang string `gorm:"size:16" json:"preferred_subtitle_lang,omitempty"`
PreferredAudioLang string `gorm:"size:16" json:"preferred_audio_lang,omitempty"`
AutoplayNext bool `gorm:"default:true" json:"autoplay_next"`
SkipIntro bool `gorm:"default:false" json:"skip_intro"`
// SegmentSource 决定片头/片尾数据从哪来:auto(优先文件内嵌章节,没有可用
// 章节时回落到 TheIntroDB)/ theintrodb / ffprobe。空值按 auto 处理。
SegmentSource string `gorm:"size:16;default:'auto'" json:"segment_source,omitempty"`
AllowedLibraryIDs string `gorm:"type:text;default:'[]'" json:"allowed_library_ids"`
TotalWatchTime int64 `gorm:"default:0" json:"total_watch_time"`
LastActiveAt *time.Time `json:"last_active_at,omitempty"`
}
@@ -0,0 +1,61 @@
package repository
import (
"context"
"errors"
"time"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"github.com/truewhile/MeBox/internal/model"
)
// MediaProbeRepository 持久化 ffprobe 全量探测的结果缓存。
type MediaProbeRepository struct{ db *gorm.DB }
// Get 返回某媒体的探测缓存,未探测过时返回 (nil, nil)。
func (r *MediaProbeRepository) Get(ctx context.Context, mediaID string) (*model.MediaProbe, error) {
var row model.MediaProbe
err := r.db.WithContext(ctx).
Where("media_id = ?", mediaID).
First(&row).Error
if errors.Is(err, gorm.ErrRecordNotFound) {
return nil, nil
}
if err != nil {
return nil, err
}
return &row, nil
}
// Upsert 写入一次成功的探测结果(含裁剪后的媒体信息)。
func (r *MediaProbeRepository) Upsert(ctx context.Context, row *model.MediaProbe) error {
onConflict := clause.OnConflict{
Columns: []clause.Column{{Name: "media_id"}},
DoUpdates: clause.AssignmentColumns([]string{
"signature", "source",
"container", "duration_sec", "bit_rate",
"width", "height", "video_codec", "audio_codec",
"video_streams", "audio_streams", "subtitle_streams", "chapter_count",
"payload", "probed_at", "last_error",
"deleted_at",
}),
}
return r.db.WithContext(ctx).Clauses(onConflict).Create(row).Error
}
// MarkFailure 记录一次失败的探测。它只更新「时间 + 错误信息」,刻意不碰
// payload 与其它字段:一次失败的重探不该把上一次成功拿到的媒体信息抹掉。
func (r *MediaProbeRepository) MarkFailure(ctx context.Context, mediaID, message string, probedAt time.Time) error {
row := model.MediaProbe{MediaID: mediaID, ProbedAt: probedAt, LastError: message}
onConflict := clause.OnConflict{
Columns: []clause.Column{{Name: "media_id"}},
DoUpdates: clause.Assignments(map[string]any{
"probed_at": probedAt,
"last_error": message,
"deleted_at": nil,
}),
}
return r.db.WithContext(ctx).Clauses(onConflict).Create(&row).Error
}
@@ -0,0 +1,105 @@
package repository
import (
"testing"
"time"
"github.com/truewhile/MeBox/internal/model"
)
func TestMediaProbeUpsertReplacesAndGetReturnsLatest(t *testing.T) {
repos := newSegmentTestRepos(t)
ctx := t.Context()
first := &model.MediaProbe{
MediaID: "m-1", Signature: "sig-1", Source: "local",
Container: "matroska,webm", DurationSec: 1451, ChapterCount: 2,
Payload: `{"container":"matroska,webm"}`, ProbedAt: time.Now(),
}
if err := repos.MediaProbe.Upsert(ctx, first); err != nil {
t.Fatalf("upsert #1: %v", err)
}
second := &model.MediaProbe{
MediaID: "m-1", Signature: "sig-2", Source: "strm",
Container: "mp4", DurationSec: 900, ChapterCount: 0,
Payload: `{"container":"mp4"}`, ProbedAt: time.Now(),
}
if err := repos.MediaProbe.Upsert(ctx, second); err != nil {
t.Fatalf("upsert #2: %v", err)
}
got, err := repos.MediaProbe.Get(ctx, "m-1")
if err != nil {
t.Fatal(err)
}
if got == nil {
t.Fatal("probe row missing")
}
if got.Container != "mp4" || got.DurationSec != 900 || got.Signature != "sig-2" || got.Payload != `{"container":"mp4"}` {
t.Fatalf("probe = %#v, want the second upsert to win", got)
}
// 重复探测不能累积重复行:media_id 是唯一索引。
var count int64
if err := repos.DB.Model(&model.MediaProbe{}).Where("media_id = ?", "m-1").Count(&count).Error; err != nil {
t.Fatal(err)
}
if count != 1 {
t.Fatalf("rows = %d, want 1", count)
}
}
func TestMediaProbeGetReturnsNilWhenMissing(t *testing.T) {
repos := newSegmentTestRepos(t)
got, err := repos.MediaProbe.Get(t.Context(), "nope")
if err != nil {
t.Fatalf("Get: %v", err)
}
if got != nil {
t.Fatalf("probe = %#v, want nil", got)
}
}
func TestMediaProbeMarkFailureKeepsPreviousPayload(t *testing.T) {
repos := newSegmentTestRepos(t)
ctx := t.Context()
if err := repos.MediaProbe.Upsert(ctx, &model.MediaProbe{
MediaID: "m-1", Container: "matroska,webm", DurationSec: 1451,
ChapterCount: 2, Payload: `{"container":"matroska,webm"}`, ProbedAt: time.Now(),
}); err != nil {
t.Fatal(err)
}
if err := repos.MediaProbe.MarkFailure(ctx, "m-1", "ffprobe full: exit status 1", time.Now()); err != nil {
t.Fatal(err)
}
got, err := repos.MediaProbe.Get(ctx, "m-1")
if err != nil {
t.Fatal(err)
}
if got == nil || got.LastError != "ffprobe full: exit status 1" {
t.Fatalf("probe = %#v, want the recorded failure", got)
}
// 一次失败的重探不该把上一次成功拿到的媒体信息抹掉。
if got.Container != "matroska,webm" || got.DurationSec != 1451 || got.Payload == "" || got.ChapterCount != 2 {
t.Fatalf("a failed re-probe wiped the previous summary: %#v", got)
}
}
func TestMediaProbeMarkFailureInsertsRowWhenAbsent(t *testing.T) {
repos := newSegmentTestRepos(t)
ctx := t.Context()
// 没有历史成功记录时,失败也必须落一行:否则冷却期没有时间戳,
// 每次播放都会重跑一次注定失败的探测。
if err := repos.MediaProbe.MarkFailure(ctx, "m-2", "boom", time.Now()); err != nil {
t.Fatal(err)
}
got, err := repos.MediaProbe.Get(ctx, "m-2")
if err != nil {
t.Fatal(err)
}
if got == nil || got.LastError != "boom" || got.ProbedAt.IsZero() {
t.Fatalf("probe = %#v, want a failure row with a timestamp", got)
}
}
@@ -23,6 +23,18 @@ func (r *MediaSegmentRepository) ListByMedia(ctx context.Context, mediaID string
return rows, err
}
// ListByMediaSource returns the segments contributed by one source only, ordered
// by start. Used to read the ffprobe-extracted chapters independently of the
// community-database rows, so the player can pick between them.
func (r *MediaSegmentRepository) ListByMediaSource(ctx context.Context, mediaID, source string) ([]model.MediaSegment, error) {
rows := make([]model.MediaSegment, 0, 4)
err := r.db.WithContext(ctx).
Where("media_id = ? AND source = ?", mediaID, source).
Order("start_ms asc").
Find(&rows).Error
return rows, err
}
// ReplaceForMedia swaps the segments contributed by one source in a single
// transaction, so a provider refresh can never leave a half-updated set.
//
+2
View File
@@ -16,6 +16,7 @@ type Container struct {
Series *SeriesRepository
History *HistoryRepository
MediaSegment *MediaSegmentRepository
MediaProbe *MediaProbeRepository
Favorite *FavoriteRepository
Playlist *PlaylistRepository
Setting *SettingRepository
@@ -47,6 +48,7 @@ func New(db *gorm.DB) *Container {
Series: &SeriesRepository{db: db},
History: &HistoryRepository{db: db},
MediaSegment: &MediaSegmentRepository{db: db},
MediaProbe: &MediaProbeRepository{db: db},
Favorite: &FavoriteRepository{db: db},
Playlist: &PlaylistRepository{db: db},
Setting: &SettingRepository{db: db},
+283
View File
@@ -0,0 +1,283 @@
package service
import (
"context"
"encoding/json"
"errors"
"fmt"
"os/exec"
"strconv"
"strings"
"time"
)
// ffprobeFullTimeout 是一次全量探测(含章节)的超时。远端直链实测 3~4 秒,
// 留足余量的同时不能让一个坏源把工作协程永久占住。
const ffprobeFullTimeout = 60 * time.Second
// ProbeInput 是喂给 ffprobe 的输入:本地文件给路径,STRM / 云盘给已解析的最终
// 直链加绑定请求头。解析由调用方负责(复用播放链路的换链逻辑),否则会踩到
// 115 CDN 的防盗链 403。
type ProbeInput struct {
Source string
Headers map[string]string
}
// ProbeStream 是一路轨道的关键字段,供详情页展示。
type ProbeStream struct {
Index int `json:"index"`
Type string `json:"type"`
Codec string `json:"codec,omitempty"`
Profile string `json:"profile,omitempty"`
Language string `json:"language,omitempty"`
Title string `json:"title,omitempty"`
Width int `json:"width,omitempty"`
Height int `json:"height,omitempty"`
PixFmt string `json:"pix_fmt,omitempty"`
FrameRate string `json:"frame_rate,omitempty"`
Channels int `json:"channels,omitempty"`
ChannelLayout string `json:"channel_layout,omitempty"`
SampleRate int `json:"sample_rate,omitempty"`
BitRate int64 `json:"bit_rate,omitempty"`
Default bool `json:"default,omitempty"`
Forced bool `json:"forced,omitempty"`
}
// ProbeChapter 是一个内嵌章节区间。
type ProbeChapter struct {
Index int `json:"index"`
StartMs int64 `json:"start_ms"`
EndMs int64 `json:"end_ms"`
Title string `json:"title,omitempty"`
}
// FullProbeResult 是一次全量探测的裁剪结果。
type FullProbeResult struct {
Container string
DurationSec int
BitRate int64
Streams []ProbeStream
Chapters []ProbeChapter
}
// mediaProbePayload 是落库的媒体信息结构。它只包含白名单字段——刻意不接受
// ffprobe 的原始 JSON,因为其中的 format.filename 是播放直链(含签名与
// pickcode),原样保存会外泄。
type mediaProbePayload struct {
Container string `json:"container,omitempty"`
DurationSec int `json:"duration_sec"`
BitRate int64 `json:"bit_rate,omitempty"`
Streams []ProbeStream `json:"streams"`
Chapters []ProbeChapter `json:"chapters,omitempty"`
}
// StreamsOfType 返回指定类型的轨道,供详情页分组展示。
func (r *FullProbeResult) StreamsOfType(kind string) []ProbeStream {
if r == nil {
return nil
}
out := make([]ProbeStream, 0, len(r.Streams))
for _, s := range r.Streams {
if s.Type == kind {
out = append(out, s)
}
}
return out
}
// PayloadJSON 序列化落库用的媒体信息。
func (r *FullProbeResult) PayloadJSON() (string, error) {
if r == nil {
return "", nil
}
streams := r.Streams
if streams == nil {
streams = []ProbeStream{}
}
body, err := json.Marshal(mediaProbePayload{
Container: r.Container,
DurationSec: r.DurationSec,
BitRate: r.BitRate,
Streams: streams,
Chapters: r.Chapters,
})
if err != nil {
return "", err
}
return string(body), nil
}
// ProbeFull 跑一次全量探测:容器信息 + 全部轨道 + 内嵌章节。
//
// 与 Probe 的区别:Probe 只取扫描需要的几个字段、对着媒体行的本地路径跑;
// ProbeFull 接受调用方解析好的输入(远端直链 + 绑定请求头),并额外抓章节。
func (f *FFprobeService) ProbeFull(ctx context.Context, input ProbeInput) (*FullProbeResult, error) {
if f == nil || f.cfg == nil {
return nil, errors.New("ffprobe service nil")
}
source := strings.TrimSpace(input.Source)
if source == "" {
return nil, errors.New("empty probe source")
}
token, err := f.acquire(ctx)
if err != nil {
return nil, err
}
defer f.release(token)
bin, err := resolveLocalExecutable(f.cfg.App.FFprobePath, "ffprobe")
if err != nil {
return nil, fmt.Errorf("ffprobe unavailable: %w", err)
}
probeCtx, cancel := context.WithTimeout(ctx, ffprobeFullTimeout)
defer cancel()
args := []string{"-v", "error"}
if headerText := ffmpegHeaderText(input.Headers); headerText != "" {
args = append(args, "-headers", headerText)
}
args = append(args,
"-print_format", "json",
"-show_format",
"-show_streams",
"-show_chapters",
source,
)
out, err := exec.CommandContext(probeCtx, bin, args...).Output() // #nosec G204 -- bin is resolved by resolveLocalExecutable before execution.
if err != nil {
return nil, fmt.Errorf("ffprobe full: %w", err)
}
return parseFullProbeJSON(out)
}
// probeNumber 兼容 ffprobe 把数值输出成字符串或裸数字两种形态
// (duration 是字符串,chapter 的 start_time 也可能是数字)。解析不出来就保持
// 零值:这些字段都只是展示用,不该因为一个格式差异让整次探测失败。
type probeNumber float64
func (p *probeNumber) UnmarshalJSON(data []byte) error {
text := strings.TrimSpace(strings.Trim(string(data), `"`))
if text == "" || text == "null" {
return nil
}
value, err := strconv.ParseFloat(text, 64)
if err != nil {
return nil
}
*p = probeNumber(value)
return nil
}
func (p probeNumber) float() float64 { return float64(p) }
func (p probeNumber) int() int { return int(float64(p)) }
func (p probeNumber) int64() int64 { return int64(float64(p)) }
// rawFullProbe 镜像 ffprobe -show_format -show_streams -show_chapters 的输出。
// 只声明用得到的字段;format.filename 刻意不声明,避免它进入任何落库路径。
type rawFullProbe struct {
Format struct {
FormatName string `json:"format_name"`
Duration probeNumber `json:"duration"`
BitRate probeNumber `json:"bit_rate"`
} `json:"format"`
Streams []struct {
Index int `json:"index"`
CodecType string `json:"codec_type"`
CodecName string `json:"codec_name"`
Profile string `json:"profile"`
Width int `json:"width"`
Height int `json:"height"`
PixFmt string `json:"pix_fmt"`
AvgFrameRate string `json:"avg_frame_rate"`
Channels int `json:"channels"`
ChannelLayout string `json:"channel_layout"`
SampleRate probeNumber `json:"sample_rate"`
BitRate probeNumber `json:"bit_rate"`
Tags struct {
Language string `json:"language"`
Title string `json:"title"`
} `json:"tags"`
Disposition struct {
Default int `json:"default"`
Forced int `json:"forced"`
} `json:"disposition"`
} `json:"streams"`
Chapters []struct {
StartTime probeNumber `json:"start_time"`
EndTime probeNumber `json:"end_time"`
Tags struct {
Title string `json:"title"`
} `json:"tags"`
} `json:"chapters"`
}
func parseFullProbeJSON(data []byte) (*FullProbeResult, error) {
var raw rawFullProbe
if err := json.Unmarshal(data, &raw); err != nil {
return nil, fmt.Errorf("parse ffprobe json: %w", err)
}
result := &FullProbeResult{
Container: strings.TrimSpace(raw.Format.FormatName),
DurationSec: raw.Format.Duration.int(),
BitRate: raw.Format.BitRate.int64(),
Streams: make([]ProbeStream, 0, len(raw.Streams)),
Chapters: make([]ProbeChapter, 0, len(raw.Chapters)),
}
for _, s := range raw.Streams {
stream := ProbeStream{
Index: s.Index,
Type: s.CodecType,
Codec: s.CodecName,
Profile: s.Profile,
Language: strings.TrimSpace(s.Tags.Language),
Title: strings.TrimSpace(s.Tags.Title),
Width: s.Width,
Height: s.Height,
PixFmt: s.PixFmt,
FrameRate: normalizeFrameRate(s.AvgFrameRate),
Channels: s.Channels,
ChannelLayout: s.ChannelLayout,
SampleRate: s.SampleRate.int(),
BitRate: s.BitRate.int64(),
Default: s.Disposition.Default != 0,
Forced: s.Disposition.Forced != 0,
}
result.Streams = append(result.Streams, stream)
}
for index, chapter := range raw.Chapters {
result.Chapters = append(result.Chapters, ProbeChapter{
Index: index,
StartMs: secondsToMillis(chapter.StartTime.float()),
EndMs: secondsToMillis(chapter.EndTime.float()),
Title: strings.TrimSpace(chapter.Tags.Title),
})
}
return result, nil
}
// secondsToMillis 把 ffprobe 的秒(浮点)转成毫秒整数。
func secondsToMillis(seconds float64) int64 {
if seconds <= 0 {
return 0
}
return int64(seconds*1000 + 0.5)
}
// normalizeFrameRate 把 "24000/1001" 这类分数帧率换算成可读形式;"0/0"
// (未知)返回空串。
func normalizeFrameRate(raw string) string {
raw = strings.TrimSpace(raw)
if raw == "" {
return ""
}
parts := strings.SplitN(raw, "/", 2)
if len(parts) != 2 {
return raw
}
num, errNum := strconv.ParseFloat(parts[0], 64)
den, errDen := strconv.ParseFloat(parts[1], 64)
if errNum != nil || errDen != nil || den == 0 || num <= 0 {
return ""
}
return strconv.FormatFloat(num/den, 'f', 3, 64)
}
+149
View File
@@ -0,0 +1,149 @@
package service
import (
"strings"
"testing"
)
const fullProbeFixture = `{
"format": {
"format_name": "matroska,webm",
"duration": "1451.024000",
"bit_rate": "8000000",
"filename": "https://cdn.example.com/secret/movie.mkv?d=vip-abc-pickcode&token=xyz"
},
"streams": [
{"index": 0, "codec_type": "video", "codec_name": "hevc", "profile": "Main 10",
"width": 3840, "height": 2160, "pix_fmt": "yuv420p10le", "avg_frame_rate": "24000/1001",
"bit_rate": "7800000", "disposition": {"default": 1}},
{"index": 1, "codec_type": "audio", "codec_name": "eac3", "channels": 6,
"channel_layout": "5.1(side)", "sample_rate": "48000",
"tags": {"language": "eng", "title": "Surround"}, "disposition": {"default": 1}},
{"index": 2, "codec_type": "subtitle", "codec_name": "ass",
"tags": {"language": "chi"}, "disposition": {"forced": 1}}
],
"chapters": [
{"start_time": "0.000000", "end_time": "95.000000", "tags": {"title": "Chapter 01"}},
{"start_time": "228.664000", "end_time": "246.143000", "tags": {"title": "Opening"}}
]
}`
func TestParseFullProbeJSONExtractsStreamsAndChapters(t *testing.T) {
got, err := parseFullProbeJSON([]byte(fullProbeFixture))
if err != nil {
t.Fatalf("parseFullProbeJSON: %v", err)
}
if got.Container != "matroska,webm" || got.DurationSec != 1451 || got.BitRate != 8_000_000 {
t.Fatalf("container/duration/bitrate = %q/%d/%d", got.Container, got.DurationSec, got.BitRate)
}
if len(got.Streams) != 3 {
t.Fatalf("streams = %#v, want 3", got.Streams)
}
video := got.Streams[0]
if video.Type != "video" || video.Codec != "hevc" || video.Width != 3840 || video.Height != 2160 {
t.Fatalf("video stream = %#v", video)
}
if video.FrameRate != "23.976" {
t.Fatalf("frame rate = %q, want 23.976 (converted from 24000/1001)", video.FrameRate)
}
if !video.Default {
t.Fatal("video stream should be flagged default")
}
audio := got.Streams[1]
if audio.Codec != "eac3" || audio.Language != "eng" || audio.Title != "Surround" || audio.Channels != 6 || audio.SampleRate != 48000 {
t.Fatalf("audio stream = %#v", audio)
}
sub := got.Streams[2]
if sub.Type != "subtitle" || sub.Language != "chi" || !sub.Forced {
t.Fatalf("subtitle stream = %#v", sub)
}
if len(got.Chapters) != 2 {
t.Fatalf("chapters = %#v, want 2", got.Chapters)
}
if got.Chapters[1].StartMs != 228_664 || got.Chapters[1].EndMs != 246_143 || got.Chapters[1].Title != "Opening" {
t.Fatalf("chapter[1] = %#v", got.Chapters[1])
}
if videoCount := len(got.StreamsOfType("video")); videoCount != 1 {
t.Fatalf("StreamsOfType(video) = %d, want 1", videoCount)
}
if audioCount := len(got.StreamsOfType("audio")); audioCount != 1 {
t.Fatalf("StreamsOfType(audio) = %d, want 1", audioCount)
}
}
// 落库的 payload 绝不能带 ffprobe 的 format.filename:那是解析后的播放直链,
// 里面是网盘签名和 pickcode,存进数据库等于把可直接下载的链接留下来。
func TestFullProbePayloadNeverLeaksSourceURL(t *testing.T) {
got, err := parseFullProbeJSON([]byte(fullProbeFixture))
if err != nil {
t.Fatalf("parseFullProbeJSON: %v", err)
}
payload, err := got.PayloadJSON()
if err != nil {
t.Fatalf("PayloadJSON: %v", err)
}
for _, needle := range []string{"filename", "cdn.example.com", "pickcode", "token=", "http"} {
if strings.Contains(payload, needle) {
t.Fatalf("payload 泄漏了 %q:\n%s", needle, payload)
}
}
// 技术信息本身必须保留。
for _, needle := range []string{"matroska,webm", "hevc", "3840", "Opening"} {
if !strings.Contains(payload, needle) {
t.Fatalf("payload 缺少 %q:\n%s", needle, payload)
}
}
}
// ffprobe 对 duration / start_time 有时输出字符串、有时输出裸数字,两种都要能读。
func TestParseFullProbeJSONAcceptsNumericAndStringTimes(t *testing.T) {
got, err := parseFullProbeJSON([]byte(`{
"format": {"format_name": "mp4", "duration": 125.5},
"streams": [],
"chapters": [{"start_time": 12.5, "end_time": 20}]
}`))
if err != nil {
t.Fatalf("parseFullProbeJSON: %v", err)
}
if got.DurationSec != 125 {
t.Fatalf("duration = %d, want 125", got.DurationSec)
}
if len(got.Chapters) != 1 || got.Chapters[0].StartMs != 12_500 || got.Chapters[0].EndMs != 20_000 {
t.Fatalf("chapters = %#v", got.Chapters)
}
}
func TestParseFullProbeJSONToleratesMissingSections(t *testing.T) {
got, err := parseFullProbeJSON([]byte(`{}`))
if err != nil {
t.Fatalf("parseFullProbeJSON: %v", err)
}
if got.DurationSec != 0 || len(got.Streams) != 0 || len(got.Chapters) != 0 {
t.Fatalf("result = %#v, want empty", got)
}
payload, err := got.PayloadJSON()
if err != nil {
t.Fatalf("PayloadJSON: %v", err)
}
// 空轨道必须序列化成 [],不能是 null——详情页前端按数组消费。
if !strings.Contains(payload, `"streams":[]`) {
t.Fatalf("payload = %s, want an empty streams array", payload)
}
}
func TestNormalizeFrameRate(t *testing.T) {
cases := map[string]string{
"24000/1001": "23.976",
"25/1": "25.000",
"0/0": "",
"": "",
"25": "25",
}
for input, want := range cases {
if got := normalizeFrameRate(input); got != want {
t.Errorf("normalizeFrameRate(%q) = %q, want %q", input, got, want)
}
}
}
+349
View File
@@ -0,0 +1,349 @@
package service
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"os"
"strings"
"sync"
"time"
"go.uber.org/zap"
"github.com/truewhile/MeBox/internal/helper"
"github.com/truewhile/MeBox/internal/model"
"github.com/truewhile/MeBox/internal/repository"
)
// 播放档案里可选的片头/片尾数据来源。
const (
// SegmentSourceAuto 优先用文件内嵌章节(对这个片源最准),没有可用章节时
// 回落到 TheIntroDB。
SegmentSourceAuto = "auto"
// SegmentSourceFFprobe 只用文件内嵌章节。
SegmentSourceFFprobe = "ffprobe"
// SegmentSourceTheIntroDB 只用社区数据库。它同时是 media_segments.source 的
// 取值——两处必须一致,所以直接复用提供方的常量。
SegmentSourceTheIntroDB = IntroDBSource
)
const (
// mediaProbeTimeout 是一次后台探测的总预算(含把 strm 目标解析成直链)。
mediaProbeTimeout = 90 * time.Second
// mediaProbeFailureRetry 是探测失败后允许重新探测的间隔。失败的结果也会落库,
// 否则每次播放都会为一个坏源重跑一次。
mediaProbeFailureRetry = 6 * time.Hour
// mediaProbeErrorLimit 限制落库的错误信息长度(列宽 512 字节,错误里可能
// 带 URL,截断同时避免超长)。
mediaProbeErrorLimit = 200
)
// NormalizeSegmentSource 把任意输入收敛到合法取值;未知值一律按 auto 处理
// (历史档案里这一列可能还是空串)。
func NormalizeSegmentSource(raw string) string {
switch strings.ToLower(strings.TrimSpace(raw)) {
case SegmentSourceFFprobe:
return SegmentSourceFFprobe
case SegmentSourceTheIntroDB:
return SegmentSourceTheIntroDB
default:
return SegmentSourceAuto
}
}
// mediaProber 是 MediaProbeService 需要的探测能力。抽成接口是为了在测试里注入
// 桩,避免依赖真实 ffprobe 二进制。
type mediaProber interface {
ProbeFull(ctx context.Context, input ProbeInput) (*FullProbeResult, error)
}
// MediaProbeService 用 ffprobe 提取媒体的完整信息(容器、每路轨道、内嵌章节),
// 把结果落库缓存,并把章节映射成可跳过的片头/片尾区间。
//
// 核心约束:一次探测要 3~4 秒(远端直链更慢,跨洋要跑三次 HTTP 事务),所以
// 只允许异步跑。播放链路永远只读缓存,拿不到就下次再来——绝不能让一次探测挡在
// 起播路径上。
type MediaProbeService struct {
log *zap.Logger
repo *repository.Container
probe mediaProber
// resolve 把 strm 播放目标解析成最终直链(含绑定 UA 的请求头)。与转码、
// 内嵌字幕发现走同一条换链路径,否则会踩到网盘 CDN 的防盗链 403。
resolve func(ctx context.Context, raw, userAgent string) (*StrmPlayResult, error)
mu sync.Mutex
inFlight map[string]struct{}
}
// NewMediaProbeService is the constructor.
func NewMediaProbeService(log *zap.Logger, repo *repository.Container, probe *FFprobeService) *MediaProbeService {
svc := &MediaProbeService{
log: log,
repo: repo,
inFlight: make(map[string]struct{}),
}
if probe != nil {
svc.probe = probe
}
return svc
}
// SetPlayTargetResolver injects the strm → direct-link resolver.
func (s *MediaProbeService) SetPlayTargetResolver(resolve func(ctx context.Context, raw, userAgent string) (*StrmPlayResult, error)) *MediaProbeService {
if s != nil {
s.resolve = resolve
}
return s
}
// EnsureAsync 保证这部媒体的探测已排上队,并立刻返回。
//
// 已经有同一条媒体的探测在跑、或服务未配置好时返回 false;这次新排上一条返回
// true——调用方据此告诉客户端「稍后再拉一次」。
func (s *MediaProbeService) EnsureAsync(m *model.Media) bool {
if s == nil || s.probe == nil || s.repo == nil || m == nil || strings.TrimSpace(m.ID) == "" {
return false
}
if !s.reserve(m.ID) {
return false
}
// 复制一份媒体行:调用方的对象可能属于请求作用域,后台协程不该继续引用它。
snapshot := *m
helper.Go(s.log, "service.mediaProbe", func() {
defer s.release(snapshot.ID)
s.run(context.Background(), &snapshot)
})
return true
}
func (s *MediaProbeService) reserve(mediaID string) bool {
s.mu.Lock()
defer s.mu.Unlock()
if s.inFlight == nil {
s.inFlight = make(map[string]struct{})
}
if _, running := s.inFlight[mediaID]; running {
return false
}
s.inFlight[mediaID] = struct{}{}
return true
}
func (s *MediaProbeService) release(mediaID string) {
s.mu.Lock()
delete(s.inFlight, mediaID)
s.mu.Unlock()
}
// run 执行一次探测并落库。它跑在后台,没有调用方能接收错误,所以任何失败都只
// 记录、不外抛。
func (s *MediaProbeService) run(ctx context.Context, m *model.Media) {
ctx, cancel := context.WithTimeout(ctx, mediaProbeTimeout)
defer cancel()
input, err := s.probeInput(ctx, m)
if err != nil {
s.markFailure(ctx, m.ID, err)
return
}
result, err := s.probe.ProbeFull(ctx, input)
if err != nil {
s.markFailure(ctx, m.ID, err)
return
}
if err := s.persistProbe(ctx, m, result); err != nil {
s.markFailure(ctx, m.ID, err)
}
}
// probeInput 把媒体行解析成 ffprobe 能直接打开的输入。
//
// 本地文件给路径;STRM / 云盘先解析成最终直链并带上绑定的请求头——直链与 UA
// 必须配套,用错会被 CDN 拒绝。
func (s *MediaProbeService) probeInput(ctx context.Context, m *model.Media) (ProbeInput, error) {
if m == nil {
return ProbeInput{}, ErrMediaNotFound
}
if !isStrmMediaRow(m) {
if _, err := os.Stat(m.Path); err != nil {
return ProbeInput{}, ErrMediaNotFound
}
return ProbeInput{Source: m.Path}, nil
}
raw := strings.TrimSpace(m.STRMURL)
if raw == "" && strings.HasSuffix(strings.ToLower(strings.TrimSpace(m.Path)), ".strm") {
parsed, err := readLocalSTRMTarget(m.Path)
if err == nil {
raw = strings.TrimSpace(parsed)
}
}
if raw == "" {
return ProbeInput{}, errors.New("strm play target missing")
}
if s.resolve != nil {
resolved, err := s.resolve(ctx, raw, "")
if err != nil {
return ProbeInput{}, err
}
in, err := transcodeInputFromPlayResult(resolved)
if err != nil {
return ProbeInput{}, err
}
return ProbeInput{Source: in.Source, Headers: in.Headers}, nil
}
if isHTTPPlaybackTarget(raw) {
return ProbeInput{Source: raw}, nil
}
return ProbeInput{}, errors.New("strm probe source unavailable")
}
// persistProbe 把一次成功的探测落库:先写片段区间,再写探测行。
//
// 顺序很重要:探测行是「已经探过」的标记,客户端靠它决定要不要继续轮询。先写
// 它会让客户端在区间还没落库时就停止等待。
func (s *MediaProbeService) persistProbe(ctx context.Context, m *model.Media, result *FullProbeResult) error {
if s.repo == nil || s.repo.MediaProbe == nil || s.repo.MediaSegment == nil {
return errors.New("media probe repository not wired")
}
rows := chaptersToSegments(result.Chapters)
for i := range rows {
rows[i].MediaID = m.ID
rows[i].SeriesID = m.SeriesID
}
if err := s.repo.MediaSegment.ReplaceForMedia(ctx, m.ID, SegmentSourceFFprobe, rows); err != nil {
return err
}
payload, err := result.PayloadJSON()
if err != nil {
return err
}
row := &model.MediaProbe{
MediaID: m.ID,
Signature: mediaProbeSignature(m),
Source: mediaProbeInputKind(m),
Container: result.Container,
DurationSec: result.DurationSec,
BitRate: result.BitRate,
Width: firstStreamDimension(result, "video", true),
Height: firstStreamDimension(result, "video", false),
VideoCodec: firstStreamCodec(result, "video"),
AudioCodec: firstStreamCodec(result, "audio"),
VideoStreams: len(result.StreamsOfType("video")),
AudioStreams: len(result.StreamsOfType("audio")),
SubtitleStreams: len(result.StreamsOfType("subtitle")),
ChapterCount: len(result.Chapters),
Payload: payload,
ProbedAt: time.Now(),
}
if err := s.repo.MediaProbe.Upsert(ctx, row); err != nil {
return err
}
s.backfillDuration(ctx, m, result.DurationSec)
return nil
}
// backfillDuration 把探测到的时长补进 media.duration_sec。STRM / 云盘媒体在扫描
// 阶段拿不到时长,而末段区间(end_ms = 0)要靠它才能换算出真实结束时间。
func (s *MediaProbeService) backfillDuration(ctx context.Context, m *model.Media, durationSec int) {
if s.repo == nil || s.repo.DB == nil || m == nil || durationSec <= 0 || m.DurationSec > 0 {
return
}
err := s.repo.DB.WithContext(ctx).Model(&model.Media{}).
Where("id = ? AND duration_sec <= 0", m.ID).
Update("duration_sec", durationSec).Error
if err != nil && s.log != nil {
s.log.Debug("backfill probed duration failed", zap.String("media_id", m.ID), zap.Error(err))
}
}
func (s *MediaProbeService) markFailure(ctx context.Context, mediaID string, probeErr error) {
if s == nil || s.repo == nil || probeErr == nil {
return
}
message := truncateProbeError(probeErr)
if err := s.repo.MediaProbe.MarkFailure(ctx, mediaID, message, time.Now()); err != nil && s.log != nil {
s.log.Debug("record media probe failure failed", zap.String("media_id", mediaID), zap.Error(err))
}
if s.log != nil {
s.log.Debug("media probe failed", zap.String("media_id", mediaID), zap.Error(probeErr))
}
}
// truncateProbeError 限制错误信息长度。错误里可能带被拒绝的直链,落库时截断,
// 避免超长并减少敏感内容。
func truncateProbeError(err error) string {
message := strings.TrimSpace(err.Error())
runes := []rune(message)
if len(runes) > mediaProbeErrorLimit {
return string(runes[:mediaProbeErrorLimit])
}
return message
}
// mediaProbeInputKind 记录输入形态,供详情页判断「这个时长是本地读的还是远端读的」。
func mediaProbeInputKind(m *model.Media) string {
if m == nil {
return ""
}
if isStrmMediaRow(m) {
return "strm"
}
return "local"
}
// mediaProbeSignature 是「探的是哪个文件」的指纹。
//
// 本地文件用路径 + 大小 + 修改时间;STRM / 云盘没有本地文件,用固化的播放目标,
// 且绝不能用解析后的直链(每次签名都不同,缓存会永远失效)。整体做哈希:
// 既固定长度,也不把路径或 pickcode 再抄一份进数据库。
func mediaProbeSignature(m *model.Media) string {
if m == nil {
return ""
}
var base string
if isStrmMediaRow(m) {
target := strings.TrimSpace(m.STRMURL)
if target == "" {
target = strings.TrimSpace(m.Path)
}
base = "strm|" + target
} else if info, err := os.Stat(m.Path); err == nil {
base = fmt.Sprintf("local|%s|%d|%d", m.Path, info.Size(), info.ModTime().Unix())
} else {
base = "local|" + m.Path
}
sum := sha256.Sum256([]byte(base))
return hex.EncodeToString(sum[:16])
}
func firstStreamCodec(result *FullProbeResult, kind string) string {
if result == nil {
return ""
}
for _, stream := range result.Streams {
if stream.Type == kind && stream.Codec != "" {
return stream.Codec
}
}
return ""
}
// firstStreamDimension 取第一路指定类型轨道的宽(width=true)或高。
func firstStreamDimension(result *FullProbeResult, kind string, width bool) int {
if result == nil {
return 0
}
for _, stream := range result.Streams {
if stream.Type != kind {
continue
}
if width {
return stream.Width
}
return stream.Height
}
return 0
}
+72
View File
@@ -0,0 +1,72 @@
package service
import (
"regexp"
"strings"
"github.com/truewhile/MeBox/internal/model"
)
// chapterMinMs 是章节被判为「可跳过区间」的最短长度。太短的章节没有可跳过的
// 内容,多半是章节切分噪声。
const chapterMinMs = 2000
// chapterKindPatterns 是章节标题 → 片段类型的映射表。
//
// 只认「有语义的」标题。大量发布组的章节叫 Chapter 01/02…,那种一律不猜:
// auto 档会优先采用章节数据,猜错会让跳过按钮指到错误的位置,比拿不到数据更糟。
// 匹配不到就返回空串,让 auto 回落到 TheIntroDB。
//
// 西文关键词要求词边界(前后不能是字母),否则 "Options" 会被当成片头(op)、
// "Weekend" 会被当成片尾(end)。中日文关键词按子串匹配(CJK 不是 [a-z],所以
// 同一条边界规则对它们天然成立)。
var chapterKindPatterns = []struct {
kind string
pattern *regexp.Regexp
}{
{model.SegmentKindIntro, regexp.MustCompile(`(?i)(^|[^a-z])(intro|opening|op|ncop)([^a-z]|$)|片头|オープニング`)},
{model.SegmentKindRecap, regexp.MustCompile(`(?i)(^|[^a-z])(recap|previously)([^a-z]|$)|前情|回顾|あらすじ`)},
{model.SegmentKindCredits, regexp.MustCompile(`(?i)(^|[^a-z])(credits|ending|end|outro|nced|ed)([^a-z]|$)|片尾|エンディング`)},
{model.SegmentKindPreview, regexp.MustCompile(`(?i)(^|[^a-z])(preview|next\s+episode|next\s+time)([^a-z]|$)|预告|次回|予告`)},
}
// chapterTitleKind 返回章节标题对应的片段类型;识别不出来时返回空串。
func chapterTitleKind(title string) string {
title = strings.TrimSpace(title)
if title == "" {
return ""
}
for _, entry := range chapterKindPatterns {
if entry.pattern.MatchString(title) {
return entry.kind
}
}
return ""
}
// chaptersToSegments 把内嵌章节映射成可跳过的区间。
//
// 只有标题能识别出类型时才产出区间;EndMs 为 0 的章节落成 0(= 延续到片尾),
// 与提供方契约一致,由客户端按媒体总时长补齐。
func chaptersToSegments(chapters []ProbeChapter) []model.MediaSegment {
rows := make([]model.MediaSegment, 0, 4)
for _, chapter := range chapters {
kind := chapterTitleKind(chapter.Title)
if kind == "" || chapter.StartMs < 0 {
continue
}
if chapter.EndMs > 0 && chapter.EndMs-chapter.StartMs < chapterMinMs {
continue
}
rows = append(rows, model.MediaSegment{
Kind: kind,
StartMs: chapter.StartMs,
EndMs: chapter.EndMs,
Source: SegmentSourceFFprobe,
})
}
return rows
}
@@ -0,0 +1,111 @@
package service
import (
"testing"
"github.com/truewhile/MeBox/internal/model"
)
// 章节标题只有「有语义」时才认。西文关键词要词边界,否则 Options / Weekend
// 这类标题会被误判成片头/片尾;中文日文按子串匹配。
func TestChapterTitleKind(t *testing.T) {
cases := []struct {
title string
want string
}{
{"Intro", model.SegmentKindIntro},
{"Opening", model.SegmentKindIntro},
{"Opening Title", model.SegmentKindIntro},
{"OP", model.SegmentKindIntro},
{"NCOP", model.SegmentKindIntro},
{"片头", model.SegmentKindIntro},
{"オープニング", model.SegmentKindIntro},
{"Recap", model.SegmentKindRecap},
{"Previously on", model.SegmentKindRecap},
{"前情提要", model.SegmentKindRecap},
{"回顾", model.SegmentKindRecap},
{"あらすじ", model.SegmentKindRecap},
{"End Credits", model.SegmentKindCredits},
{"Credits", model.SegmentKindCredits},
{"Ending", model.SegmentKindCredits},
{"ED", model.SegmentKindCredits},
{"Outro", model.SegmentKindCredits},
{"片尾", model.SegmentKindCredits},
{"エンディング", model.SegmentKindCredits},
{"Next Episode", model.SegmentKindPreview},
{"Next time", model.SegmentKindPreview},
{"Preview", model.SegmentKindPreview},
{"预告", model.SegmentKindPreview},
{"次回予告", model.SegmentKindPreview},
// 无语义或容易误判的标题一律不认:猜错比没有数据更糟。
{"", ""},
{" ", ""},
{"Chapter 01", ""},
{"Chapter 12", ""},
{"Options", ""},
{"Weekend", ""},
{"Main Menu", ""},
{"Menu", ""},
}
for _, tc := range cases {
if got := chapterTitleKind(tc.title); got != tc.want {
t.Errorf("chapterTitleKind(%q) = %q, want %q", tc.title, got, tc.want)
}
}
}
func TestChaptersToSegmentsOnlyKeepsSemanticChapters(t *testing.T) {
rows := chaptersToSegments([]ProbeChapter{
{Index: 0, StartMs: 0, EndMs: 95_000, Title: "Chapter 01"},
{Index: 1, StartMs: 228_664, EndMs: 246_143, Title: "Opening"},
{Index: 2, StartMs: 246_500, EndMs: 247_400, Title: "Recap"}, // 不足 chapterMinMs
{Index: 3, StartMs: 3_431_000, EndMs: 0, Title: "End Credits"},
})
if len(rows) != 2 {
t.Fatalf("rows = %#v, want 2 (Opening + End Credits)", rows)
}
if rows[0].Kind != model.SegmentKindIntro || rows[0].StartMs != 228_664 || rows[0].EndMs != 246_143 {
t.Fatalf("first row = %#v", rows[0])
}
if rows[1].Kind != model.SegmentKindCredits || rows[1].StartMs != 3_431_000 || rows[1].EndMs != 0 {
t.Fatalf("second row = %#v", rows[1])
}
for _, row := range rows {
if row.Source != SegmentSourceFFprobe {
t.Fatalf("source = %q, want %q", row.Source, SegmentSourceFFprobe)
}
}
}
func TestChaptersToSegmentsIgnoresUntitledAndNegativeChapters(t *testing.T) {
rows := chaptersToSegments([]ProbeChapter{
{Index: 0, StartMs: 0, EndMs: 60_000, Title: ""},
{Index: 1, StartMs: -1, EndMs: 60_000, Title: "Intro"},
})
if len(rows) != 0 {
t.Fatalf("rows = %#v, want none", rows)
}
}
func TestNormalizeSegmentSourceFallsBackToAuto(t *testing.T) {
cases := map[string]string{
"": SegmentSourceAuto,
" ": SegmentSourceAuto,
"auto": SegmentSourceAuto,
"AUTO": SegmentSourceAuto,
"ffprobe": SegmentSourceFFprobe,
" FFprobe ": SegmentSourceFFprobe,
"theintrodb": SegmentSourceTheIntroDB,
"TheIntroDB": SegmentSourceTheIntroDB,
"aniskip": SegmentSourceAuto, // 未知值按 auto 处理
}
for input, want := range cases {
if got := NormalizeSegmentSource(input); got != want {
t.Errorf("NormalizeSegmentSource(%q) = %q, want %q", input, got, want)
}
}
}
+325
View File
@@ -0,0 +1,325 @@
package service
import (
"context"
"errors"
"os"
"path/filepath"
"strings"
"sync"
"testing"
"time"
"go.uber.org/zap"
"github.com/truewhile/MeBox/internal/model"
"github.com/truewhile/MeBox/internal/repository"
)
// stubProber 是 mediaProber 的测试桩:可以拦住探测(gate)用来验证并发去重,
// 也可以直接返回预置结果或错误。
type stubProber struct {
gate chan struct{}
result *FullProbeResult
err error
mu sync.Mutex
calls int
last ProbeInput
}
func (s *stubProber) ProbeFull(_ context.Context, input ProbeInput) (*FullProbeResult, error) {
s.mu.Lock()
s.calls++
s.last = input
gate := s.gate
result := s.result
err := s.err
s.mu.Unlock()
if gate != nil {
<-gate
}
return result, err
}
func (s *stubProber) callCount() int {
s.mu.Lock()
defer s.mu.Unlock()
return s.calls
}
func (s *stubProber) lastInput() ProbeInput {
s.mu.Lock()
defer s.mu.Unlock()
return s.last
}
func chapterProbeResult() *FullProbeResult {
return &FullProbeResult{
Container: "matroska,webm",
DurationSec: 1451,
BitRate: 8_000_000,
Streams: []ProbeStream{
{Index: 0, Type: "video", Codec: "hevc", Width: 3840, Height: 2160},
{Index: 1, Type: "audio", Codec: "eac3"},
{Index: 2, Type: "subtitle", Codec: "ass"},
},
Chapters: []ProbeChapter{
{Index: 0, StartMs: 0, EndMs: 95_000, Title: "Chapter 01"},
{Index: 1, StartMs: 228_664, EndMs: 246_143, Title: "Opening"},
{Index: 2, StartMs: 3_431_000, EndMs: 0, Title: "End Credits"},
},
}
}
// newProbeFixture 建一条本地媒体(真实落在临时目录里,因为 probeInput 会 stat
// 它)和一个可注入桩的 MediaProbeService。
func newProbeFixture(t *testing.T) (*MediaProbeService, *repository.Container, *model.Media) {
t.Helper()
repos := repository.New(newServiceTestDB(t))
dir := t.TempDir()
path := filepath.Join(dir, "S01E01.mkv")
if err := os.WriteFile(path, []byte("not-really-a-video"), 0o600); err != nil {
t.Fatal(err)
}
m := &model.Media{
Base: model.Base{ID: "ep-1"},
LibraryID: "lib-anime",
SeriesID: "s-1",
Title: "某剧",
Path: path,
SeasonNum: 1,
EpisodeNum: 1,
}
if err := repos.DB.Create(m).Error; err != nil {
t.Fatal(err)
}
svc := NewMediaProbeService(zap.NewNop(), repos, nil)
return svc, repos, m
}
func waitForCondition(t *testing.T, timeout time.Duration, cond func() bool) {
t.Helper()
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
if cond() {
return
}
time.Sleep(5 * time.Millisecond)
}
t.Fatal("condition was not met before the deadline")
}
func TestMediaProbeRunPersistsSegmentsAndProbeRow(t *testing.T) {
svc, repos, m := newProbeFixture(t)
prober := &stubProber{result: chapterProbeResult()}
svc.probe = prober
// 直接调 run(而不是 EnsureAsync)以获得确定性的断言:这里要验证的是落库,
// 不是调度。
svc.run(context.Background(), m)
if got := prober.lastInput().Source; got != m.Path {
t.Fatalf("probe source = %q, want the local path %q", got, m.Path)
}
rows, err := repos.MediaSegment.ListByMediaSource(t.Context(), m.ID, SegmentSourceFFprobe)
if err != nil {
t.Fatal(err)
}
if len(rows) != 2 {
t.Fatalf("segments = %#v, want the 2 semantic chapters", rows)
}
if rows[0].Kind != model.SegmentKindIntro || rows[0].SeriesID != "s-1" {
t.Fatalf("segment[0] = %#v, want an intro carrying the series id", rows[0])
}
probe, err := repos.MediaProbe.Get(t.Context(), m.ID)
if err != nil {
t.Fatal(err)
}
if probe == nil {
t.Fatal("probe row missing")
}
if probe.Container != "matroska,webm" || probe.DurationSec != 1451 || probe.ChapterCount != 3 {
t.Fatalf("probe summary = %#v", probe)
}
if probe.VideoStreams != 1 || probe.AudioStreams != 1 || probe.SubtitleStreams != 1 {
t.Fatalf("stream counts = %d/%d/%d", probe.VideoStreams, probe.AudioStreams, probe.SubtitleStreams)
}
if probe.Width != 3840 || probe.Height != 2160 || probe.VideoCodec != "hevc" || probe.AudioCodec != "eac3" {
t.Fatalf("probe primaries = %#v", probe)
}
if probe.Source != "local" || probe.Signature == "" {
t.Fatalf("probe source/signature = %q/%q", probe.Source, probe.Signature)
}
if probe.LastError != "" {
t.Fatalf("last error = %q, want empty", probe.LastError)
}
if !strings.Contains(probe.Payload, "hevc") {
t.Fatalf("payload should carry the parsed streams: %s", probe.Payload)
}
// 时长回填:STRM / 云盘媒体扫描时拿不到时长,末段区间要靠它换算结束时间。
var refreshed model.Media
if err := repos.DB.First(&refreshed, "id = ?", m.ID).Error; err != nil {
t.Fatal(err)
}
if refreshed.DurationSec != 1451 {
t.Fatalf("media duration = %d, want the probed 1451", refreshed.DurationSec)
}
}
func TestMediaProbeFailureKeepsPreviousProbeAndSegments(t *testing.T) {
svc, repos, m := newProbeFixture(t)
prober := &stubProber{result: chapterProbeResult()}
svc.probe = prober
svc.run(context.Background(), m)
prober.err = errors.New("ffprobe full: exit status 1")
svc.run(context.Background(), m)
probe, err := repos.MediaProbe.Get(t.Context(), m.ID)
if err != nil {
t.Fatal(err)
}
if probe == nil || probe.LastError == "" {
t.Fatal("a failed re-probe must record the error")
}
// 关键:失败只更新时间与错误信息,上一次成功的媒体信息与片段必须留着。
if probe.DurationSec != 1451 || probe.Container != "matroska,webm" {
t.Fatalf("a failed re-probe wiped the previous summary: %#v", probe)
}
rows, err := repos.MediaSegment.ListByMediaSource(t.Context(), m.ID, SegmentSourceFFprobe)
if err != nil {
t.Fatal(err)
}
if len(rows) != 2 {
t.Fatalf("segments = %#v, want the previous 2 kept", rows)
}
}
// 探测有结论后就不再重探;失败要等冷却期过去才允许重试。
func TestMediaProbeSettledSkipsReprobeButStaleFailureRetries(t *testing.T) {
if !mediaProbeSettled(&model.MediaProbe{ProbedAt: time.Now()}) {
t.Fatal("a completed probe is settled")
}
if !mediaProbeSettled(&model.MediaProbe{ProbedAt: time.Now(), LastError: "boom"}) {
t.Fatal("a recent failure is settled: it must not be retried on every play")
}
if mediaProbeSettled(&model.MediaProbe{ProbedAt: time.Now().Add(-mediaProbeFailureRetry - time.Minute), LastError: "boom"}) {
t.Fatal("a failure past the cooldown should be retried")
}
if mediaProbeSettled(nil) {
t.Fatal("a missing probe row is not settled")
}
}
func TestMediaProbeEnsureAsyncDedupesInFlightProbes(t *testing.T) {
svc, _, m := newProbeFixture(t)
gate := make(chan struct{})
prober := &stubProber{result: chapterProbeResult(), gate: gate}
svc.probe = prober
if !svc.EnsureAsync(m) {
t.Fatal("first EnsureAsync should enqueue a probe")
}
// 同一条媒体在跑的时候不能重复排队:否则每次播放请求都会再起一次 3 秒探测。
if svc.EnsureAsync(m) {
t.Fatal("second EnsureAsync must report that a probe is already running")
}
close(gate)
waitForCondition(t, 5*time.Second, func() bool {
svc.mu.Lock()
defer svc.mu.Unlock()
return len(svc.inFlight) == 0
})
if got := prober.callCount(); got != 1 {
t.Fatalf("probe calls = %d, want 1", got)
}
// 跑完之后允许再次排队(例如失败冷却期到了之后的重试)。
if !svc.EnsureAsync(m) {
t.Fatal("EnsureAsync should enqueue again once the previous run finished")
}
waitForCondition(t, 5*time.Second, func() bool { return prober.callCount() == 2 })
}
func TestMediaProbeEnsureAsyncWithoutProberIsNotPending(t *testing.T) {
svc, _, m := newProbeFixture(t)
// probe 未注入(ffprobe 不可用时就是这样):不能谎报「提取中」,否则客户端
// 会白轮询一轮。
if svc.EnsureAsync(m) {
t.Fatal("EnsureAsync must return false without a prober")
}
}
func TestMediaProbeInputRejectsMissingLocalFile(t *testing.T) {
svc, _, _ := newProbeFixture(t)
_, err := svc.probeInput(t.Context(), &model.Media{
Base: model.Base{ID: "missing"}, Path: filepath.Join(t.TempDir(), "nope.mkv"),
})
if !errors.Is(err, ErrMediaNotFound) {
t.Fatalf("err = %v, want ErrMediaNotFound", err)
}
}
func TestMediaProbeInputResolvesStrmWithHeaders(t *testing.T) {
svc := &MediaProbeService{}
svc.SetPlayTargetResolver(func(_ context.Context, raw, userAgent string) (*StrmPlayResult, error) {
if raw != "https://pan.example.com/api/strm/play/115?v=1" {
t.Fatalf("resolver got raw = %q", raw)
}
if userAgent != "" {
t.Fatalf("userAgent = %q, want empty for a background probe", userAgent)
}
return &StrmPlayResult{RedirectURL: "https://cdn.example.com/a.mkv"}, nil
})
input, err := svc.probeInput(t.Context(), &model.Media{
Base: model.Base{ID: "strm-1"},
Path: "/media/a.mkv.strm",
STRMURL: "https://pan.example.com/api/strm/play/115?v=1",
})
if err != nil {
t.Fatalf("probeInput: %v", err)
}
if input.Source != "https://cdn.example.com/a.mkv" {
t.Fatalf("source = %q, want the resolved direct link", input.Source)
}
}
// 签名不能把播放目标(含 pickcode)再抄一份进数据库,所以整体做哈希。
func TestMediaProbeSignatureHidesPlayTargetAndTracksFileChanges(t *testing.T) {
dir := t.TempDir()
path := filepath.Join(dir, "movie.mkv")
if err := os.WriteFile(path, []byte("v1"), 0o600); err != nil {
t.Fatal(err)
}
local := &model.Media{Base: model.Base{ID: "l-1"}, Path: path}
first := mediaProbeSignature(local)
if len(first) != 32 {
t.Fatalf("signature = %q, want a 32-char hex digest", first)
}
if strings.Contains(first, "movie.mkv") {
t.Fatalf("signature leaked the path: %q", first)
}
if err := os.WriteFile(path, []byte("v2-changed-size"), 0o600); err != nil {
t.Fatal(err)
}
if mediaProbeSignature(local) == first {
t.Fatal("the signature must change when the file changes")
}
strm := &model.Media{
Base: model.Base{ID: "s-1"},
Path: "/media/movie.mkv.strm",
STRMURL: "https://pan.example.com/api/strm/play/115?pickcode=secret-pickcode",
}
strmSig := mediaProbeSignature(strm)
if strings.Contains(strmSig, "secret-pickcode") || strings.Contains(strmSig, "pan.example.com") {
t.Fatalf("strm signature leaked the play target: %q", strmSig)
}
if strmSig == "" {
t.Fatal("strm signature must not be empty")
}
}
+138 -9
View File
@@ -44,6 +44,9 @@ type MediaSegmentService struct {
log *zap.Logger
repo *repository.Container
introdb *IntroDBService
// probe 负责从文件内嵌章节里提取片头/片尾(异步、落库)。未注入时只用
// TheIntroDB。
probe *MediaProbeService
}
// NewMediaSegmentService is the constructor.
@@ -59,16 +62,96 @@ func (s *MediaSegmentService) SetIntroDB(p *IntroDBService) *MediaSegmentService
return s
}
// ListForPlayback returns the segments known for a media item, refreshing from
// the provider when the cache is stale.
//
// 它不做任何阻塞起播的事情——调用方是在播放已经开始之后用一次独立请求进来的,
// 抓取失败也只是少一个「跳过片头」按钮,绝不能让播放报错。
func (s *MediaSegmentService) ListForPlayback(ctx context.Context, m *model.Media) ([]model.MediaSegment, error) {
if s == nil || s.repo == nil || m == nil || m.ID == "" {
return nil, nil
// SetProbe wires the in-file chapter extractor. Without it the ffprobe source
// simply yields nothing and auto falls back to TheIntroDB.
func (s *MediaSegmentService) SetProbe(p *MediaProbeService) *MediaSegmentService {
if s != nil {
s.probe = p
}
cached, err := s.repo.MediaSegment.ListByMedia(ctx, m.ID)
return s
}
// SegmentsResult 是播放器一次查询的结果。
type SegmentsResult struct {
// Segments 是按当前来源选定、可以直接用来跳过的区间。
Segments []model.MediaSegment
// Pending 为 true 表示 ffprobe 提取还在后台跑:这次可能还没有章节数据,
// 客户端过几秒再拉一次就能拿到;那时如果片头还没播完,跳过按钮会自动出现。
Pending bool
}
// SegmentsForPlayback 返回播放器该用的片段,并按需触发数据补齐。
//
// 它绝不做阻塞起播的事:TheIntroDB 的抓取沿用原来的「以调用方 deadline 为预算」,
// ffprobe 提取则完全异步。抓取或提取失败只是少一个「跳过片头」按钮、或晚几秒
// 出现,绝不能让播放报错。
func (s *MediaSegmentService) SegmentsForPlayback(ctx context.Context, m *model.Media, source string) (SegmentsResult, error) {
if s == nil || s.repo == nil || m == nil || m.ID == "" {
return SegmentsResult{}, nil
}
source = NormalizeSegmentSource(source)
result := SegmentsResult{}
// 只用社区库时连探测都不该触发:没必要为一次用不上的章节提取去跑 ffprobe。
var chapterRows []model.MediaSegment
needsProbe := false
if source != SegmentSourceTheIntroDB {
chapterRows, needsProbe = s.chapterSegments(ctx, m)
}
result.Pending = needsProbe
// 用 defer 保证无论走哪条分支、包括中途出错,后台提取都会排上队;同时它一定在
// 所有前台数据库读写之后才启动(见 triggerAsyncProbe 的说明)。
defer triggerAsyncProbe(s, m, needsProbe)
if source == SegmentSourceTheIntroDB {
introRows, err := s.introDBSegments(ctx, m)
if err != nil {
return result, err
}
result.Segments = introRows
return result, nil
}
// ffprobe 档、以及 auto 档下已经有可用章节的情况,都整体采用章节数据。
//
// 章节与社区库的数据刻意不合并:两边对同一集的判定会互相矛盾(同一集的片尾
// 起点能差上百秒),只能按 media 整体二选一。auto 走到这里说明章节可用,也就
// 不必再去打一次用不上的社区库。
if source == SegmentSourceFFprobe || len(chapterRows) > 0 {
result.Segments = chapterRows
return result, nil
}
// auto 且没有可用章节:回落到社区库;Pending 保留,客户端会再拉一次。
introRows, err := s.introDBSegments(ctx, m)
if err != nil {
return result, err
}
result.Segments = introRows
return result, nil
}
// triggerAsyncProbe 在所有前台数据库读写都结束之后再起后台提取。
//
// 顺序很关键:后台探测自己也要写库,若在本次请求的写事务还没结束时启动,两个写
// 事务会抢同一把锁(SQLite 下就是 SQLITE_BUSY)。
func triggerAsyncProbe(s *MediaSegmentService, m *model.Media, needsProbe bool) {
if !needsProbe || s == nil || s.probe == nil {
return
}
s.probe.EnsureAsync(m)
}
// ListForPlayback 供 Emby / Jellyfin 兼容接口使用:按 auto 档取数据(章节优先,
// 回落社区库)。第三方客户端不会轮询,所以这里只返回当前能拿到的部分。
func (s *MediaSegmentService) ListForPlayback(ctx context.Context, m *model.Media) ([]model.MediaSegment, error) {
result, err := s.SegmentsForPlayback(ctx, m, SegmentSourceAuto)
if err != nil {
return nil, err
}
return result.Segments, nil
}
// introDBSegments 读社区库的片段,缓存过期时按调用方的预算抓一次并落库。
func (s *MediaSegmentService) introDBSegments(ctx context.Context, m *model.Media) ([]model.MediaSegment, error) {
cached, err := s.repo.MediaSegment.ListByMediaSource(ctx, m.ID, IntroDBSource)
if err != nil {
return nil, err
}
@@ -91,6 +174,52 @@ func (s *MediaSegmentService) ListForPlayback(ctx context.Context, m *model.Medi
return refreshed, nil
}
// chapterSegments 读 ffprobe 提取出的章节区间,并报告「是否还需要等一次提取结果」。
//
// 它只读、不启动提取:调用方要等所有前台数据库读写都结束之后再起后台任务,否则
// 后台写事务会和本次请求的写事务抢同一把锁。探测失败的结果也会落库,所以不会
// 每次播放都为同一个坏源重跑。
func (s *MediaSegmentService) chapterSegments(ctx context.Context, m *model.Media) ([]model.MediaSegment, bool) {
if s == nil || s.probe == nil || s.repo == nil {
return nil, false
}
rows, err := s.repo.MediaSegment.ListByMediaSource(ctx, m.ID, SegmentSourceFFprobe)
if err != nil {
s.debug("list ffprobe segments failed", m.ID, err)
return nil, false
}
cached, err := s.repo.MediaProbe.Get(ctx, m.ID)
if err != nil {
s.debug("get media probe failed", m.ID, err)
return nil, false
}
if mediaProbeSettled(cached) {
return rows, false
}
// 还没探过、或失败已过冷却期:值得让客户端稍后再来一次。
return rows, true
}
// mediaProbeSettled 判断这部媒体的探测是否已经「有结论」——成功过,或者失败但还在
// 冷却期内。有结论就不必再探,客户端也不用继续轮询;失败且已过冷却期时返回
// false,让下一次播放重试。
func mediaProbeSettled(row *model.MediaProbe) bool {
if row == nil {
return false
}
if strings.TrimSpace(row.LastError) == "" {
return true
}
return time.Since(row.ProbedAt) < mediaProbeFailureRetry
}
func (s *MediaSegmentService) debug(message, mediaID string, err error) {
if s == nil || s.log == nil {
return
}
s.log.Debug(message, zap.String("media_id", mediaID), zap.Error(err))
}
// refresh 向提供方查询并落库,返回 (rows, 是否真的发起过查询, error)。
//
// attempted=false 表示这部媒体缺少可查询的外部 ID(最常见的原因是还没刮削,
@@ -0,0 +1,223 @@
package service
import (
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"sync/atomic"
"testing"
"time"
"go.uber.org/zap"
"github.com/truewhile/MeBox/internal/model"
"github.com/truewhile/MeBox/internal/repository"
)
// segmentSourceFixture 提供一个带社区库假服务端、可选探测桩、以及一条已有本地
// 文件的媒体行,用来验证「数据来源选择」这一层。
type segmentSourceFixture struct {
svc *MediaSegmentService
repos *repository.Container
calls *int32
prober *stubProber
media *model.Media
}
func newSegmentSourceFixture(t *testing.T, prober *stubProber) *segmentSourceFixture {
t.Helper()
repos := repository.New(newServiceTestDB(t))
var calls int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
atomic.AddInt32(&calls, 1)
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(introDBMoviePayload))
}))
t.Cleanup(server.Close)
svc := NewMediaSegmentService(zap.NewNop(), repos).
SetIntroDB(NewIntroDBService(zap.NewNop()).SetBaseURL(server.URL).SetRetryDelay(0))
dir := t.TempDir()
path := filepath.Join(dir, "inception.mkv")
if err := os.WriteFile(path, []byte("data"), 0o600); err != nil {
t.Fatal(err)
}
m := &model.Media{Base: model.Base{ID: "mv-1"}, Title: "盗梦空间", Path: path, TMDbID: 27205}
if err := repos.DB.Create(m).Error; err != nil {
t.Fatal(err)
}
if prober != nil {
probeSvc := NewMediaProbeService(zap.NewNop(), repos, nil)
probeSvc.probe = prober
svc.SetProbe(probeSvc)
}
return &segmentSourceFixture{svc: svc, repos: repos, calls: &calls, prober: prober, media: m}
}
// seedChapters 模拟「已经提取过并且拿到了章节」:写入章节区间 + 探测行(探测行
// 是「已经探过」的标记,客户端靠它停止轮询)。
func (f *segmentSourceFixture) seedChapters(t *testing.T, rows ...model.MediaSegment) {
t.Helper()
seeded := make([]model.MediaSegment, 0, len(rows))
for _, row := range rows {
row.MediaID = f.media.ID
row.Source = SegmentSourceFFprobe
seeded = append(seeded, row)
}
if err := f.repos.MediaSegment.ReplaceForMedia(t.Context(), f.media.ID, SegmentSourceFFprobe, seeded); err != nil {
t.Fatal(err)
}
if err := f.repos.MediaProbe.Upsert(t.Context(), &model.MediaProbe{
MediaID: f.media.ID,
ProbedAt: time.Now(),
ChapterCount: len(seeded),
}); err != nil {
t.Fatal(err)
}
}
func chapterIntroRow() model.MediaSegment {
return model.MediaSegment{Kind: model.SegmentKindIntro, StartMs: 228_664, EndMs: 246_143}
}
// auto 档在有章节时必须整体采用章节,并且不打社区库。
func TestSegmentsForPlaybackAutoPrefersChapters(t *testing.T) {
f := newSegmentSourceFixture(t, &stubProber{result: chapterProbeResult()})
f.seedChapters(t, chapterIntroRow())
for _, source := range []string{SegmentSourceAuto, SegmentSourceFFprobe, "未知值"} {
result, err := f.svc.SegmentsForPlayback(t.Context(), f.media, source)
if err != nil {
t.Fatalf("source %q: %v", source, err)
}
if len(result.Segments) != 1 || result.Segments[0].StartMs != 228_664 {
t.Fatalf("source %q segments = %#v, want the chapter row", source, result.Segments)
}
if result.Segments[0].Source != SegmentSourceFFprobe {
t.Fatalf("source %q returned a %q row", source, result.Segments[0].Source)
}
if result.Pending {
t.Fatalf("source %q reported pending although the probe is cached", source)
}
}
if got := atomic.LoadInt32(f.calls); got != 0 {
t.Fatalf("provider calls = %d, want 0: chapters must not trigger TheIntroDB", got)
}
}
// theintrodb 档必须忽略章节,并且不触发任何探测。
func TestSegmentsForPlaybackIntroDBOnlyIgnoresChapters(t *testing.T) {
prober := &stubProber{result: chapterProbeResult()}
f := newSegmentSourceFixture(t, prober)
f.seedChapters(t, chapterIntroRow())
result, err := f.svc.SegmentsForPlayback(t.Context(), f.media, SegmentSourceTheIntroDB)
if err != nil {
t.Fatal(err)
}
if len(result.Segments) != 1 || result.Segments[0].Source != IntroDBSource {
t.Fatalf("segments = %#v, want the community row", result.Segments)
}
if result.Segments[0].StartMs != 0 || result.Segments[0].EndMs != 38_000 {
t.Fatalf("segments = %#v, want the movie intro from the provider", result.Segments)
}
if result.Pending {
t.Fatal("theintrodb-only must never report pending")
}
// 只用社区库时不该为一次用不上的章节提取去跑 3 秒 ffprobe。
if got := prober.callCount(); got != 0 {
t.Fatalf("probe calls = %d, want 0", got)
}
}
// auto 档没有章节时回落到社区库,同时告诉客户端「章节还在提取」。
func TestSegmentsForPlaybackAutoFallsBackAndReportsPending(t *testing.T) {
prober := &stubProber{result: chapterProbeResult()}
f := newSegmentSourceFixture(t, prober)
result, err := f.svc.SegmentsForPlayback(t.Context(), f.media, SegmentSourceAuto)
if err != nil {
t.Fatal(err)
}
if len(result.Segments) != 1 || result.Segments[0].Source != IntroDBSource {
t.Fatalf("segments = %#v, want the community fallback", result.Segments)
}
if !result.Pending {
t.Fatal("auto without chapters must report that a probe is running")
}
// 异步提取确实排上了队(这里只验证调度;落库由 media_probe_test.go 覆盖)。
waitForCondition(t, 5*time.Second, func() bool { return prober.callCount() >= 1 })
}
// ffprobe 档没有章节数据时返回空,但会起一次提取并报告 pending。
func TestSegmentsForPlaybackFFprobeOnlyReturnsEmptyWhilePending(t *testing.T) {
prober := &stubProber{result: chapterProbeResult()}
f := newSegmentSourceFixture(t, prober)
result, err := f.svc.SegmentsForPlayback(t.Context(), f.media, SegmentSourceFFprobe)
if err != nil {
t.Fatal(err)
}
if len(result.Segments) != 0 {
t.Fatalf("segments = %#v, want none before extraction finishes", result.Segments)
}
if !result.Pending {
t.Fatal("ffprobe-only must report pending while the extraction runs")
}
if got := atomic.LoadInt32(f.calls); got != 0 {
t.Fatalf("provider calls = %d, want 0", got)
}
}
// 提取失败也要有个结果:不能在每次播放时无限重试,客户端也不该一直轮询。
func TestSegmentsForPlaybackStopsPendingAfterFailedProbe(t *testing.T) {
prober := &stubProber{result: chapterProbeResult()}
f := newSegmentSourceFixture(t, prober)
if err := f.repos.MediaProbe.MarkFailure(t.Context(), f.media.ID, "probe exploded", time.Now()); err != nil {
t.Fatal(err)
}
result, err := f.svc.SegmentsForPlayback(t.Context(), f.media, SegmentSourceFFprobe)
if err != nil {
t.Fatal(err)
}
if result.Pending {
t.Fatal("a recent failure must stop the client from polling")
}
if got := prober.callCount(); got != 0 {
t.Fatalf("probe calls = %d, want 0 inside the failure cooldown", got)
}
}
// 没有注入探测服务时(ffprobe 不可用的精简部署)不能谎报 pending。
func TestSegmentsForPlaybackWithoutProbeNeverReportsPending(t *testing.T) {
f := newSegmentSourceFixture(t, nil)
result, err := f.svc.SegmentsForPlayback(t.Context(), f.media, SegmentSourceFFprobe)
if err != nil {
t.Fatal(err)
}
if result.Pending {
t.Fatal("pending must be false when no prober is wired")
}
if len(result.Segments) != 0 {
t.Fatalf("segments = %#v, want none", result.Segments)
}
}
// ListForPlayback 仍供 Emby / Jellyfin 兼容接口使用,按 auto 档取数。
func TestListForPlaybackUsesAutoSelection(t *testing.T) {
f := newSegmentSourceFixture(t, &stubProber{result: chapterProbeResult()})
f.seedChapters(t, chapterIntroRow())
rows, err := f.svc.ListForPlayback(t.Context(), f.media)
if err != nil {
t.Fatal(err)
}
if len(rows) != 1 || rows[0].Source != SegmentSourceFFprobe {
t.Fatalf("rows = %#v, want the chapter row via auto", rows)
}
}
+2
View File
@@ -94,6 +94,7 @@ func (s *PlayProfileService) Create(ctx context.Context, in PlayProfileInput) (*
PreferredAudioLang: in.PreferredAudioLang,
AutoplayNext: in.AutoplayNext,
SkipIntro: in.SkipIntro,
SegmentSource: NormalizeSegmentSource(in.SegmentSource),
AllowedLibraryIDs: string(libsBlob),
}
if in.RequirePIN && in.PIN != "" {
@@ -153,6 +154,7 @@ func (s *PlayProfileService) updateExisting(ctx context.Context, row *model.Play
"preferred_audio_lang": in.PreferredAudioLang,
"autoplay_next": in.AutoplayNext,
"skip_intro": in.SkipIntro,
"segment_source": NormalizeSegmentSource(in.SegmentSource),
"allowed_library_ids": string(libsBlob),
}
if in.RequirePIN && in.PIN != "" {
+14 -12
View File
@@ -13,18 +13,20 @@ import (
// PlayProfileInput is the create/update payload accepted by the API.
// PIN is hashed only when non-empty so omitting it preserves the existing PIN on update.
type PlayProfileInput struct {
UserID string `json:"user_id"`
Name string `json:"name"`
IsDefault bool `json:"is_default"`
ContentRatingLimit string `json:"content_rating_limit"`
AllowAdult bool `json:"allow_adult"`
RequirePIN bool `json:"require_pin"`
PIN string `json:"pin,omitempty"`
PreferredSubtitleLang string `json:"preferred_subtitle_lang"`
PreferredAudioLang string `json:"preferred_audio_lang"`
AutoplayNext bool `json:"autoplay_next"`
SkipIntro bool `json:"skip_intro"`
AllowedLibraryIDs []string `json:"allowed_library_ids"`
UserID string `json:"user_id"`
Name string `json:"name"`
IsDefault bool `json:"is_default"`
ContentRatingLimit string `json:"content_rating_limit"`
AllowAdult bool `json:"allow_adult"`
RequirePIN bool `json:"require_pin"`
PIN string `json:"pin,omitempty"`
PreferredSubtitleLang string `json:"preferred_subtitle_lang"`
PreferredAudioLang string `json:"preferred_audio_lang"`
AutoplayNext bool `json:"autoplay_next"`
SkipIntro bool `json:"skip_intro"`
// SegmentSource 是片头/片尾数据来源:auto | theintrodb | ffprobe。
SegmentSource string `json:"segment_source"`
AllowedLibraryIDs []string `json:"allowed_library_ids"`
}
// ProfileView is the public shape for React forms.
+1
View File
@@ -36,6 +36,7 @@ type Container struct {
Scraper *ScraperService
Playback *PlaybackService
Segments *MediaSegmentService
MediaProbe *MediaProbeService
ImageProxy *ImageProxy
Watcher *WatcherService
Subtitle *SubtitleService
+5
View File
@@ -147,6 +147,11 @@ func (b *serviceContainerBuilder) initContentServices() {
b.c.Transcoder.SetStrmPlayTargetResolver(b.c.Strm.ResolvePlayTarget)
b.c.Transcoder.SetProbe(b.c.FFprobe)
b.c.Subtitle.SetStrmPlayTargetResolver(b.c.Strm.ResolvePlayTarget)
// 播放时的媒体信息提取(ffprobe 章节 → 跳过片头/片尾):完全异步,播放链路
// 只读缓存。换链复用与转码/字幕同一条路径,避免 CDN 防盗链 403。
b.c.MediaProbe = NewMediaProbeService(b.log, b.repos, b.c.FFprobe).
SetPlayTargetResolver(b.c.Strm.ResolvePlayTargetWithUA)
b.c.Segments.SetProbe(b.c.MediaProbe)
// 播放链路:/Videos/{id}/stream 与 /api/stream/{id} 在服务端完成换链后直接
// 302 到最终直链,客户端少跟随一次 302(高延迟线路上省一个往返)。
b.c.Stream.SetStrmPlayTargetResolver(b.c.Strm.ResolvePlayTargetWithUA)
+2 -1
View File
@@ -1,5 +1,5 @@
import { api } from './client'
import type { PlayProfile } from '../types'
import type { PlayProfile, SegmentSource } from '../types'
// Payload accepted by create / update.
export interface PlayProfileInput {
@@ -14,6 +14,7 @@ export interface PlayProfileInput {
preferred_audio_lang?: string
autoplay_next: boolean
skip_intro: boolean
segment_source: SegmentSource
allowed_library_ids: string[]
}
+28 -8
View File
@@ -75,6 +75,13 @@ type PlaybackProgressSession = {
// 自动跳过片头后,「已跳过 · 撤销」提示停留的时长。
const SKIP_NOTICE_MS = 6000
// ffprobe 章节提取在服务端是异步的:pending 为 true 时按这个间隔重试,最多这么
// 多次(远端探测实测 3~4 秒,本地更快;40 秒的预算足以覆盖慢速 CDN 与大文件)。
// 重试期间播放完全不受影响;拿到数据时如果片头已经播过去了,resolveActiveSkip
// 自然不会再提示,不会出现「点一下就跳过头」的按钮。
const SKIP_SEGMENTS_POLL_MS = 4000
const SKIP_SEGMENTS_MAX_POLLS = 10
function normalizePlayerVolume(value: unknown): number {
const parsed = Number(value)
if (!Number.isFinite(parsed)) return 1
@@ -732,6 +739,7 @@ export function PlayerPage() {
useEffect(() => {
if (!mediaId) return
let cancelled = false
let pollTimer: ReturnType<typeof setTimeout> | null = null
setRawSkipSegments([])
setAutoSkipIntro(false)
setActiveSkip(null)
@@ -739,16 +747,28 @@ export function PlayerPage() {
setDismissedSkipKinds([])
setAutoSuppressedKinds([])
// 片段数据与播放来源无关,播放开始后异步补抓即可,绝不挡在起播路径上。
playbackAPI
.segments(mediaId)
.then((res) => {
if (cancelled) return
setAutoSkipIntro(Boolean(res.auto_skip))
setRawSkipSegments(res.segments ?? [])
})
.catch(() => undefined)
//
// 服务端的 ffprobe 章节提取是异步的:pending 为 true 说明这次还没结果,隔几秒
// 再拉一次。提取完成后如果片头还没播完,跳过按钮会自己出现;已经过了片头时间
// 的话 resolveActiveSkip 不会提示,所以不会出现「点一下就跳过头」的按钮。
const load = (attempt: number) => {
playbackAPI
.segments(mediaId)
.then((res) => {
if (cancelled) return
setAutoSkipIntro(Boolean(res.auto_skip))
setRawSkipSegments(res.segments ?? [])
const hasSegments = (res.segments ?? []).length > 0
if (res.pending && !hasSegments && attempt < SKIP_SEGMENTS_MAX_POLLS) {
pollTimer = setTimeout(() => load(attempt + 1), SKIP_SEGMENTS_POLL_MS)
}
})
.catch(() => undefined)
}
load(0)
return () => {
cancelled = true
if (pollTimer) clearTimeout(pollTimer)
}
}, [mediaId])
+1
View File
@@ -37,6 +37,7 @@ export function ProfileFormModal({
preferred_audio_lang: editing?.preferred_audio_lang ?? '',
autoplay_next: editing?.autoplay_next ?? true,
skip_intro: editing?.skip_intro ?? false,
segment_source: editing?.segment_source ?? 'auto',
allowed_library_ids: editing?.allowed_library_ids ?? [],
}))
const [saving, setSaving] = useState(false)
+11
View File
@@ -126,6 +126,17 @@ export function ProfilePreferenceFields({
checked={form.skip_intro}
onChange={(value) => update({ skip_intro: value })}
/>
<Field label="片头片尾数据来源">
<select
className="input-base"
value={form.segment_source}
onChange={(event) => update({ segment_source: event.target.value as PlayProfileInput['segment_source'] })}
>
<option value="auto">自动(优先用文件内嵌章节,没有可用章节时用 TheIntroDB)</option>
<option value="theintrodb">TheIntroDB 社区数据库(依赖刮削出的 TMDb ID)</option>
<option value="ffprobe">ffprobe 本地提取(读取文件内嵌章节,不联网)</option>
</select>
</Field>
</>
)
}
+4
View File
@@ -10,9 +10,13 @@ export interface PlayProfile {
preferred_audio_lang?: string
autoplay_next: boolean
skip_intro: boolean
segment_source?: SegmentSource
allowed_library_ids: string[]
total_watch_time: number
last_active_at?: string
created_at: string
updated_at: string
}
/** 片头/片尾数据来源:auto 优先文件内嵌章节,没有则回落 TheIntroDB。 */
export type SegmentSource = 'auto' | 'theintrodb' | 'ffprobe'
+11
View File
@@ -43,5 +43,16 @@ export interface PlaybackSegmentsResponse {
segments: PlaybackSegment[]
/** 当前生效播放档案的「自动跳过片头」开关。 */
auto_skip: boolean
/**
* true 表示 ffprobe 章节提取还在服务端后台跑:这次可能还没有章节数据,隔几秒
* 再拉一次就能拿到。拿到时若片头还没播完,跳过按钮会自动出现;已经过了片头
* 时间则不会提示。
*/
pending?: boolean
/** 本次实际生效的数据来源。 */
source?: PlaybackSegmentSource
}
/** 片头/片尾数据来源,与后端 play_profiles.segment_source 一致。 */
export type PlaybackSegmentSource = 'auto' | 'theintrodb' | 'ffprobe'