From e273b5e87d54e042464367ad11d7a022cb4d4bdb Mon Sep 17 00:00:00 2001 From: truewhile <779943132@qq.com> Date: Thu, 24 Sep 2026 11:43:17 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../plans/2026-09-24-segment-coverage.md | 36 +++ .../2026-09-24-segment-coverage-design.md | 52 +++++ internal/repository/media_repository.go | 25 +++ .../repository/media_segment_repository.go | 43 ++++ .../media_segment_repository_test.go | 53 +++++ internal/service/media_segment.go | 212 ++++++++++++++++-- internal/service/media_segment_test.go | 118 ++++++++++ internal/service/scheduler.go | 15 +- internal/service/scheduler_segment_jobs.go | 52 +++++ .../service/scheduler_segment_jobs_test.go | 60 +++++ internal/service/service_builder.go | 1 + web/src/pages/settingsGroupGeneral.ts | 7 + 12 files changed, 653 insertions(+), 21 deletions(-) create mode 100644 docs/superpowers/plans/2026-09-24-segment-coverage.md create mode 100644 docs/superpowers/specs/2026-09-24-segment-coverage-design.md create mode 100644 internal/service/scheduler_segment_jobs.go create mode 100644 internal/service/scheduler_segment_jobs_test.go diff --git a/docs/superpowers/plans/2026-09-24-segment-coverage.md b/docs/superpowers/plans/2026-09-24-segment-coverage.md new file mode 100644 index 0000000..d036761 --- /dev/null +++ b/docs/superpowers/plans/2026-09-24-segment-coverage.md @@ -0,0 +1,36 @@ +# Segment Coverage Implementation Plan + +> **For Claude:** REQUIRED SUB-SKILL: Use superpowers:executing-plans to implement this plan task-by-task. + +**Goal:** IntroDB prewarm + same-season intro propagation + shorter miss TTL + playback-only duration probe (option C). + +**Architecture:** Extend `MediaSegmentService` for merge/propagate/prewarm; add scheduler job `segment_prewarm`; repository queries for candidates and season siblings. + +**Tech Stack:** Go, GORM/SQLite, existing SchedulerService + IntroDB client. + +--- + +### Task 1: Negative cache TTL + ledgerFresh tests + +**Files:** `internal/service/media_segment.go`, `internal/service/media_segment_test.go` + +- Change `segmentMissingTTL` to 24h; update tests for fresh/stale miss. + +### Task 2: Playback merge + propagation + +**Files:** `media_segment.go`, `media_segment_repository.go`, `media_repository.go`, tests + +- Constant `PropagatedSource = "propagated"`. +- `ListForPlayback` merges sources with priority manual > theintrodb > propagated. +- After successful IntroDB refresh with intro, propagate to same-season siblings (shared `tm_db_id` or `series_id`). + +### Task 3: Prewarm API + scheduler job + +**Files:** `media_segment.go`, `media_segment_repository.go`, `scheduler.go`, `scheduler_segment_jobs.go`, `service_builder.go`, tests + +- `Prewarm(ctx, limit)` rate-limited refreshes. +- Job `segment_prewarm` every 6h. + +### Task 4: Verify + +- `go test` for affected packages. diff --git a/docs/superpowers/specs/2026-09-24-segment-coverage-design.md b/docs/superpowers/specs/2026-09-24-segment-coverage-design.md new file mode 100644 index 0000000..4c9b6eb --- /dev/null +++ b/docs/superpowers/specs/2026-09-24-segment-coverage-design.md @@ -0,0 +1,52 @@ +# Segment Coverage Improvements (1–4, option C) + +## Goal + +Raise intro/outro skip availability without Chromaprint or full-library STRM probing. + +## Scope (approved) + +1. **IntroDB prewarm** — background job for all queryable media (TMDb + season/episode). +2. **Duration** — playback-path only (no full STRM scan); keep `ListForPlayback` → `EnsureAsync`. +3. **Same-season intro propagation** — copy intro to sibling episodes after a hit. +4. **Negative cache** — miss TTL 24h; 403/429/timeout must not write miss ledger. + +## Non-goals + +Chromaprint, manual segment UI, library-name filters, batch STRM duration scan. + +## Behavior + +### Prewarm (`segment_prewarm`) + +- Scheduler job, interval 6h, initial delay 6h (avoid restart spike). +- **Opt-in via `segment.prewarm_enabled` (default off)** — system settings toggle「后台预热 IntroDB 片头片段」. +- Candidates: `tm_db_id > 0`, season/episode or movie, ledger missing or stale. +- Order: recently played first, then others. +- Rate: ~1 req/s, stop on context cancel; respect IntroDB 429 retry already in client. +- Reuse `MediaSegmentService` refresh + propagation. +- Manual `RunNow` can bypass the toggle (same pattern as organize). + +### Propagation + +- After TheIntroDB refresh finds at least one `intro` span, write the same intro window to same-season siblings that lack a `theintrodb` intro. +- Source tag: `propagated`. +- Do not propagate credits/preview/recap. +- Playback merge priority per kind: `manual` > `theintrodb` > `propagated`. + +### Negative cache + +- `segmentMissingTTL`: 24h (was 7d). +- `segmentFoundTTL`: 30d (unchanged). +- Provider errors (non-404) continue to skip ledger writes. + +### Duration (C) + +- No new batch probe job. +- Existing `defer ensureMediaProbe` on playback remains the only STRM duration path. + +## Success criteria + +- Prewarm increases `media_segment_fetches` over time without blocking play. +- One IntroDB intro hit can surface skip on sibling episodes via `propagated`. +- Misses are retried within ~24h (sooner if recently played and prewarm runs). diff --git a/internal/repository/media_repository.go b/internal/repository/media_repository.go index 23c49e0..673bf86 100644 --- a/internal/repository/media_repository.go +++ b/internal/repository/media_repository.go @@ -232,6 +232,31 @@ func (r *MediaRepository) ExistsSiblingWithTMDbID(ctx context.Context, m *model. return count > 0 } +// ListSeasonSiblings returns other episodes in the same season as m. +// Prefers series_id; falls back to shared library+title+tm_db_id (anime scrape path). +func (r *MediaRepository) ListSeasonSiblings(ctx context.Context, m *model.Media) ([]model.Media, error) { + if r == nil || m == nil || m.SeasonNum <= 0 || m.ID == "" { + return nil, nil + } + query := r.db.WithContext(ctx).Model(&model.Media{}). + Where("season_num = ? AND id <> ?", m.SeasonNum, m.ID) + if seriesID := strings.TrimSpace(m.SeriesID); seriesID != "" { + query = query.Where("series_id = ?", seriesID) + } else { + libraryID := strings.TrimSpace(m.LibraryID) + title := strings.TrimSpace(m.Title) + if libraryID == "" || title == "" || m.TMDbID <= 0 { + return nil, nil + } + query = query.Where("library_id = ? AND title = ? AND tm_db_id = ?", libraryID, title, m.TMDbID) + } + var rows []model.Media + if err := query.Find(&rows).Error; err != nil { + return nil, err + } + return rows, nil +} + // ListByLibrary returns paginated media items for a library. func (r *MediaRepository) ListByLibrary(ctx context.Context, libraryID string, offset, limit int) ([]model.Media, int64, error) { return r.ListByLibraryFiltered(ctx, libraryID, offset, limit, MediaQueryFilter{IncludeNSFW: true}) diff --git a/internal/repository/media_segment_repository.go b/internal/repository/media_segment_repository.go index bd0a900..a63a6d0 100644 --- a/internal/repository/media_segment_repository.go +++ b/internal/repository/media_segment_repository.go @@ -3,6 +3,7 @@ package repository import ( "context" "errors" + "time" "gorm.io/gorm" "gorm.io/gorm/clause" @@ -84,3 +85,45 @@ func (r *MediaSegmentRepository) UpsertFetch(ctx context.Context, row *model.Med } return r.db.WithContext(ctx).Clauses(onConflict).Create(row).Error } + +// ListPrewarmCandidates returns media that should be refreshed from source: +// no ledger, expired miss (fetched_at < missBefore), or expired hit (fetched_at < hitBefore). +// Recently played items come first so hot titles recover coverage sooner. +func (r *MediaSegmentRepository) ListPrewarmCandidates( + ctx context.Context, + source string, + missBefore, hitBefore time.Time, + limit int, +) ([]model.Media, error) { + if r == nil || limit <= 0 { + return nil, nil + } + rows := make([]model.Media, 0, limit) + // 可查询:有 TMDb,且是剧集(有季集)或电影(无季集)。 + err := r.db.WithContext(ctx).Raw(` +SELECT m.* +FROM media m +LEFT JOIN media_segment_fetches f + ON f.media_id = m.id AND f.source = ? AND f.deleted_at IS NULL +LEFT JOIN ( + SELECT media_id, MAX(updated_at) AS last_played + FROM playback_histories + WHERE deleted_at IS NULL + GROUP BY media_id +) ph ON ph.media_id = m.id +WHERE m.deleted_at IS NULL + AND m.tm_db_id > 0 + AND ( + (m.season_num > 0 AND m.episode_num > 0) + OR (COALESCE(m.season_num, 0) = 0 AND COALESCE(m.episode_num, 0) = 0) + ) + AND ( + f.id IS NULL + OR (f.found = 0 AND f.fetched_at < ?) + OR (f.found = 1 AND f.fetched_at < ?) + ) +ORDER BY ph.last_played DESC +LIMIT ? +`, source, missBefore, hitBefore, limit).Scan(&rows).Error + return rows, err +} diff --git a/internal/repository/media_segment_repository_test.go b/internal/repository/media_segment_repository_test.go index c4ac6b7..b9db972 100644 --- a/internal/repository/media_segment_repository_test.go +++ b/internal/repository/media_segment_repository_test.go @@ -152,3 +152,56 @@ func TestUpsertFetchKeepsOneRowPerMediaAndSource(t *testing.T) { t.Fatalf("ledger should be updated in place, got %#v", got) } } + +func TestListSeasonSiblingsSharesLibraryTitleTMDb(t *testing.T) { + repos := newSegmentTestRepos(t) + ctx := t.Context() + eps := []*model.Media{ + {Base: model.Base{ID: "a"}, LibraryID: "lib", Title: "Show", Path: "/a", SeasonNum: 1, EpisodeNum: 1, TMDbID: 99}, + {Base: model.Base{ID: "b"}, LibraryID: "lib", Title: "Show", Path: "/b", SeasonNum: 1, EpisodeNum: 2, TMDbID: 99}, + {Base: model.Base{ID: "c"}, LibraryID: "lib", Title: "Show", Path: "/c", SeasonNum: 2, EpisodeNum: 1, TMDbID: 99}, + } + for _, ep := range eps { + if err := repos.DB.Create(ep).Error; err != nil { + t.Fatal(err) + } + } + got, err := repos.Media.ListSeasonSiblings(ctx, eps[0]) + if err != nil { + t.Fatal(err) + } + if len(got) != 1 || got[0].ID != "b" { + t.Fatalf("siblings = %#v, want only same-season ep b", got) + } +} + +func TestListPrewarmCandidatesOrdersRecentPlaysFirst(t *testing.T) { + repos := newSegmentTestRepos(t) + ctx := t.Context() + now := time.Now() + for _, m := range []*model.Media{ + {Base: model.Base{ID: "cold"}, Path: "/cold", TMDbID: 1}, + {Base: model.Base{ID: "hot"}, Path: "/hot", TMDbID: 2}, + } { + if err := repos.DB.Create(m).Error; err != nil { + t.Fatal(err) + } + } + if err := repos.DB.Create(&model.PlaybackHistory{ + Base: model.Base{ID: "ph1"}, UserID: "u", MediaID: "hot", + }).Error; err != nil { + t.Fatal(err) + } + got, err := repos.MediaSegment.ListPrewarmCandidates( + ctx, "theintrodb", now.Add(-time.Hour), now.Add(-time.Hour), 10, + ) + if err != nil { + t.Fatal(err) + } + if len(got) < 2 { + t.Fatalf("candidates = %#v, want both", got) + } + if got[0].ID != "hot" { + t.Fatalf("first = %s, want hot (recently played)", got[0].ID) + } +} diff --git a/internal/service/media_segment.go b/internal/service/media_segment.go index 3e53031..b345fb3 100644 --- a/internal/service/media_segment.go +++ b/internal/service/media_segment.go @@ -18,9 +18,21 @@ import ( // 片段数据的缓存时长。命中过说明社区库里已有记录、数据很少变动,可以放很久; // 未命中说明这部片还没人贡献,隔一段时间再试一次即可——负缓存是必须的,否则 // 每次播放一部没有片段数据的影片都会打一次外网。 +// +// 未命中 TTL 刻意短于命中:社区库在持续补充,尤其是热门剧,24h 重试一次比 +// 锁死 7 天更能跟上贡献节奏;预热任务也会优先扫最近播放过的 miss。 const ( segmentFoundTTL = 30 * 24 * time.Hour - segmentMissingTTL = 7 * 24 * time.Hour + segmentMissingTTL = 24 * time.Hour + + // SegmentSourceManual / Propagated 与 IntroDBSource 并列,播放时按优先级合并。 + SegmentSourceManual = "manual" + SegmentSourcePropagated = "propagated" + + // segmentPrewarmDefaultLimit 是一次预热任务最多处理的媒体数,避免单次跑太久。 + segmentPrewarmDefaultLimit = 200 + // segmentPrewarmInterval 是预热请求之间的间隔,压低对 TheIntroDB 的 429。 + segmentPrewarmInterval = time.Second ) // SegmentView 是播放器消费的最小片段结构,避免把库内字段(source 等)暴露给前端。 @@ -44,14 +56,29 @@ type MediaSegmentService struct { log *zap.Logger repo *repository.Container introdb *IntroDBService - // probe 负责从文件内嵌章节里提取片头/片尾(异步、落库)。未注入时只用 - // TheIntroDB。 + // probe 负责异步补齐媒体信息(主要是 STRM 时长)。未注入时只用 TheIntroDB。 probe *MediaProbeService + // prewarmGap 是预热两次 IntroDB 请求之间的间隔;测试可设为 0。 + prewarmGap time.Duration + now func() time.Time } // NewMediaSegmentService is the constructor. func NewMediaSegmentService(log *zap.Logger, repo *repository.Container) *MediaSegmentService { - return &MediaSegmentService{log: log, repo: repo} + return &MediaSegmentService{ + log: log, + repo: repo, + prewarmGap: segmentPrewarmInterval, + now: time.Now, + } +} + +// SetPrewarmGap overrides the delay between prewarm fetches (tests use 0). +func (s *MediaSegmentService) SetPrewarmGap(gap time.Duration) *MediaSegmentService { + if s != nil { + s.prewarmGap = gap + } + return s } // SetIntroDB wires the provider. Without it the service only reads cached rows. @@ -62,8 +89,8 @@ func (s *MediaSegmentService) SetIntroDB(p *IntroDBService) *MediaSegmentService return s } -// SetProbe wires the in-file chapter extractor. Without it the ffprobe source -// simply yields nothing and auto falls back to TheIntroDB. +// SetProbe wires the async media probe (duration backfill for STRM). Without it +// open-ended credits still work once duration is known from elsewhere. func (s *MediaSegmentService) SetProbe(p *MediaProbeService) *MediaSegmentService { if s != nil { s.probe = p @@ -121,28 +148,73 @@ func mediaProbeSettled(row *model.MediaProbe) bool { } // introDBSegments 读社区库的片段,缓存过期时按调用方的预算抓一次并落库。 +// 返回值是按来源优先级合并后的结果(manual > theintrodb > propagated)。 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 - } ledger, err := s.repo.MediaSegment.GetFetch(ctx, m.ID, IntroDBSource) if err != nil { + // 调用方预算耗尽时仍读本地缓存,绝不能让播放接口报错。 + if ctx.Err() != nil { + return s.mergeForPlayback(context.WithoutCancel(ctx), m.ID) + } return nil, err } - if ledger != nil && ledgerFresh(ledger) { - return cached, nil + if ledger == nil || !ledgerFresh(ledger) { + if _, _, err := s.refresh(ctx, m); err != nil { + // 社区库不可达或返回异常:沿用已有缓存,不影响播放;也不写负缓存。 + logIntroDBFailure(s.log, 0, err) + } } - refreshed, attempted, err := s.refresh(ctx, m) + // 合并读库只碰本地,脱离调用方 deadline,避免外网抓取耗尽预算后读缓存也失败。 + return s.mergeForPlayback(context.WithoutCancel(ctx), m.ID) +} + +// mergeForPlayback 合并同一媒体上多来源片段。同 kind 只保留优先级最高的来源。 +func (s *MediaSegmentService) mergeForPlayback(ctx context.Context, mediaID string) ([]model.MediaSegment, error) { + rows, err := s.repo.MediaSegment.ListByMedia(ctx, mediaID) if err != nil { - // 社区库不可达或返回异常:沿用已有缓存,不影响播放。 - logIntroDBFailure(s.log, 0, err) - return cached, nil + return nil, err } - if !attempted { - return cached, nil + return preferSegmentsBySource(rows), nil +} + +// segmentSourcePriority 数值越大越优先。未知来源视为最低,避免挡住已知源。 +func segmentSourcePriority(source string) int { + switch source { + case SegmentSourceManual: + return 3 + case IntroDBSource: + return 2 + case SegmentSourcePropagated: + return 1 + default: + return 0 } - return refreshed, nil +} + +// preferSegmentsBySource 按 kind 选取最高优先级来源的全部区间。 +func preferSegmentsBySource(rows []model.MediaSegment) []model.MediaSegment { + bestPri := make(map[string]int, 4) + byKind := make(map[string][]model.MediaSegment, 4) + for _, row := range rows { + pri := segmentSourcePriority(row.Source) + cur, seen := bestPri[row.Kind] + if !seen || pri > cur { + bestPri[row.Kind] = pri + byKind[row.Kind] = []model.MediaSegment{row} + continue + } + if pri == cur { + byKind[row.Kind] = append(byKind[row.Kind], row) + } + } + out := make([]model.MediaSegment, 0, len(rows)) + for _, kind := range []string{ + model.SegmentKindIntro, model.SegmentKindRecap, + model.SegmentKindCredits, model.SegmentKindPreview, + } { + out = append(out, byKind[kind]...) + } + return out } func (s *MediaSegmentService) debug(message, mediaID string, err error) { @@ -210,9 +282,111 @@ func (s *MediaSegmentService) refresh(ctx context.Context, m *model.Media) ([]mo }); err != nil { return nil, true, err } + s.propagateIntroToSeason(fetchCtx, m, rows) return rows, true, nil } +// propagateIntroToSeason 把本集 IntroDB 命中的 intro 复制到同季还没有 +// theintrodb intro 的兄弟集。片尾/预告不传播:各集时长与片尾位置经常不同。 +func (s *MediaSegmentService) propagateIntroToSeason(ctx context.Context, m *model.Media, rows []model.MediaSegment) { + if s == nil || s.repo == nil || m == nil || m.SeasonNum <= 0 { + return + } + intros := introSpansFrom(rows) + if len(intros) == 0 { + return + } + siblings, err := s.repo.Media.ListSeasonSiblings(ctx, m) + if err != nil { + s.debug("list season siblings failed", m.ID, err) + return + } + for i := range siblings { + sib := &siblings[i] + existing, err := s.repo.MediaSegment.ListByMediaSource(ctx, sib.ID, IntroDBSource) + if err != nil { + s.debug("list sibling introdb segments failed", sib.ID, err) + continue + } + if hasSegmentKind(existing, model.SegmentKindIntro) { + continue + } + copied := make([]model.MediaSegment, 0, len(intros)) + for _, intro := range intros { + copied = append(copied, model.MediaSegment{ + MediaID: sib.ID, + SeriesID: sib.SeriesID, + Kind: model.SegmentKindIntro, + StartMs: intro.StartMs, + EndMs: intro.EndMs, + Source: SegmentSourcePropagated, + }) + } + if err := s.repo.MediaSegment.ReplaceForMedia(ctx, sib.ID, SegmentSourcePropagated, copied); err != nil { + s.debug("propagate intro failed", sib.ID, err) + } + } +} + +func introSpansFrom(rows []model.MediaSegment) []model.MediaSegment { + out := make([]model.MediaSegment, 0, 1) + for _, row := range rows { + if row.Kind == model.SegmentKindIntro { + out = append(out, row) + } + } + return out +} + +func hasSegmentKind(rows []model.MediaSegment, kind string) bool { + for _, row := range rows { + if row.Kind == kind { + return true + } + } + return false +} + +// Prewarm 批量向 TheIntroDB 补齐过期/缺失账本的可查询媒体。返回实际发起过查询的数量。 +func (s *MediaSegmentService) Prewarm(ctx context.Context, limit int) (int, error) { + if s == nil || s.repo == nil || s.introdb == nil { + return 0, nil + } + if limit <= 0 { + limit = segmentPrewarmDefaultLimit + } + now := time.Now() + if s.now != nil { + now = s.now() + } + candidates, err := s.repo.MediaSegment.ListPrewarmCandidates( + ctx, IntroDBSource, now.Add(-segmentMissingTTL), now.Add(-segmentFoundTTL), limit, + ) + if err != nil { + return 0, err + } + attempted := 0 + for i := range candidates { + if err := ctx.Err(); err != nil { + return attempted, err + } + m := &candidates[i] + _, did, err := s.refresh(ctx, m) + if err != nil { + logIntroDBFailure(s.log, m.TMDbID, err) + } + if did { + attempted++ + } + if i+1 < len(candidates) && s.prewarmGap > 0 { + if !waitForIntroDBRetry(ctx, s.prewarmGap) { + return attempted, ctx.Err() + } + } + } + return attempted, nil +} + // queryIDs resolves the provider query key. Movies use their own TMDb id. // // 剧集需要「剧集级」TMDb id 加季/集。优先取 Series.TMDbID;但有些刮削路径 diff --git a/internal/service/media_segment_test.go b/internal/service/media_segment_test.go index 6d52984..6cf78f3 100644 --- a/internal/service/media_segment_test.go +++ b/internal/service/media_segment_test.go @@ -385,6 +385,123 @@ func TestListForPlaybackRespectsCallerDeadline(t *testing.T) { } } +func TestPreferSegmentsBySourcePrefersManualThenIntroDBThenPropagated(t *testing.T) { + rows := []model.MediaSegment{ + {Kind: model.SegmentKindIntro, StartMs: 1, EndMs: 2, Source: SegmentSourcePropagated}, + {Kind: model.SegmentKindIntro, StartMs: 10, EndMs: 20, Source: IntroDBSource}, + {Kind: model.SegmentKindIntro, StartMs: 100, EndMs: 200, Source: SegmentSourceManual}, + {Kind: model.SegmentKindCredits, StartMs: 1000, EndMs: 0, Source: SegmentSourcePropagated}, + {Kind: model.SegmentKindCredits, StartMs: 2000, EndMs: 0, Source: IntroDBSource}, + } + got := preferSegmentsBySource(rows) + if len(got) != 2 { + t.Fatalf("got %#v, want manual intro + introdb credits", got) + } + if got[0].Source != SegmentSourceManual || got[0].StartMs != 100 { + t.Fatalf("intro = %#v, want manual", got[0]) + } + if got[1].Source != IntroDBSource || got[1].StartMs != 2000 { + t.Fatalf("credits = %#v, want theintrodb", got[1]) + } +} + +func TestListForPlaybackPropagatesIntroToSeasonSiblings(t *testing.T) { + svc, repos, _ := newSegmentServiceFixture(t, writeJSONBody(introDBTVPayload)) + svc.SetPrewarmGap(0) + ctx := t.Context() + episodes := []*model.Media{ + {Base: model.Base{ID: "ep-1"}, LibraryID: "lib-anime", Title: "便·当", + Path: "/anime/S01E01.mkv", SeasonNum: 1, EpisodeNum: 1, TMDbID: 61970}, + {Base: model.Base{ID: "ep-2"}, LibraryID: "lib-anime", Title: "便·当", + Path: "/anime/S01E02.mkv", SeasonNum: 1, EpisodeNum: 2, TMDbID: 61970}, + } + for _, ep := range episodes { + if err := repos.DB.Create(ep).Error; err != nil { + t.Fatal(err) + } + } + + rows, err := svc.ListForPlayback(ctx, episodes[0]) + if err != nil { + t.Fatal(err) + } + if len(rows) != 2 { + t.Fatalf("ep1 rows = %#v, want intro+credits from provider", rows) + } + + sib, err := repos.MediaSegment.ListByMediaSource(ctx, "ep-2", SegmentSourcePropagated) + if err != nil { + t.Fatal(err) + } + if len(sib) != 1 || sib[0].Kind != model.SegmentKindIntro { + t.Fatalf("sibling propagated = %#v, want intro only", sib) + } + if sib[0].StartMs != 228_664 || sib[0].EndMs != 246_143 { + t.Fatalf("propagated window = %#v", sib[0]) + } + + // 兄弟集播放时应直接看到传播来的 intro(即使自己还没打过 IntroDB)。 + if err := repos.MediaSegment.UpsertFetch(ctx, &model.MediaSegmentFetch{ + MediaID: "ep-2", Source: IntroDBSource, FetchedAt: time.Now(), Found: false, + }); err != nil { + t.Fatal(err) + } + merged, err := svc.ListForPlayback(ctx, episodes[1]) + if err != nil { + t.Fatal(err) + } + if len(merged) != 1 || merged[0].Kind != model.SegmentKindIntro || merged[0].Source != SegmentSourcePropagated { + t.Fatalf("ep2 playback = %#v, want propagated intro", merged) + } +} + +func TestPrewarmFetchesStaleMissesPreferringRecentPlays(t *testing.T) { + svc, repos, calls := newSegmentServiceFixture(t, writeJSONBody(introDBMoviePayload)) + svc.SetPrewarmGap(0) + ctx := t.Context() + + old := &model.Media{Base: model.Base{ID: "mv-old"}, Path: "/a.mkv", TMDbID: 111} + hot := &model.Media{Base: model.Base{ID: "mv-hot"}, Path: "/b.mkv", TMDbID: 27205} + for _, m := range []*model.Media{old, hot} { + if err := repos.DB.Create(m).Error; err != nil { + t.Fatal(err) + } + } + // 两条都是过期 miss,热播的应优先被预热。 + stale := time.Now().Add(-segmentMissingTTL - time.Hour) + for _, id := range []string{"mv-old", "mv-hot"} { + if err := repos.MediaSegment.UpsertFetch(ctx, &model.MediaSegmentFetch{ + MediaID: id, Source: IntroDBSource, FetchedAt: stale, Found: false, + }); err != nil { + t.Fatal(err) + } + } + if err := repos.DB.Create(&model.PlaybackHistory{ + Base: model.Base{ID: "h-1"}, UserID: "u1", MediaID: "mv-hot", + }).Error; err != nil { + t.Fatal(err) + } + + n, err := svc.Prewarm(ctx, 1) + if err != nil { + t.Fatal(err) + } + if n != 1 { + t.Fatalf("attempted = %d, want 1", n) + } + if got := atomic.LoadInt32(calls); got != 1 { + t.Fatalf("provider calls = %d, want 1", got) + } + ledger, err := repos.MediaSegment.GetFetch(ctx, "mv-hot", IntroDBSource) + if err != nil || ledger == nil || !ledger.Found { + t.Fatalf("hot ledger = %#v err=%v, want found", ledger, err) + } + oldLedger, err := repos.MediaSegment.GetFetch(ctx, "mv-old", IntroDBSource) + if err != nil || oldLedger == nil || oldLedger.Found { + t.Fatalf("old ledger should remain a miss, got %#v err=%v", oldLedger, err) + } +} + func TestLedgerFreshUsesLongerTTLWhenDataWasFound(t *testing.T) { now := time.Now() found := &model.MediaSegmentFetch{FetchedAt: now.Add(-segmentMissingTTL), Found: true} @@ -399,3 +516,4 @@ func TestLedgerFreshUsesLongerTTLWhenDataWasFound(t *testing.T) { t.Fatal("a missing ledger must not be considered fresh") } } + diff --git a/internal/service/scheduler.go b/internal/service/scheduler.go index 0386831..a6323a0 100644 --- a/internal/service/scheduler.go +++ b/internal/service/scheduler.go @@ -8,6 +8,7 @@ // organize_source opt-in — organize the configured staging folder. // transcode_cleanup every 24 h — purge HLS transcode artefacts // older than 24 h. +// segment_prewarm every 6 h — fill IntroDB skip segments for queryable media. // // Each job runs at most once at a time (an in-flight run blocks the // next tick). All work happens on a long-lived background context so @@ -42,6 +43,8 @@ type SchedulerService struct { imagesPolicyProvider func() ImageCachePolicy + segments *MediaSegmentService + mu sync.Mutex stopCh chan struct{} jobs []*scheduledJob @@ -145,6 +148,14 @@ func (s *SchedulerService) Start(ctx context.Context) { run: s.jobCleanImageCache, }, } + // 片头预热只在注入了 Segments 时注册,避免测试跑无转外网任务。 + if s.segments != nil { + s.jobs = append(s.jobs, &scheduledJob{ + name: "segment_prewarm", + interval: segmentPrewarmJobInterval, + run: s.jobSegmentPrewarm, + }) + } // 到期提醒只在配置了巡检器时注册,避免测试与未启用通知的部署跑空转任务。 if s.expiryWatcher != nil { s.jobs = append(s.jobs, &scheduledJob{ @@ -155,8 +166,8 @@ func (s *SchedulerService) Start(ctx context.Context) { } for _, j := range s.jobs { initialDelay := 15 * time.Second - if j.name == "library_scan" || j.name == "organize_source" { - // 重启后不立即整库重扫/整理下载目录:更新窗口恰是登录高峰, + if j.name == "library_scan" || j.name == "organize_source" || j.name == "segment_prewarm" { + // 重启后不立即整库重扫/整理/预热:更新窗口恰是登录高峰, // 15 秒即全量 walk + ffprobe 曾把 CPU/磁盘打满导致无法登录。 // 首轮等满一个完整周期再跑,平时节奏不变。 initialDelay = j.interval diff --git a/internal/service/scheduler_segment_jobs.go b/internal/service/scheduler_segment_jobs.go new file mode 100644 index 0000000..bc0f38f --- /dev/null +++ b/internal/service/scheduler_segment_jobs.go @@ -0,0 +1,52 @@ +package service + +import ( + "context" + "time" + + "go.uber.org/zap" +) + +const ( + segmentPrewarmJobInterval = 6 * time.Hour + // SegmentPrewarmEnabledKey 控制是否后台预热 TheIntroDB 片段;默认关闭。 + SegmentPrewarmEnabledKey = "segment.prewarm_enabled" +) + +// SetSegments wires the IntroDB prewarm job. Without it the job is not registered. +func (s *SchedulerService) SetSegments(segments *MediaSegmentService) { + if s != nil { + s.segments = segments + } +} + +func (s *SchedulerService) jobSegmentPrewarm(ctx context.Context) error { + if s == nil || s.segments == nil { + return nil + } + manual, _ := ctx.Value(schedulerManualRunKey{}).(bool) + if !manual && !s.segmentPrewarmEnabled(ctx) { + return nil + } + n, err := s.segments.Prewarm(ctx, segmentPrewarmDefaultLimit) + if s.log != nil { + s.log.Info("segment prewarm finished", + zap.Int("attempted", n), + zap.Error(err), + ) + } + return err +} + +// segmentPrewarmEnabled reports whether the operator opted into background +// IntroDB prewarm. Defaults to false so playback-on-demand remains the only path. +func (s *SchedulerService) segmentPrewarmEnabled(ctx context.Context) bool { + if s.repo == nil || s.repo.Setting == nil { + return false + } + v, err := s.repo.Setting.Get(ctx, SegmentPrewarmEnabledKey) + if err != nil { + return false + } + return parseBoolSetting(v, false) +} diff --git a/internal/service/scheduler_segment_jobs_test.go b/internal/service/scheduler_segment_jobs_test.go new file mode 100644 index 0000000..0e02c3b --- /dev/null +++ b/internal/service/scheduler_segment_jobs_test.go @@ -0,0 +1,60 @@ +package service + +import ( + "context" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + + "go.uber.org/zap" + + "github.com/truewhile/MeBox/internal/model" + "github.com/truewhile/MeBox/internal/repository" +) + +func TestSegmentPrewarmDefaultsOff(t *testing.T) { + repos := repository.New(newServiceTestDB(t, &model.Setting{})) + scheduler := NewSchedulerService(zap.NewNop(), repos, nil, nil, nil, nil, "") + if scheduler.segmentPrewarmEnabled(t.Context()) { + t.Fatal("prewarm must default to off") + } +} + +func TestJobSegmentPrewarmSkippedWhenDisabled(t *testing.T) { + var calls int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + atomic.AddInt32(&calls, 1) + w.WriteHeader(http.StatusNotFound) + })) + t.Cleanup(server.Close) + + repos := repository.New(newServiceTestDB(t)) + if err := repos.DB.Create(&model.Media{ + Base: model.Base{ID: "mv-1"}, Path: "/a.mkv", TMDbID: 1, + }).Error; err != nil { + t.Fatal(err) + } + segments := NewMediaSegmentService(zap.NewNop(), repos). + SetIntroDB(NewIntroDBService(zap.NewNop()).SetBaseURL(server.URL)). + SetPrewarmGap(0) + scheduler := NewSchedulerService(zap.NewNop(), repos, nil, nil, nil, nil, "") + scheduler.SetSegments(segments) + + if err := scheduler.jobSegmentPrewarm(t.Context()); err != nil { + t.Fatal(err) + } + if got := atomic.LoadInt32(&calls); got != 0 { + t.Fatalf("provider calls = %d, want 0 when prewarm is disabled", got) + } + + if err := repos.Setting.Set(t.Context(), SegmentPrewarmEnabledKey, "true"); err != nil { + t.Fatal(err) + } + if err := scheduler.jobSegmentPrewarm(context.Background()); err != nil { + t.Fatal(err) + } + if got := atomic.LoadInt32(&calls); got != 1 { + t.Fatalf("provider calls = %d, want 1 after enabling prewarm", got) + } +} diff --git a/internal/service/service_builder.go b/internal/service/service_builder.go index 9ab8860..8b2a6a7 100644 --- a/internal/service/service_builder.go +++ b/internal/service/service_builder.go @@ -192,6 +192,7 @@ func (b *serviceContainerBuilder) initAccessAndStorageServices() { ) b.c.Scheduler.SetTaskTracker(b.c.Tasks) b.c.Scheduler.SetOrganizePipeline(b.c.OrganizePipeline) + b.c.Scheduler.SetSegments(b.c.Segments) b.c.Scheduler.SetImageCachePolicyProvider(func() ImageCachePolicy { if b.cfg == nil { return ImageCachePolicy{} diff --git a/web/src/pages/settingsGroupGeneral.ts b/web/src/pages/settingsGroupGeneral.ts index a9e98d2..f8c5a48 100644 --- a/web/src/pages/settingsGroupGeneral.ts +++ b/web/src/pages/settingsGroupGeneral.ts @@ -23,6 +23,13 @@ export const generalSettingsGroup: SettingGroup = { hint: '默认关闭。开启后宿主机不再进行任何 FFmpeg 转码,所有播放交给第三方客户端(Infuse / VLC / Emby 客户端等)或浏览器本地解码直连(direct play / 302 直链),大幅降低宿主机 CPU 占用。若客户端不支持源编码可能无法播放。', defaultValue: 'false', }, + { + key: 'segment.prewarm_enabled', + label: '后台预热 IntroDB 片头片段', + type: 'toggle', + hint: '默认关闭。开启后每 6 小时在后台向 TheIntroDB 拉取可识别媒体的片头/片尾区间并写入本地;关闭时仅在播放时按需请求。', + defaultValue: 'false', + }, { key: 'transcode.enabled', label: '启用转码',