diff --git a/internal/handler/playback_segments.go b/internal/handler/playback_segments.go index 35d5abe..ef06393 100644 --- a/internal/handler/playback_segments.go +++ b/internal/handler/playback_segments.go @@ -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) +} diff --git a/internal/model/media_probe.go b/internal/model/media_probe.go new file mode 100644 index 0000000..8336081 --- /dev/null +++ b/internal/model/media_probe.go @@ -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"` +} diff --git a/internal/model/model.go b/internal/model/model.go index fa5f7dd..fc735fc 100644 --- a/internal/model/model.go +++ b/internal/model/model.go @@ -40,6 +40,7 @@ func AllModels() []interface{} { &PlaybackHistory{}, &MediaSegment{}, &MediaSegmentFetch{}, + &MediaProbe{}, &Favorite{}, &Playlist{}, &PlaylistItem{}, diff --git a/internal/model/playback_collection.go b/internal/model/playback_collection.go index 5e27722..5c7bbf5 100644 --- a/internal/model/playback_collection.go +++ b/internal/model/playback_collection.go @@ -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"` } diff --git a/internal/repository/media_probe_repository.go b/internal/repository/media_probe_repository.go new file mode 100644 index 0000000..42cca1e --- /dev/null +++ b/internal/repository/media_probe_repository.go @@ -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 +} diff --git a/internal/repository/media_probe_repository_test.go b/internal/repository/media_probe_repository_test.go new file mode 100644 index 0000000..08a21fa --- /dev/null +++ b/internal/repository/media_probe_repository_test.go @@ -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) + } +} diff --git a/internal/repository/media_segment_repository.go b/internal/repository/media_segment_repository.go index df66e11..bd0a900 100644 --- a/internal/repository/media_segment_repository.go +++ b/internal/repository/media_segment_repository.go @@ -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. // diff --git a/internal/repository/repository.go b/internal/repository/repository.go index 343f44c..5dd90a2 100644 --- a/internal/repository/repository.go +++ b/internal/repository/repository.go @@ -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}, diff --git a/internal/service/ffprobe_full.go b/internal/service/ffprobe_full.go new file mode 100644 index 0000000..ef3e411 --- /dev/null +++ b/internal/service/ffprobe_full.go @@ -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) +} diff --git a/internal/service/ffprobe_full_test.go b/internal/service/ffprobe_full_test.go new file mode 100644 index 0000000..1f95a76 --- /dev/null +++ b/internal/service/ffprobe_full_test.go @@ -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) + } + } +} diff --git a/internal/service/media_probe.go b/internal/service/media_probe.go new file mode 100644 index 0000000..b1fba0b --- /dev/null +++ b/internal/service/media_probe.go @@ -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 +} diff --git a/internal/service/media_probe_chapters.go b/internal/service/media_probe_chapters.go new file mode 100644 index 0000000..b6fc7e5 --- /dev/null +++ b/internal/service/media_probe_chapters.go @@ -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 +} diff --git a/internal/service/media_probe_chapters_test.go b/internal/service/media_probe_chapters_test.go new file mode 100644 index 0000000..bd9099d --- /dev/null +++ b/internal/service/media_probe_chapters_test.go @@ -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) + } + } +} diff --git a/internal/service/media_probe_test.go b/internal/service/media_probe_test.go new file mode 100644 index 0000000..f8e4327 --- /dev/null +++ b/internal/service/media_probe_test.go @@ -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") + } +} diff --git a/internal/service/media_segment.go b/internal/service/media_segment.go index f4fecec..e3e3536 100644 --- a/internal/service/media_segment.go +++ b/internal/service/media_segment.go @@ -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(最常见的原因是还没刮削, diff --git a/internal/service/media_segment_source_test.go b/internal/service/media_segment_source_test.go new file mode 100644 index 0000000..63fd040 --- /dev/null +++ b/internal/service/media_segment_source_test.go @@ -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) + } +} diff --git a/internal/service/play_profile.go b/internal/service/play_profile.go index d182809..4c3d5a3 100644 --- a/internal/service/play_profile.go +++ b/internal/service/play_profile.go @@ -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 != "" { diff --git a/internal/service/play_profile_model.go b/internal/service/play_profile_model.go index 77f8908..7e910f9 100644 --- a/internal/service/play_profile_model.go +++ b/internal/service/play_profile_model.go @@ -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. diff --git a/internal/service/service.go b/internal/service/service.go index 905d0de..95b1c5b 100644 --- a/internal/service/service.go +++ b/internal/service/service.go @@ -36,6 +36,7 @@ type Container struct { Scraper *ScraperService Playback *PlaybackService Segments *MediaSegmentService + MediaProbe *MediaProbeService ImageProxy *ImageProxy Watcher *WatcherService Subtitle *SubtitleService diff --git a/internal/service/service_builder.go b/internal/service/service_builder.go index 8ffcfd4..9ab8860 100644 --- a/internal/service/service_builder.go +++ b/internal/service/service_builder.go @@ -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) diff --git a/web/src/api/play_profiles.ts b/web/src/api/play_profiles.ts index 158a73c..8899635 100644 --- a/web/src/api/play_profiles.ts +++ b/web/src/api/play_profiles.ts @@ -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[] } diff --git a/web/src/pages/PlayerPage.tsx b/web/src/pages/PlayerPage.tsx index 4890ad2..5bc33f6 100644 --- a/web/src/pages/PlayerPage.tsx +++ b/web/src/pages/PlayerPage.tsx @@ -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 | 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]) diff --git a/web/src/pages/ProfileFormModal.tsx b/web/src/pages/ProfileFormModal.tsx index 3eca77c..fb54aec 100644 --- a/web/src/pages/ProfileFormModal.tsx +++ b/web/src/pages/ProfileFormModal.tsx @@ -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) diff --git a/web/src/pages/ProfileFormSections.tsx b/web/src/pages/ProfileFormSections.tsx index 41de151..d55bcaa 100644 --- a/web/src/pages/ProfileFormSections.tsx +++ b/web/src/pages/ProfileFormSections.tsx @@ -126,6 +126,17 @@ export function ProfilePreferenceFields({ checked={form.skip_intro} onChange={(value) => update({ skip_intro: value })} /> + + + ) } diff --git a/web/src/types/playProfiles.ts b/web/src/types/playProfiles.ts index 5e05a60..028ce87 100644 --- a/web/src/types/playProfiles.ts +++ b/web/src/types/playProfiles.ts @@ -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' diff --git a/web/src/types/playback.ts b/web/src/types/playback.ts index 83558e4..96a15fb 100644 --- a/web/src/types/playback.ts +++ b/web/src/types/playback.ts @@ -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' +