Compare commits

...

5 Commits

Author SHA1 Message Date
truewhile e65accf2cd Optimize Emby caching and image resize concurrency 2026-09-12 12:46:30 +08:00
truewhile f54ef2228c Merge pull request #32 from truewhile/cursor/fix-emby-hero-backdrop
修复 Emby 大海报在剧集图缓存被淘汰后显示占位图
2026-09-12 11:47:38 +08:00
truewhile 3c2bfb14ca 修复 Emby 大海报在剧集图缓存被淘汰后显示占位图
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-12 11:43:51 +08:00
truewhile 8f64f961ac 缓存弹幕抓取结果并合并并发请求 2026-09-12 11:02:24 +08:00
truewhile 6797394c5b Serve font assets with caching and test missing-font handling 2026-09-12 10:33:51 +08:00
21 changed files with 957 additions and 99 deletions
+29
View File
@@ -76,6 +76,9 @@ func TestServeSPAServesAssetsImmutableAndBypassesAPIRoutes(t *testing.T) {
if err := os.MkdirAll(filepath.Join(webDir, "assets"), 0o755); err != nil {
t.Fatal(err)
}
if err := os.MkdirAll(filepath.Join(webDir, "fonts"), 0o755); err != nil {
t.Fatal(err)
}
if err := os.MkdirAll(filepath.Join(webDir, "brand"), 0o755); err != nil {
t.Fatal(err)
}
@@ -85,6 +88,9 @@ func TestServeSPAServesAssetsImmutableAndBypassesAPIRoutes(t *testing.T) {
if err := os.WriteFile(filepath.Join(webDir, "assets", "app.js"), []byte("console.log('ok')"), 0o644); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(webDir, "fonts", "geist-400.woff2"), []byte("wOF2-test-font"), 0o644); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(webDir, "brand", "mebox-logo.svg"), []byte("<svg></svg>"), 0o644); err != nil {
t.Fatal(err)
}
@@ -105,6 +111,29 @@ func TestServeSPAServesAssetsImmutableAndBypassesAPIRoutes(t *testing.T) {
t.Fatalf("asset Cache-Control = %q, want immutable", got)
}
fontReq := httptest.NewRequest(http.MethodGet, "/fonts/geist-400.woff2", nil)
fontResp := httptest.NewRecorder()
router.ServeHTTP(fontResp, fontReq)
if fontResp.Code != http.StatusOK {
t.Fatalf("font status = %d, want 200", fontResp.Code)
}
if got := fontResp.Header().Get("Cache-Control"); !strings.Contains(got, "max-age=86400") {
t.Fatalf("font Cache-Control = %q, want max-age=86400", got)
}
if got := fontResp.Body.String(); got != "wOF2-test-font" {
t.Fatalf("font body = %q, want wOF2-test-font", got)
}
missingFontReq := httptest.NewRequest(http.MethodGet, "/fonts/missing.woff2", nil)
missingFontResp := httptest.NewRecorder()
router.ServeHTTP(missingFontResp, missingFontReq)
if missingFontResp.Code != http.StatusNotFound {
t.Fatalf("missing font status = %d, want 404", missingFontResp.Code)
}
if strings.Contains(missingFontResp.Body.String(), "index") {
t.Fatalf("missing font should not serve SPA index: %q", missingFontResp.Body.String())
}
brandReq := httptest.NewRequest(http.MethodGet, "/brand/mebox-logo.svg", nil)
brandResp := httptest.NewRecorder()
router.ServeHTTP(brandResp, brandReq)
+12 -5
View File
@@ -59,6 +59,13 @@ func serveSPA(r *gin.Engine, root fs.FS) {
c.Next()
})
assets.GET("/*filepath", serveFSDir(root, "assets"))
fonts := r.Group("/fonts")
fonts.Use(middleware.GzipStatic())
fonts.Use(func(c *gin.Context) {
c.Header("Cache-Control", "public, max-age=86400")
c.Next()
})
fonts.GET("/*filepath", serveFSDir(root, "fonts"))
brand := r.Group("/brand")
brand.Use(func(c *gin.Context) {
setNoCacheHeaders(c)
@@ -70,11 +77,11 @@ func serveSPA(r *gin.Engine, root fs.FS) {
r.GET(rootFile, serveFSFile(root, name))
r.HEAD(rootFile, serveFSFile(root, name))
}
r.NoRoute(middleware.GzipStatic(), func(c *gin.Context) {
if handler.TryHandleEmbyNormalizedRoute(c, r) {
return
}
path := c.Request.URL.Path
r.NoRoute(middleware.GzipStatic(), func(c *gin.Context) {
if handler.TryHandleEmbyNormalizedRoute(c, r) {
return
}
path := c.Request.URL.Path
if shouldBypassSPAFallback(path) {
c.Status(http.StatusNotFound)
return
+1 -1
View File
@@ -20,6 +20,7 @@ require (
go.uber.org/zap v1.27.0
golang.org/x/crypto v0.49.0
golang.org/x/image v0.37.0
golang.org/x/sync v0.20.0
golang.org/x/sys v0.42.0
golang.org/x/time v0.15.0
gopkg.in/yaml.v3 v3.0.1
@@ -88,7 +89,6 @@ require (
golang.org/x/arch v0.25.0 // indirect
golang.org/x/exp v0.0.0-20251023183803-a4bb9ffd2546 // indirect
golang.org/x/net v0.52.0 // indirect
golang.org/x/sync v0.20.0 // indirect
golang.org/x/text v0.35.0 // indirect
google.golang.org/protobuf v1.36.11 // indirect
gopkg.in/ini.v1 v1.67.0 // indirect
+35
View File
@@ -824,3 +824,38 @@ func TestDanmakuFetchEmbyRemoteStreamFailedFallsBackToSearch(t *testing.T) {
require.Equal(t, "降级搜索番剧", res.AnimeTitle)
require.Contains(t, res.Raw, "降级搜索弹幕")
}
// 即使刮削元数据完整,也应先走准确率最高的 hash 层。
func TestDanmakuFetchPrefersHashOverCompleteMetadata(t *testing.T) {
videoPath, wantHash := writeDanmakuTestVideo(t, "测试动画.第01话.mkv")
var seen string
official := danmakuOfficialServer(t,
`{"success":true,"isMatched":true,"matches":[{"episodeId":25484,"animeId":1001,"animeTitle":"官方测试动画","episodeTitle":"第1话"}]}`,
`<?xml version="1.0"?><i><d p="0.5,1,16777215,user1">hash命中</d></i>`,
&seen)
overrideDanmakuOfficialBase(t, official.URL)
source := newDanmakuSourceServerWithSearch(t,
`{"hasMore":false,"animes":[{"animeId":2002,"animeTitle":"官方测试动画","episodes":[{"episodeId":25484,"episodeTitle":"第1话"}]}]}`)
svc := newDanmakuTestService(t)
ctx := context.Background()
require.NoError(t, svc.repo.Setting.Set(ctx, DanmakuSourceKey, source.URL()))
m := model.Media{
Title: "测试动画",
Path: videoPath,
SizeBytes: 123,
EpisodeNum: 1,
EpisodeTitle: "起始与终结的序章",
Year: 2011,
}
m.ID = "hash-first"
require.NoError(t, svc.repo.DB.Create(&m).Error)
res, err := svc.Fetch(ctx, "hash-first", "", "")
require.NoError(t, err)
require.Equal(t, "hash", res.MatchMode)
require.NotEmpty(t, res.Raw)
require.Contains(t, seen, wantHash)
}
+163 -5
View File
@@ -20,6 +20,8 @@ import (
"time"
"unicode"
"golang.org/x/sync/singleflight"
"go.uber.org/zap"
"github.com/truewhile/MeBox/internal/model"
@@ -136,6 +138,12 @@ type DanmakuService struct {
hashCacheMu sync.Mutex
hashCache map[string]string // stamp → 16MB-prefix MD5
// resultCache 缓存整条弹幕抓取结果,避免同一集重复播放时重走
// 「16MB 哈希 + 上游搜索 + 评论拉取」这条高延迟链路。
resultCacheMu sync.Mutex
resultCache map[string]danmakuResultCacheEntry
fetchGroup singleflight.Group
}
func danmakuHTTPClient() *http.Client {
@@ -147,10 +155,11 @@ func NewDanmakuService(log *zap.Logger, repo *repository.Container) *DanmakuServ
log = zap.NewNop()
}
return &DanmakuService{
log: log,
repo: repo,
client: danmakuHTTPClient(),
hashCache: make(map[string]string),
log: log,
repo: repo,
client: danmakuHTTPClient(),
hashCache: make(map[string]string),
resultCache: make(map[string]danmakuResultCacheEntry),
}
}
@@ -170,6 +179,121 @@ func (s *DanmakuService) SetRemoteMediaResolver(resolve DanmakuRemoteMediaResolv
}
}
// 弹幕抓取结果缓存:同一集在 TTL 内重复播放时直接返回,避免重复执行
// 「16MB 前置哈希 → 上游搜索 → 评论拉取」这条高延迟链路。弹幕库内容变化
// 很慢,而用户通常会在短时间内反复切集/回看,因此 TTL 取 6 小时。
const (
danmakuResultCacheTTL = 6 * time.Hour
danmakuResultCacheMaxEntries = 64
// 候选列表由上游搜索决定,相对稳定但可能随源更新,缓存 10 分钟。
danmakuResultCacheCandidateTTL = 10 * time.Minute
// 空结果可能只是上游临时抖动,只短暂缓存,避免长时间看不到弹幕。
danmakuResultCacheEmptyTTL = time.Minute
)
type danmakuResultCacheEntry struct {
result *DanmakuFetchResult
expiresAt time.Time
storedAt time.Time
}
func danmakuResultCacheKey(source, mediaID, keyword, episodeID string, merge bool) string {
return strings.Join([]string{
strings.TrimSpace(source),
mediaID,
strings.TrimSpace(keyword),
strings.TrimSpace(episodeID),
strconv.FormatBool(merge),
}, "\x00")
}
func (s *DanmakuService) resultCacheGet(key string) (*DanmakuFetchResult, bool) {
if s == nil || key == "" {
return nil, false
}
now := time.Now()
s.resultCacheMu.Lock()
defer s.resultCacheMu.Unlock()
entry, ok := s.resultCache[key]
if !ok {
return nil, false
}
if now.After(entry.expiresAt) {
delete(s.resultCache, key)
return nil, false
}
return cloneDanmakuFetchResult(entry.result), true
}
func (s *DanmakuService) resultCachePut(key string, result *DanmakuFetchResult) {
if s == nil || key == "" || result == nil {
return
}
now := time.Now()
s.resultCacheMu.Lock()
defer s.resultCacheMu.Unlock()
if s.resultCache == nil {
s.resultCache = make(map[string]danmakuResultCacheEntry)
}
ttl := danmakuResultCacheTTLFor(result)
if ttl <= 0 {
return
}
if _, exists := s.resultCache[key]; !exists && len(s.resultCache) >= danmakuResultCacheMaxEntries {
oldestKey := ""
var oldest time.Time
for k, entry := range s.resultCache {
if oldestKey == "" || entry.storedAt.Before(oldest) {
oldestKey, oldest = k, entry.storedAt
}
}
delete(s.resultCache, oldestKey)
}
s.resultCache[key] = danmakuResultCacheEntry{
result: cloneDanmakuFetchResult(result),
expiresAt: now.Add(ttl),
storedAt: now,
}
}
// danmakuResultCacheTTLFor 按结果完整性选择缓存时长:拿到弹幕正文才值得
// 长缓存;只有候选列表时短缓存;空结果只缓存一分钟。
func danmakuResultCacheTTLFor(result *DanmakuFetchResult) time.Duration {
if result == nil {
return 0
}
if strings.TrimSpace(result.Raw) != "" {
return danmakuResultCacheTTL
}
if len(result.Candidates) > 0 {
return danmakuResultCacheCandidateTTL
}
return danmakuResultCacheEmptyTTL
}
// cloneDanmakuFetchResult 深拷贝切片字段,避免缓存命中后调用方修改共享数据。
func cloneDanmakuFetchResult(in *DanmakuFetchResult) *DanmakuFetchResult {
if in == nil {
return nil
}
out := *in
out.Candidates = cloneDanmakuAnimeList(in.Candidates)
out.Alternatives = cloneDanmakuAnimeList(in.Alternatives)
return &out
}
func cloneDanmakuAnimeList(in []DanmakuAnime) []DanmakuAnime {
if in == nil {
return nil
}
out := make([]DanmakuAnime, len(in))
for i, anime := range in {
out[i] = anime
out[i].Episodes = append([]DanmakuEpisode(nil), anime.Episodes...)
}
return out
}
// Config reads danmaku settings from the runtime settings table.
func (s *DanmakuService) Config(ctx context.Context) DanmakuRenderConfig {
cfg := DanmakuRenderConfig{
@@ -257,8 +381,42 @@ func (s *DanmakuService) Fetch(ctx context.Context, mediaID, keyword, episodeID
return s.FetchWithOptions(ctx, mediaID, keyword, episodeID, DanmakuFetchOptions{})
}
// FetchWithOptions 是 Fetch 的带偏好版本。
// FetchWithOptions 是 Fetch 的带偏好版本。结果按「配置源 + 媒体 + 关键词 +
// 指定集 + 合并开关」缓存,并用 singleflight 合并并发请求,避免同一集被重复抓取。
func (s *DanmakuService) FetchWithOptions(ctx context.Context, mediaID, keyword, episodeID string, opts DanmakuFetchOptions) (*DanmakuFetchResult, error) {
if s == nil {
return nil, errors.New("danmaku service unavailable")
}
cfg := s.Config(ctx)
if !cfg.Enabled {
return &DanmakuFetchResult{DanmakuRenderConfig: cfg, SourceType: "auto"}, nil
}
key := danmakuResultCacheKey(cfg.Source, mediaID, keyword, episodeID, opts.MergeSources)
if cached, ok := s.resultCacheGet(key); ok {
return cached, nil
}
value, err, _ := s.fetchGroup.Do(key, func() (any, error) {
// 等待期间可能已有同一 key 的请求写入缓存。
if cached, ok := s.resultCacheGet(key); ok {
return cached, nil
}
res, err := s.fetchWithOptionsUncached(ctx, mediaID, keyword, episodeID, opts)
if err != nil {
// 与原实现一致:失败时仍把已填充的渲染配置/匹配信息交给调用方。
return res, err
}
s.resultCachePut(key, res)
return res, nil
})
res, _ := value.(*DanmakuFetchResult)
if err != nil {
return cloneDanmakuFetchResult(res), err
}
return cloneDanmakuFetchResult(res), nil
}
// fetchWithOptionsUncached 是未命中缓存时执行的原始抓取流程。
func (s *DanmakuService) fetchWithOptionsUncached(ctx context.Context, mediaID, keyword, episodeID string, opts DanmakuFetchOptions) (*DanmakuFetchResult, error) {
res := &DanmakuFetchResult{DanmakuRenderConfig: s.Config(ctx), SourceType: "auto"}
if !res.Enabled {
return res, nil
+78
View File
@@ -6,6 +6,8 @@ import (
"net/http"
"net/http/httptest"
"path/filepath"
"sync"
"sync/atomic"
"testing"
"github.com/glebarez/sqlite"
@@ -326,3 +328,79 @@ func TestDanmakuFetchDetectsJSONSource(t *testing.T) {
require.Equal(t, "xml", res2.SourceType)
require.Contains(t, res2.Raw, "弹幕A")
}
// 同一集在缓存 TTL 内重复请求时不应再访问上游(搜索和评论都只发一次)。
func TestDanmakuFetchCachesRepeatedRequests(t *testing.T) {
var searchCalls, commentCalls int32
mux := http.NewServeMux()
mux.HandleFunc("/api/v2/search/episodes", func(w http.ResponseWriter, r *http.Request) {
atomic.AddInt32(&searchCalls, 1)
w.Header().Set("Content-Type", "application/json")
fmt.Fprint(w, `{"hasMore":false,"animes":[{"animeId":1001,"animeTitle":"测试动画","episodes":[{"episodeId":25484,"episodeTitle":"第1话"}]}]}`)
})
mux.HandleFunc("/api/v2/comment/25484", func(w http.ResponseWriter, r *http.Request) {
atomic.AddInt32(&commentCalls, 1)
w.Header().Set("Content-Type", "application/xml")
fmt.Fprint(w, `<?xml version="1.0"?><i><d p="0.5,1,16777215,user1">缓存测试</d></i>`)
})
srv := httptest.NewServer(mux)
defer srv.Close()
svc := newDanmakuTestService(t)
ctx := context.Background()
require.NoError(t, svc.repo.Setting.Set(ctx, DanmakuSourceKey, srv.URL))
seedDanmakuMedia(t, svc, "cache-media", "测试动画", "", 1)
first, err := svc.Fetch(ctx, "cache-media", "", "")
require.NoError(t, err)
require.NotEmpty(t, first.Raw)
second, err := svc.Fetch(ctx, "cache-media", "", "")
require.NoError(t, err)
require.Equal(t, first.Raw, second.Raw)
require.EqualValues(t, 1, atomic.LoadInt32(&searchCalls))
require.EqualValues(t, 1, atomic.LoadInt32(&commentCalls))
}
// 并发的同一集请求应被 singleflight 合并,上游只被访问一次。
func TestDanmakuFetchCoalescesConcurrentRequests(t *testing.T) {
var searchCalls int32
mux := http.NewServeMux()
mux.HandleFunc("/api/v2/search/episodes", func(w http.ResponseWriter, r *http.Request) {
atomic.AddInt32(&searchCalls, 1)
w.Header().Set("Content-Type", "application/json")
fmt.Fprint(w, `{"hasMore":false,"animes":[{"animeId":1001,"animeTitle":"测试动画","episodes":[{"episodeId":25484,"episodeTitle":"第1话"}]}]}`)
})
mux.HandleFunc("/api/v2/comment/25484", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/xml")
fmt.Fprint(w, `<?xml version="1.0"?><i><d p="0.5,1,16777215,user1">并发测试</d></i>`)
})
srv := httptest.NewServer(mux)
defer srv.Close()
svc := newDanmakuTestService(t)
ctx := context.Background()
require.NoError(t, svc.repo.Setting.Set(ctx, DanmakuSourceKey, srv.URL))
seedDanmakuMedia(t, svc, "concurrent-media", "测试动画", "", 1)
const workers = 8
start := make(chan struct{})
errs := make(chan error, workers)
var wg sync.WaitGroup
for i := 0; i < workers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
<-start
_, err := svc.Fetch(ctx, "concurrent-media", "", "")
errs <- err
}()
}
close(start)
wg.Wait()
close(errs)
for err := range errs {
require.NoError(t, err)
}
require.EqualValues(t, 1, atomic.LoadInt32(&searchCalls))
}
+20 -8
View File
@@ -32,17 +32,14 @@ func (e *EmbyService) ImageURL(ctx context.Context, id, imageType string) (strin
}
return backdrop
}
if strings.HasPrefix(id, embyVirtualSeasonPrefix) {
if strings.HasPrefix(id, embyVirtualSeasonPrefix) || strings.HasPrefix(id, embyVirtualSeriesPrefix) {
if raw, ok := e.cachedArtworkURL(id, imageType); ok {
return raw, nil
}
return "", nil
}
if strings.HasPrefix(id, embyVirtualSeriesPrefix) {
if raw, ok := e.cachedArtworkURL(id, imageType); ok {
return raw, nil
}
return "", nil
// Latest JSON can outlive (or be served after) the in-memory artwork
// map. Rebuild from the library instead of handing the client a 1x1
// placeholder that it then caches as a successful image.
return e.resolveVirtualArtwork(ctx, id, imageType, pick)
}
m, err := e.repo.Media.FindByID(ctx, id)
if err == nil && m != nil {
@@ -71,6 +68,21 @@ func (e *EmbyService) ImageURL(ctx context.Context, id, imageType string) (strin
return "", nil
}
func (e *EmbyService) resolveVirtualArtwork(ctx context.Context, id, imageType string, pick func(primary, backdrop string) string) (string, error) {
if strings.HasPrefix(id, embyVirtualSeasonPrefix) {
season, ok, err := e.findSeasonGroup(ctx, id, "")
if err != nil || !ok {
return "", err
}
return pick(season.Series.PosterURL, season.Series.BackdropURL), nil
}
series, ok, err := e.findSeriesGroup(ctx, id, "")
if err != nil || !ok {
return "", err
}
return pick(series.PosterURL, series.BackdropURL), nil
}
// imageInfoTypes 是 GET /Items/{Id}/Images 会报告的图片类型。只列 MeBox
// 真正存储的两类:ImageURL 对 Thumb / Logo / Banner 等其余类型会回退到
// 主图,若一并列出会让客户端以为存在这些图并去请求,实际拿到的却是主图。
+22
View File
@@ -61,6 +61,9 @@ type EmbyService struct {
libraryCoverMu sync.Mutex
libraryCoverCache map[string]embyArtworkCacheEntry
peopleMu sync.RWMutex
peopleCache map[string]embyPeopleCacheEntry
}
// NewEmbyService is the constructor.
@@ -117,6 +120,18 @@ const (
embyVirtualCacheTTL = 10 * time.Minute
embyVisibilityCacheTTL = 30 * time.Second
embySeriesGroupingLimit = maxMediaSearchLimit
// Virtual artwork used to be wiped entirely once the in-memory map crossed
// a few thousand entries. A homepage refresh asks Latest for every library
// at once, so that wipe dropped the series the client was about to paint.
// Caps are sized for that fan-out; overflow evicts the oldest entries only.
embyVirtualSeriesCap = 8000
embyVirtualSeasonCap = 16000
embyVirtualArtworkCap = 24000
// Clients that already cached the 1x1 placeholder treat a stable tag as
// immutable. Virtual ids are the ones that served that placeholder, so
// only those tags get a suffix that forces a refetch.
embyVirtualPrimaryTagSuffix = "-p2"
embyVirtualBackdropTagSuffix = "-bd2"
)
var (
@@ -131,6 +146,13 @@ type embyVisibilityCacheEntry struct {
expiresAt time.Time
}
// embyPeopleCacheEntry avoids re-statting/decoding the same NFO for every
// list refresh. TV clients commonly request the same posters/items repeatedly.
type embyPeopleCacheEntry struct {
people []map[string]any
expiresAt time.Time
}
// Items paginates media in Emby's hierarchy. Episodic libraries are exposed as
// Series -> Season -> Episode so Infuse/Vidhub/SenPlayer stop treating every
// episode as a separate movie card. 带 embyremote~ 前缀的 ParentID / 搜索自动
+20
View File
@@ -3,12 +3,24 @@ package service
import (
"context"
"strings"
"time"
"github.com/truewhile/MeBox/internal/model"
"gorm.io/gorm"
)
func (e *EmbyService) ItemCounts(ctx context.Context, userID string) (map[string]any, error) {
cacheKey := e.embyItemsCacheKey("counts-v1", ItemsParams{UserID: userID})
var cached embyCountsCacheValue
if e.cache != nil && e.cache.GetJSON(ctx, cacheKey, &cached) {
return map[string]any{
"MovieCount": cached.MovieCount,
"SeriesCount": int(cached.SeriesCount),
"EpisodeCount": cached.EpisodeCount,
"ItemCount": cached.ItemCount,
}, nil
}
base := func() *gorm.DB {
q := e.repo.DB.WithContext(ctx).Model(&model.Media{}).Where("deleted_at IS NULL")
return e.applyUserMediaVisibility(ctx, q, userID)
@@ -34,6 +46,14 @@ func (e *EmbyService) ItemCounts(ctx context.Context, userID string) (map[string
return nil, err
}
if e.cache != nil {
e.cache.SetJSON(ctx, cacheKey, embyCountsCacheValue{
MovieCount: movieCount,
SeriesCount: int64(seriesCount),
EpisodeCount: episodeCount,
ItemCount: itemCount,
}, time.Duration(e.mediaCacheTTLSeconds())*time.Second)
}
return map[string]any{
"MovieCount": movieCount,
"SeriesCount": seriesCount,
+15 -5
View File
@@ -9,13 +9,22 @@ import (
)
type embyItemsCacheValue struct {
Items []map[string]any `json:"items"`
TotalRecordCount int64 `json:"total_record_count"`
StartIndex int `json:"start_index"`
Items []map[string]any `json:"items"`
TotalRecordCount int64 `json:"total_record_count"`
StartIndex int `json:"start_index"`
Artwork map[string]embyArtworkRef `json:"artwork,omitempty"`
}
type embyLatestCacheValue struct {
Items []map[string]any `json:"items"`
Items []map[string]any `json:"items"`
Artwork map[string]embyArtworkRef `json:"artwork,omitempty"`
}
type embyCountsCacheValue struct {
MovieCount int64 `json:"movie_count"`
SeriesCount int64 `json:"series_count"`
EpisodeCount int64 `json:"episode_count"`
ItemCount int64 `json:"item_count"`
}
func (e *EmbyService) embyItemsCacheKey(kind string, p ItemsParams) string {
@@ -43,7 +52,8 @@ func (e *EmbyService) embyItemsCacheKey(kind string, p ItemsParams) string {
}
func (e *EmbyService) embyLatestCacheKey(userID, parentID string, limit int) string {
sum := sha256.Sum256([]byte(strings.Join([]string{"latest", userID, parentID, strconv.Itoa(limit)}, "|")))
// v2: payload tags for virtual artwork changed so clients drop cached placeholders.
sum := sha256.Sum256([]byte(strings.Join([]string{"latest-v2", userID, parentID, strconv.Itoa(limit)}, "|")))
return "media:emby:" + hex.EncodeToString(sum[:])
}
+56 -8
View File
@@ -132,15 +132,16 @@ func (e *EmbyService) LatestItems(ctx context.Context, userID, parentID string,
cacheKey := e.embyLatestCacheKey(userID, parentID, limit)
var cached embyLatestCacheValue
if e.cache != nil && e.cache.GetJSON(ctx, cacheKey, &cached) {
e.rememberArtworkRefs(cached.Artwork)
return cached.Items, nil
}
q := e.repo.DB.WithContext(ctx).Model(&model.Media{}).Where("deleted_at IS NULL")
q = e.applyUserMediaVisibility(ctx, q, userID)
if parentID != "" {
if episodic, err := e.libraryIsEpisodic(ctx, parentID); err == nil && episodic {
out, err := e.latestSeriesItemsForLibrary(ctx, userID, parentID, limit)
out, artwork, err := e.latestSeriesItemsForLibrary(ctx, userID, parentID, limit)
if err == nil && e.cache != nil {
e.cache.SetJSON(ctx, cacheKey, embyLatestCacheValue{Items: out}, time.Duration(e.embyLatestCacheTTLSeconds())*time.Second)
e.cache.SetJSON(ctx, cacheKey, embyLatestCacheValue{Items: out, Artwork: artwork}, time.Duration(e.embyLatestCacheTTLSeconds())*time.Second)
}
return out, err
}
@@ -171,7 +172,7 @@ func (e *EmbyService) LatestItems(ctx context.Context, userID, parentID string,
return out, nil
}
func (e *EmbyService) latestSeriesItemsForLibrary(ctx context.Context, userID, libraryID string, limit int) ([]map[string]any, error) {
func (e *EmbyService) latestSeriesItemsForLibrary(ctx context.Context, userID, libraryID string, limit int) ([]map[string]any, map[string]embyArtworkRef, error) {
if limit <= 0 || limit > 100 {
limit = 20
}
@@ -180,7 +181,7 @@ func (e *EmbyService) latestSeriesItemsForLibrary(ctx context.Context, userID, l
q = e.applyUserMediaVisibility(ctx, q, userID)
var rows []model.Media
if err := q.Order(mediaReleaseOrderSQL(true)).Limit(embySeriesGroupingLimit).Find(&rows).Error; err != nil {
return nil, err
return nil, nil, err
}
groups := e.seriesGroupsFromMedia(ctx, rows)
sortSeriesGroups(groups, ItemsParams{SortBy: "premieredate", SortOrder: "Descending"})
@@ -191,7 +192,7 @@ func (e *EmbyService) latestSeriesItemsForLibrary(ctx context.Context, userID, l
for _, group := range groups {
items = append(items, e.seriesPayload(group))
}
return items, nil
return items, e.artworkRefsForSeriesGroups(groups), nil
}
// ResumeItems 列出有未完成播放进度的媒体。
@@ -519,11 +520,11 @@ func (e *EmbyService) itemPayload(ctx context.Context, m *model.Media, fav bool,
if seriesID != "" {
if sEntry, ok, _ := e.payloadSeriesEntry(ctx, seriesID); ok {
if sEntry.posterURL != "" {
item["SeriesPrimaryImageTag"] = seriesID
item["SeriesPrimaryImageTag"] = embyVirtualImageTag(seriesID, embyVirtualPrimaryTagSuffix)
}
if len(backdropTags) == 0 && sEntry.backdropURL != "" {
if len(backdropTags) == 0 && (sEntry.backdropURL != "" || sEntry.posterURL != "") {
item["ParentBackdropItemId"] = seriesID
item["ParentBackdropImageTags"] = []string{seriesID + "-bd"}
item["ParentBackdropImageTags"] = []string{embyVirtualImageTag(seriesID, embyVirtualBackdropTagSuffix)}
}
}
}
@@ -537,6 +538,14 @@ func (e *EmbyService) resolveMediaPeople(ctx context.Context, m *model.Media) []
if m == nil || strings.TrimSpace(m.Path) == "" {
return []map[string]any{}
}
cacheKey := strings.TrimSpace(m.ID)
if cacheKey == "" {
cacheKey = strings.ToLower(filepath.Clean(m.Path))
}
if people, ok := e.cachedMediaPeople(cacheKey); ok {
return people
}
dir := filepath.Dir(m.Path)
candidates := make([]string, 0, 6)
seenPath := map[string]struct{}{}
@@ -609,9 +618,48 @@ func (e *EmbyService) resolveMediaPeople(ctx context.Context, m *model.Media) []
}
}
}
e.rememberMediaPeople(cacheKey, people)
return people
}
func (e *EmbyService) cachedMediaPeople(key string) ([]map[string]any, bool) {
if e == nil || strings.TrimSpace(key) == "" {
return nil, false
}
now := time.Now()
e.peopleMu.RLock()
entry, ok := e.peopleCache[key]
e.peopleMu.RUnlock()
if !ok || now.After(entry.expiresAt) {
if ok {
e.peopleMu.Lock()
delete(e.peopleCache, key)
e.peopleMu.Unlock()
}
return nil, false
}
out := make([]map[string]any, len(entry.people))
copy(out, entry.people)
return out, true
}
func (e *EmbyService) rememberMediaPeople(key string, people []map[string]any) {
if e == nil || strings.TrimSpace(key) == "" {
return
}
e.peopleMu.Lock()
defer e.peopleMu.Unlock()
if e.peopleCache == nil || len(e.peopleCache) > 8000 {
e.peopleCache = make(map[string]embyPeopleCacheEntry, 128)
}
stored := make([]map[string]any, len(people))
copy(stored, people)
e.peopleCache[key] = embyPeopleCacheEntry{
people: stored,
expiresAt: time.Now().Add(embyVirtualCacheTTL),
}
}
func embyPersonID(name, roleType string) string {
sum := sha256.Sum256([]byte(strings.ToLower(strings.TrimSpace(name)) + ":" + strings.ToLower(strings.TrimSpace(roleType))))
return "person-" + hex.EncodeToString(sum[:8])
+27 -6
View File
@@ -75,9 +75,9 @@ func (e *EmbyService) mediaItems(ctx context.Context, p ItemsParams) (map[string
orderIncludesDirection = false
case "premieredate", "productionyear":
order = mediaReleaseOrderSQL(desc)
case "datecreated", "datelastmediaadded", "datelastcontentadded":
order = "media.created_at"
orderIncludesDirection = false
case "datecreated", "datelastmediaadded", "datelastcontentadded":
order = "media.created_at"
orderIncludesDirection = false
case "dateplayed":
order = "resume.watched_at"
orderIncludesDirection = false
@@ -225,6 +225,17 @@ func (e *EmbyService) collapseMediaVersionRows(ctx context.Context, rows []model
}
func (e *EmbyService) seriesItemsForLibrary(ctx context.Context, libraryID string, p ItemsParams) (map[string]any, error) {
cacheKey := e.embyItemsCacheKey("series-items-v1", p)
var cached embyItemsCacheValue
if e.cache != nil && e.cache.GetJSON(ctx, cacheKey, &cached) {
e.rememberArtworkRefs(cached.Artwork)
return map[string]any{
"Items": cached.Items,
"TotalRecordCount": int(cached.TotalRecordCount),
"StartIndex": cached.StartIndex,
}, nil
}
q := e.repo.DB.WithContext(ctx).Model(&model.Media{}).Where("season_num > 0 OR episode_num > 0")
q = e.applyUserMediaVisibility(ctx, q, p.UserID)
if libraryID != "" {
@@ -246,9 +257,19 @@ func (e *EmbyService) seriesItemsForLibrary(ctx context.Context, libraryID strin
groups := e.seriesGroupsFromMedia(ctx, rows)
sortSeriesGroups(groups, p)
total := len(groups)
items := make([]map[string]any, 0, minInt(p.Limit, len(groups)))
for _, group := range pageSlice(groups, p.StartIndex, p.Limit) {
pageGroups := pageSlice(groups, p.StartIndex, p.Limit)
items := make([]map[string]any, 0, len(pageGroups))
for _, group := range pageGroups {
items = append(items, e.seriesPayload(group))
}
return map[string]any{"Items": items, "TotalRecordCount": total, "StartIndex": p.StartIndex}, nil
out := map[string]any{"Items": items, "TotalRecordCount": total, "StartIndex": p.StartIndex}
if e.cache != nil {
e.cache.SetJSON(ctx, cacheKey, embyItemsCacheValue{
Items: items,
TotalRecordCount: int64(total),
StartIndex: p.StartIndex,
Artwork: e.artworkRefsForSeriesGroups(pageGroups),
}, time.Duration(e.mediaCacheTTLSeconds())*time.Second)
}
return out, nil
}
+64 -14
View File
@@ -35,6 +35,17 @@ func (e *EmbyService) movieLibraryHasEpisodicContent(ctx context.Context, librar
// 与 mediaItems 的区别: 后者会把剧集结构行当散装 Episode 漏出;这里改为聚合成
// Series,从根本上消除「电影库里整部剧被拆成单集」的现象。
func (e *EmbyService) movieLibraryItems(ctx context.Context, p ItemsParams) (map[string]any, error) {
cacheKey := e.embyItemsCacheKey("movie-library-items-v1", p)
var cached embyItemsCacheValue
if e.cache != nil && e.cache.GetJSON(ctx, cacheKey, &cached) {
e.rememberArtworkRefs(cached.Artwork)
return map[string]any{
"Items": cached.Items,
"TotalRecordCount": int(cached.TotalRecordCount),
"StartIndex": cached.StartIndex,
}, nil
}
libIDs := e.mergedLibraryIDs(ctx, p.ParentID)
apply := func(q *gorm.DB) *gorm.DB {
q = e.applyUserMediaVisibility(ctx, q, p.UserID)
@@ -78,33 +89,72 @@ func (e *EmbyService) movieLibraryItems(ctx context.Context, p ItemsParams) (map
if err := movieQ.Find(&movieRows).Error; err != nil {
return nil, err
}
movieItems, err := e.payloadsForMedia(ctx, movieRows, p.UserID)
if err != nil {
return nil, err
}
// 先按版本去重,再参与排序。这里不立即构建 payload:大电影库可能有
// 数万行,而客户端一页通常只要几十条,提前构建会触发大量 NFO / 数据
// 查询并把响应时间浪费在用户根本看不到的条目上。
movieRows = e.collapseMediaVersionRows(ctx, movieRows)
// 合并: Series 卡片 + Movie 项, 统一按首播/上映日期倒序。
type entry struct {
sortAt time.Time
payload map[string]any
sortAt time.Time
media *model.Media
group *embySeriesGroup
}
entries := make([]entry, 0, len(seriesGroups)+len(movieItems))
for _, g := range seriesGroups {
entries = append(entries, entry{sortAt: embySeriesReleaseSortTime(g), payload: e.seriesPayload(g)})
entries := make([]entry, 0, len(seriesGroups)+len(movieRows))
for i := range seriesGroups {
group := &seriesGroups[i]
entries = append(entries, entry{sortAt: embySeriesReleaseSortTime(*group), group: group})
}
for _, item := range movieItems {
entries = append(entries, entry{sortAt: embyPayloadReleaseSortTime(item), payload: item})
for i := range movieRows {
media := &movieRows[i]
entries = append(entries, entry{sortAt: embyMediaReleaseSortTime(*media), media: media})
}
sort.SliceStable(entries, func(i, j int) bool {
return entries[i].sortAt.After(entries[j].sortAt)
})
total := len(entries)
paged := pageSlice(entries, p.StartIndex, p.Limit)
items := make([]map[string]any, 0, len(paged))
pageMovies := make([]model.Media, 0, len(paged))
for _, en := range paged {
items = append(items, en.payload)
if en.media != nil {
pageMovies = append(pageMovies, *en.media)
}
}
return map[string]any{"Items": items, "TotalRecordCount": total, "StartIndex": p.StartIndex}, nil
moviePayloads, err := e.payloadsForMedia(ctx, pageMovies, p.UserID)
if err != nil {
return nil, err
}
payloadByID := make(map[string]map[string]any, len(moviePayloads))
for _, item := range moviePayloads {
if id, ok := item["Id"].(string); ok {
payloadByID[id] = item
}
}
items := make([]map[string]any, 0, len(paged))
pageGroups := make([]embySeriesGroup, 0, len(paged))
for _, en := range paged {
switch {
case en.group != nil:
pageGroups = append(pageGroups, *en.group)
items = append(items, e.seriesPayload(*en.group))
case en.media != nil:
if item := payloadByID[en.media.ID]; item != nil {
items = append(items, item)
}
}
}
out := map[string]any{"Items": items, "TotalRecordCount": total, "StartIndex": p.StartIndex}
if e.cache != nil {
e.cache.SetJSON(ctx, cacheKey, embyItemsCacheValue{
Items: items,
TotalRecordCount: int64(total),
StartIndex: p.StartIndex,
Artwork: e.artworkRefsForSeriesGroups(pageGroups),
}, time.Duration(e.mediaCacheTTLSeconds())*time.Second)
}
return out, nil
}
// embyPayloadCreatedAt 从 item payload 里取 DateCreated(time.Time),用于合并排序。
+123 -5
View File
@@ -1,6 +1,7 @@
package service
import (
"sort"
"strings"
"time"
)
@@ -37,11 +38,7 @@ func (e *EmbyService) rememberSeriesGroup(group embySeriesGroup) {
if e.virtualArtwork == nil {
e.virtualArtwork = make(map[string]embyArtworkCacheEntry)
}
if len(e.virtualSeries) > 2000 || len(e.virtualSeasons) > 5000 || len(e.virtualArtwork) > 7000 {
e.virtualSeries = make(map[string]embySeriesCacheEntry)
e.virtualSeasons = make(map[string]embySeasonCacheEntry)
e.virtualArtwork = make(map[string]embyArtworkCacheEntry)
}
e.trimVirtualCachesLocked(time.Now())
e.virtualSeries[group.ID] = embySeriesCacheEntry{group: group, expiresAt: expiresAt}
e.virtualArtwork[group.ID] = embyArtworkCacheEntry{primary: group.PosterURL, backdrop: group.BackdropURL, expiresAt: expiresAt}
e.virtualArtwork[group.ID+"-bd"] = embyArtworkCacheEntry{primary: group.PosterURL, backdrop: group.BackdropURL, expiresAt: expiresAt}
@@ -65,6 +62,7 @@ func (e *EmbyService) rememberSeasonGroup(season embySeasonGroup) {
if e.virtualArtwork == nil {
e.virtualArtwork = make(map[string]embyArtworkCacheEntry)
}
e.trimVirtualCachesLocked(time.Now())
e.virtualSeasons[season.ID] = embySeasonCacheEntry{season: season, expiresAt: expiresAt}
e.virtualArtwork[season.ID] = embyArtworkCacheEntry{primary: season.Series.PosterURL, backdrop: season.Series.BackdropURL, expiresAt: expiresAt}
e.virtualArtwork[season.ID+"-bd"] = embyArtworkCacheEntry{primary: season.Series.PosterURL, backdrop: season.Series.BackdropURL, expiresAt: expiresAt}
@@ -135,3 +133,123 @@ func (e *EmbyService) cachedArtworkURL(id, imageType string) (string, bool) {
}
return entry.backdrop, entry.backdrop != ""
}
// embyArtworkRef is the poster/backdrop pair stored next to a Latest payload
// so a JSON-cache hit can refill the in-memory artwork map before the client
// asks for the image.
type embyArtworkRef struct {
Primary string `json:"primary,omitempty"`
Backdrop string `json:"backdrop,omitempty"`
}
func (e *EmbyService) artworkRefsForSeriesGroups(groups []embySeriesGroup) map[string]embyArtworkRef {
if e == nil || len(groups) == 0 {
return nil
}
refs := make(map[string]embyArtworkRef, len(groups)*2)
for _, group := range groups {
if group.PosterURL == "" && group.BackdropURL == "" {
continue
}
ref := embyArtworkRef{Primary: group.PosterURL, Backdrop: group.BackdropURL}
refs[group.ID] = ref
for _, season := range e.seasonsForSeries(group) {
refs[season.ID] = ref
}
}
if len(refs) == 0 {
return nil
}
return refs
}
func (e *EmbyService) rememberArtworkRefs(refs map[string]embyArtworkRef) {
if e == nil || len(refs) == 0 {
return
}
expiresAt := time.Now().Add(embyVirtualCacheTTL)
e.virtualMu.Lock()
defer e.virtualMu.Unlock()
if e.virtualArtwork == nil {
e.virtualArtwork = make(map[string]embyArtworkCacheEntry, len(refs))
}
e.trimVirtualArtworkLocked(time.Now())
for id, ref := range refs {
if strings.TrimSpace(id) == "" || (ref.Primary == "" && ref.Backdrop == "") {
continue
}
e.virtualArtwork[id] = embyArtworkCacheEntry{primary: ref.Primary, backdrop: ref.Backdrop, expiresAt: expiresAt}
}
}
func (e *EmbyService) trimVirtualCachesLocked(now time.Time) {
e.trimVirtualSeriesLocked(now)
e.trimVirtualSeasonsLocked(now)
e.trimVirtualArtworkLocked(now)
}
func (e *EmbyService) trimVirtualSeriesLocked(now time.Time) {
for id, entry := range e.virtualSeries {
if now.After(entry.expiresAt) {
delete(e.virtualSeries, id)
}
}
evictOldest(e.virtualSeries, embyVirtualSeriesCap, func(entry embySeriesCacheEntry) time.Time {
return entry.expiresAt
})
}
func (e *EmbyService) trimVirtualSeasonsLocked(now time.Time) {
for id, entry := range e.virtualSeasons {
if now.After(entry.expiresAt) {
delete(e.virtualSeasons, id)
}
}
evictOldest(e.virtualSeasons, embyVirtualSeasonCap, func(entry embySeasonCacheEntry) time.Time {
return entry.expiresAt
})
}
func (e *EmbyService) trimVirtualArtworkLocked(now time.Time) {
for id, entry := range e.virtualArtwork {
if now.After(entry.expiresAt) {
delete(e.virtualArtwork, id)
}
}
evictOldest(e.virtualArtwork, embyVirtualArtworkCap, func(entry embyArtworkCacheEntry) time.Time {
return entry.expiresAt
})
}
// evictOldest drops the soonest-expiring entries until the map is under cap.
// It must not replace the map: a homepage refresh remembers many series at
// once, and wiping the whole cache made the just-advertised backdrops 404
// into a 1x1 placeholder.
func evictOldest[T any](items map[string]T, cap int, expiresAt func(T) time.Time) {
excess := len(items) - cap
if cap <= 0 || excess <= 0 {
return
}
type pair struct {
id string
at time.Time
}
ordered := make([]pair, 0, len(items))
for id, entry := range items {
ordered = append(ordered, pair{id: id, at: expiresAt(entry)})
}
sort.Slice(ordered, func(i, j int) bool { return ordered[i].at.Before(ordered[j].at) })
if excess > len(ordered) {
excess = len(ordered)
}
for i := 0; i < excess; i++ {
delete(items, ordered[i].id)
}
}
func embyVirtualImageTag(id, suffix string) string {
if strings.HasPrefix(id, embyVirtualSeriesPrefix) || strings.HasPrefix(id, embyVirtualSeasonPrefix) {
return id + suffix
}
return id
}
+113 -6
View File
@@ -296,6 +296,71 @@ func TestEmbyVirtualSeriesArtworkUsesListCache(t *testing.T) {
}
}
func TestEmbyVirtualSeriesArtworkRebuildsAfterMemoryDrop(t *testing.T) {
svc := newTestEmbyService(t)
svc.cache = NewRuntimeCacheService(nil, nil)
lib := model.Library{Name: "番剧", Path: `/media/anime`, Type: "anime", Enabled: true}
if err := svc.repo.Library.Create(t.Context(), &lib); err != nil {
t.Fatalf("create library: %v", err)
}
media := model.Media{
Base: model.Base{ID: "ep-hero"},
LibraryID: lib.ID,
Title: "树海之魔",
Path: `/media/anime/树海之魔/Season 01/树海之魔 - S01E01.mkv`,
PosterURL: `/poster.jpg`,
BackdropURL: `/backdrop.jpg`,
SeasonNum: 1,
EpisodeNum: 1,
}
if err := svc.repo.DB.Create(&media).Error; err != nil {
t.Fatalf("create media: %v", err)
}
items, err := svc.LatestItems(t.Context(), "", lib.ID, 5)
if err != nil {
t.Fatalf("latest items: %v", err)
}
if len(items) != 1 {
t.Fatalf("latest len = %d, want 1", len(items))
}
seriesID, _ := items[0]["Id"].(string)
tags, _ := items[0]["BackdropImageTags"].([]string)
if seriesID == "" || len(tags) != 1 || tags[0] != seriesID+embyVirtualBackdropTagSuffix {
t.Fatalf("hero item should advertise a cache-busted backdrop tag, got id=%q tags=%#v", seriesID, items[0]["BackdropImageTags"])
}
svc.virtualMu.Lock()
svc.virtualArtwork = nil
svc.virtualSeries = nil
svc.virtualSeasons = nil
svc.virtualMu.Unlock()
backdrop, err := svc.ImageURL(t.Context(), seriesID, "Backdrop")
if err != nil {
t.Fatalf("backdrop after memory drop: %v", err)
}
if backdrop != "/backdrop.jpg" {
t.Fatalf("backdrop = %q, want rebuilt backdrop", backdrop)
}
svc.virtualMu.Lock()
svc.virtualArtwork = nil
svc.virtualMu.Unlock()
if _, err := svc.LatestItems(t.Context(), "", lib.ID, 5); err != nil {
t.Fatalf("cached latest: %v", err)
}
cancelled, cancel := context.WithCancel(t.Context())
cancel()
backdrop, err = svc.ImageURL(cancelled, seriesID, "Backdrop")
if err != nil {
t.Fatalf("backdrop from rewarmed cache: %v", err)
}
if backdrop != "/backdrop.jpg" {
t.Fatalf("rewarmed backdrop = %q, want cached backdrop", backdrop)
}
}
func TestEmbyCloudAnimeUsesSeriesNameFromChineseSeasonFolder(t *testing.T) {
svc := newTestEmbyService(t)
lib := model.Library{Name: "OpenList · 国漫", Path: `cloud://openlist/国漫`, Type: "anime", Enabled: true}
@@ -429,13 +494,13 @@ func TestInferSeriesNameFromPath(t *testing.T) {
want: "间谍过家家",
},
}
for _, tc := range tests {
got := inferSeriesNameFromPath(tc.path)
if got != tc.want {
t.Errorf("inferSeriesNameFromPath(%q) = %q, want %q", tc.path, got, tc.want)
}
}
for _, tc := range tests {
got := inferSeriesNameFromPath(tc.path)
if got != tc.want {
t.Errorf("inferSeriesNameFromPath(%q) = %q, want %q", tc.path, got, tc.want)
}
}
}
func TestEmbySeriesSortByDateLastMediaAdded(t *testing.T) {
svc := newTestEmbyService(t)
@@ -513,3 +578,45 @@ func TestEmbySeriesSortByDateLastMediaAdded(t *testing.T) {
t.Fatalf("DateLastMediaAdded = %v, want %v", items[0]["DateLastMediaAdded"], tNew)
}
}
func TestEmbySeriesLibraryListUsesRuntimeCache(t *testing.T) {
svc := newTestEmbyService(t)
svc.cache = NewRuntimeCacheService(nil, nil)
lib := model.Library{Name: "番剧", Path: `/media/anime`, Type: "anime", Enabled: true}
if err := svc.repo.Library.Create(t.Context(), &lib); err != nil {
t.Fatalf("create library: %v", err)
}
for i := 1; i <= 2; i++ {
media := model.Media{
Base: model.Base{ID: fmt.Sprintf("cache-ep-%d", i)},
LibraryID: lib.ID,
Title: "缓存测试番",
Path: fmt.Sprintf(`/media/anime/缓存测试番/Season 01/缓存测试番.S01E%02d.mkv`, i),
SeasonNum: 1,
EpisodeNum: i,
}
if err := svc.repo.DB.Create(&media).Error; err != nil {
t.Fatalf("create media: %v", err)
}
}
first, err := svc.Items(t.Context(), ItemsParams{ParentID: lib.ID, Limit: 20})
if err != nil {
t.Fatalf("first items call: %v", err)
}
if first["TotalRecordCount"] != 1 {
t.Fatalf("first series total = %#v, want 1", first["TotalRecordCount"])
}
if err := svc.repo.DB.Unscoped().Where("library_id = ?", lib.ID).Delete(&model.Media{}).Error; err != nil {
t.Fatalf("delete media: %v", err)
}
second, err := svc.Items(t.Context(), ItemsParams{ParentID: lib.ID, Limit: 20})
if err != nil {
t.Fatalf("second items call: %v", err)
}
items, _ := second["Items"].([]map[string]any)
if second["TotalRecordCount"] != 1 || len(items) != 1 {
t.Fatalf("cached series list = %#v, want the first response", second)
}
}
+28 -28
View File
@@ -5,28 +5,28 @@ func (e *EmbyService) seriesPayload(group embySeriesGroup) map[string]any {
imageTags := map[string]string{}
backdropTags := []string{}
if group.PosterURL != "" {
imageTags["Primary"] = group.ID
imageTags["Primary"] = embyVirtualImageTag(group.ID, embyVirtualPrimaryTagSuffix)
}
if group.BackdropURL != "" {
backdropTags = append(backdropTags, group.ID+"-bd")
if group.BackdropURL != "" || group.PosterURL != "" {
backdropTags = append(backdropTags, embyVirtualImageTag(group.ID, embyVirtualBackdropTagSuffix))
}
lastMediaAdded := group.DateLastMediaAdded
if lastMediaAdded.IsZero() {
lastMediaAdded = group.CreatedAt
}
item := map[string]any{
"Id": group.ID,
"Name": group.Name,
"ServerId": embyServerID,
"Type": "Series",
"MediaType": "Video",
"IsFolder": true,
"ParentId": group.LibraryID,
"ProductionYear": group.Year,
"Overview": group.Overview,
"CommunityRating": group.Rating,
"RecursiveItemCount": len(group.Episodes),
"ChildCount": len(e.seasonsForSeries(group)),
lastMediaAdded := group.DateLastMediaAdded
if lastMediaAdded.IsZero() {
lastMediaAdded = group.CreatedAt
}
item := map[string]any{
"Id": group.ID,
"Name": group.Name,
"ServerId": embyServerID,
"Type": "Series",
"MediaType": "Video",
"IsFolder": true,
"ParentId": group.LibraryID,
"ProductionYear": group.Year,
"Overview": group.Overview,
"CommunityRating": group.Rating,
"RecursiveItemCount": len(group.Episodes),
"ChildCount": len(e.seasonsForSeries(group)),
"DateCreated": group.CreatedAt,
"DateLastMediaAdded": lastMediaAdded,
"ImageTags": imageTags,
@@ -39,7 +39,7 @@ func (e *EmbyService) seriesPayload(group embySeriesGroup) map[string]any {
"UserData": emptyUserData(),
}
if group.PosterURL != "" {
item["PrimaryImageTag"] = group.ID
item["PrimaryImageTag"] = embyVirtualImageTag(group.ID, embyVirtualPrimaryTagSuffix)
}
if premiered, ok := embyPremiereDate(group.ReleaseDate); ok {
item["PremiereDate"] = premiered
@@ -52,10 +52,10 @@ func (e *EmbyService) seasonPayload(season embySeasonGroup) map[string]any {
imageTags := map[string]string{}
backdropTags := []string{}
if season.Series.PosterURL != "" {
imageTags["Primary"] = season.ID
imageTags["Primary"] = embyVirtualImageTag(season.ID, embyVirtualPrimaryTagSuffix)
}
if season.Series.BackdropURL != "" {
backdropTags = append(backdropTags, season.ID+"-bd")
if season.Series.BackdropURL != "" || season.Series.PosterURL != "" {
backdropTags = append(backdropTags, embyVirtualImageTag(season.ID, embyVirtualBackdropTagSuffix))
}
item := map[string]any{
"Id": season.ID,
@@ -75,12 +75,12 @@ func (e *EmbyService) seasonPayload(season embySeasonGroup) map[string]any {
"UserData": emptyUserData(),
}
if season.Series.PosterURL != "" {
item["PrimaryImageTag"] = season.ID
item["SeriesPrimaryImageTag"] = season.Series.ID
item["PrimaryImageTag"] = embyVirtualImageTag(season.ID, embyVirtualPrimaryTagSuffix)
item["SeriesPrimaryImageTag"] = embyVirtualImageTag(season.Series.ID, embyVirtualPrimaryTagSuffix)
}
if season.Series.BackdropURL != "" {
if season.Series.BackdropURL != "" || season.Series.PosterURL != "" {
item["ParentBackdropItemId"] = season.Series.ID
item["ParentBackdropImageTags"] = []string{season.Series.ID + "-bd"}
item["ParentBackdropImageTags"] = []string{embyVirtualImageTag(season.Series.ID, embyVirtualBackdropTagSuffix)}
}
return item
}
+23
View File
@@ -17,6 +17,7 @@ import (
"net"
"net/http"
"path/filepath"
"runtime"
"strings"
"sync"
"syscall"
@@ -35,6 +36,12 @@ type ImageProxy struct {
cacheDir string
mu sync.Mutex
// resizeSem bounds concurrent decode/resize jobs. Emby TV clients request
// poster grids in bursts; letting every request decode a source image at
// once causes CPU and memory spikes that make the whole UI feel sluggish.
resizeSemMu sync.Mutex
resizeSem chan struct{}
// libraryRootsFn returns the configured media library roots so that
// sidecar poster/artwork files stored alongside media (under arbitrary
// per-library paths) are allowed by isAllowedLocalPath. It is provided
@@ -55,6 +62,7 @@ type ImageProxy struct {
const (
imageBrowserCacheControl = "public, max-age=2592000, immutable"
imagePlaceholderCacheControl = "no-store"
imageMaxResizeConcurrency = 4
)
// NewImageProxy is the constructor.
@@ -64,6 +72,7 @@ func NewImageProxy(cfg *config.Config, log *zap.Logger) *ImageProxy {
log: log,
cacheDir: filepath.Join(cfg.Cache.CacheDir, "images"),
}
proxy.resizeSem = make(chan struct{}, imageResizeConcurrency())
// Honor HTTP(S)_PROXY env vars so deployments behind GFW can pull
// from image.tmdb.org via their HTTP proxy without extra config. On
@@ -181,6 +190,20 @@ func (p *ImageProxy) isAllowedRemoteHost(host string) bool {
return p.allowedHostsCache[host]
}
// imageResizeConcurrency keeps decode/resize concurrency within the number
// of CPU threads the process is allowed to use, capped to avoid large
// temporary RGBA buffers on tiny hosts.
func imageResizeConcurrency() int {
n := runtime.GOMAXPROCS(0)
if n < 1 {
n = 1
}
if n > imageMaxResizeConcurrency {
n = imageMaxResizeConcurrency
}
return n
}
// Prune removes oldest cached images until disk usage is within the configured limit.
func (p *ImageProxy) Prune() (PruneImageCacheResult, error) {
if p.cfg == nil || p.cfg.Cache.ImagesMaxSizeMB <= 0 {
+30 -3
View File
@@ -4,6 +4,7 @@ import (
"bytes"
"context"
"errors"
"io"
"net/http"
"os"
"path/filepath"
@@ -134,12 +135,38 @@ func (p *ImageProxy) serveCachedImage(w http.ResponseWriter, r *http.Request, ke
}
func (p *ImageProxy) removeUnusableImageCache(cachePath, failPath string) {
data, err := os.ReadFile(cachePath) // #nosec G304 -- cachePath is SHA-derived under cacheDir.
// 只读取文件头判断缓存是否可用。旧实现每次命中远程图片缓存都会把整个
// 原图读进内存再丢弃,电视端批量加载海报时会产生大量无意义的磁盘 I/O。
file, err := os.Open(cachePath) // #nosec G304 -- cachePath is SHA-derived under cacheDir.
if err != nil {
return
}
ctype := detectContentType(data)
if len(data) > 0 && isImageContentType(ctype) && !isTransparentPlaceholderData(data) {
stat, err := file.Stat()
if err != nil || stat.IsDir() || stat.Size() <= 0 {
_ = file.Close()
_ = os.Remove(cachePath)
_ = os.Remove(failPath)
return
}
headerSize := 512
if stat.Size() < int64(headerSize) {
headerSize = int(stat.Size())
}
header := make([]byte, headerSize)
n, readErr := io.ReadFull(file, header)
_ = file.Close()
if readErr != nil && readErr != io.ErrUnexpectedEOF {
_ = os.Remove(cachePath)
_ = os.Remove(failPath)
return
}
header = header[:n]
// A transparent placeholder is exactly 67 bytes; checking the header alone
// is enough for the normal image cache entries (they are much larger but
// detectContentType only inspects the same leading 512 bytes anyway).
// Close the handle before deleting: Windows refuses to delete an open file.
if n > 0 && isImageContentType(detectContentType(header)) &&
!(n == len(transparent1x1PNG) && bytes.Equal(header, transparent1x1PNG)) {
return
}
_ = os.Remove(cachePath)
+42 -5
View File
@@ -2,6 +2,7 @@ package service
import (
"bytes"
"context"
"crypto/sha256"
"encoding/hex"
"errors"
@@ -203,6 +204,28 @@ func (p *ImageProxy) resizeCachePath(key string) string {
return filepath.Join(p.cacheDir, imageResizeCacheSubdir, key+".img")
}
// acquireResizeSlot bounds CPU-heavy decode/resize work. Returning false means
// the caller should fall back to the original image instead of blocking after
// the request has already been canceled.
func (p *ImageProxy) acquireResizeSlot(ctx context.Context) (func(), bool) {
if p == nil {
return func() {}, true
}
p.resizeSemMu.Lock()
if p.resizeSem == nil {
p.resizeSem = make(chan struct{}, imageResizeConcurrency())
}
sem := p.resizeSem
p.resizeSemMu.Unlock()
select {
case sem <- struct{}{}:
return func() { <-sem }, true
case <-ctx.Done():
return nil, false
}
}
// serveResizedFromFile 从 srcPath 读取图片,按选项缩放后写出,并把结果缓存
// 到磁盘以免每次请求都重新解码。原图已满足目标尺寸时直接输出原文件。
// 返回 false 表示缩放不可用,调用方应回退到原图直出。
@@ -211,6 +234,25 @@ func (p *ImageProxy) serveResizedFromFile(w http.ResponseWriter, r *http.Request
if err != nil || stat.IsDir() || stat.Size() <= 0 {
return false
}
// 缓存命中必须发生在读原图和解码之前。否则电视端每次刷新海报墙都会
// 把已经是缩略图缓存的原图重新解码、缩放一遍,造成明显的 CPU 抖动。
key := o.resizeCacheKey(srcPath, stat)
cachePath := p.resizeCachePath(key)
if serveCachedImageFile(w, r, key, cachePath) {
return true
}
release, ok := p.acquireResizeSlot(r.Context())
if !ok {
return false
}
defer release()
// 等待并发槽期间,别的请求可能已经生成了同一张缩略图。
if serveCachedImageFile(w, r, key, cachePath) {
return true
}
data, err := os.ReadFile(srcPath) // #nosec G304 -- srcPath comes from an allowed local path or a SHA-derived cache path.
if err != nil {
return false
@@ -224,11 +266,6 @@ func (p *ImageProxy) serveResizedFromFile(w http.ResponseWriter, r *http.Request
return serveImageFile(w, r, filepath.Base(srcPath), srcPath, imageBrowserCacheControl)
}
key := o.resizeCacheKey(srcPath, stat)
cachePath := p.resizeCachePath(key)
if serveCachedImageFile(w, r, key, cachePath) {
return true
}
p.writeResizeCache(cachePath, out)
w.Header().Set("Content-Type", ctype)
w.Header().Set("Cache-Control", imageBrowserCacheControl)
+44
View File
@@ -230,3 +230,47 @@ func TestServeResizedFromFileCachesScaledResult(t *testing.T) {
t.Fatal("expected a non-empty body")
}
}
func TestServeResizedFromFileUsesCacheBeforeDecodingSource(t *testing.T) {
dir := t.TempDir()
mediaDir := dir + string(os.PathSeparator) + "media"
if err := os.MkdirAll(mediaDir, 0o755); err != nil {
t.Fatalf("mkdir: %v", err)
}
src := mediaDir + string(os.PathSeparator) + "poster.png"
original := encodeTestPNG(t, 529, 911, 255)
if err := os.WriteFile(src, original, 0o644); err != nil {
t.Fatalf("write source: %v", err)
}
proxy := &ImageProxy{cacheDir: dir + string(os.PathSeparator) + "cache"}
opts := imageResizeOptions{MaxWidth: 400, Quality: 90}
first := httptest.NewRecorder()
if !proxy.serveResizedFromFile(first, httptest.NewRequest("GET", "/x?maxWidth=400", nil), src, opts) {
t.Fatal("expected first call to be served")
}
stat, err := os.Stat(src)
if err != nil {
t.Fatalf("stat source: %v", err)
}
// Same size and mtime keep the resize cache key stable, but the source is
// now invalid image data. A correct implementation serves the cached
// thumbnail before reading/decoding the source again.
broken := bytes.Repeat([]byte{0}, len(original))
if err := os.WriteFile(src, broken, 0o644); err != nil {
t.Fatalf("overwrite source: %v", err)
}
if err := os.Chtimes(src, stat.ModTime(), stat.ModTime()); err != nil {
t.Fatalf("restore mtime: %v", err)
}
second := httptest.NewRecorder()
if !proxy.serveResizedFromFile(second, httptest.NewRequest("GET", "/x?maxWidth=400", nil), src, opts) {
t.Fatal("expected second call to be served from resize cache")
}
if !bytes.Equal(first.Body.Bytes(), second.Body.Bytes()) {
t.Fatal("expected cached thumbnail to be reused without decoding the source")
}
}
+12
View File
@@ -32,6 +32,18 @@ func mediaReleaseOrderSQL(desc bool) string {
return fmt.Sprintf("media.release_date %s, media.year %s, media.created_at %s, media.id %s", dir, dir, dir, dir)
}
// embyMediaReleaseSortTime matches the original payload-based ordering used by
// the Emby movie-library merge: release date, then year, then created_at.
func embyMediaReleaseSortTime(media model.Media) time.Time {
if t, ok := embyPremiereDate(media.ReleaseDate); ok {
return t
}
if media.Year > 0 {
return time.Date(media.Year, time.December, 31, 0, 0, 0, 0, time.UTC)
}
return media.CreatedAt
}
func embyPremiereDate(value string) (time.Time, bool) {
value = normalizeReleaseDate(value)
if value == "" {