mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-08 06:16:37 +08:00
Optimize Emby caching and image resize concurrency
This commit is contained in:
@@ -61,6 +61,9 @@ type EmbyService struct {
|
|||||||
|
|
||||||
libraryCoverMu sync.Mutex
|
libraryCoverMu sync.Mutex
|
||||||
libraryCoverCache map[string]embyArtworkCacheEntry
|
libraryCoverCache map[string]embyArtworkCacheEntry
|
||||||
|
|
||||||
|
peopleMu sync.RWMutex
|
||||||
|
peopleCache map[string]embyPeopleCacheEntry
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewEmbyService is the constructor.
|
// NewEmbyService is the constructor.
|
||||||
@@ -143,6 +146,13 @@ type embyVisibilityCacheEntry struct {
|
|||||||
expiresAt time.Time
|
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
|
// Items paginates media in Emby's hierarchy. Episodic libraries are exposed as
|
||||||
// Series -> Season -> Episode so Infuse/Vidhub/SenPlayer stop treating every
|
// Series -> Season -> Episode so Infuse/Vidhub/SenPlayer stop treating every
|
||||||
// episode as a separate movie card. 带 embyremote~ 前缀的 ParentID / 搜索自动
|
// episode as a separate movie card. 带 embyremote~ 前缀的 ParentID / 搜索自动
|
||||||
|
|||||||
@@ -3,12 +3,24 @@ package service
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"strings"
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/truewhile/MeBox/internal/model"
|
"github.com/truewhile/MeBox/internal/model"
|
||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
)
|
)
|
||||||
|
|
||||||
func (e *EmbyService) ItemCounts(ctx context.Context, userID string) (map[string]any, error) {
|
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 {
|
base := func() *gorm.DB {
|
||||||
q := e.repo.DB.WithContext(ctx).Model(&model.Media{}).Where("deleted_at IS NULL")
|
q := e.repo.DB.WithContext(ctx).Model(&model.Media{}).Where("deleted_at IS NULL")
|
||||||
return e.applyUserMediaVisibility(ctx, q, userID)
|
return e.applyUserMediaVisibility(ctx, q, userID)
|
||||||
@@ -34,6 +46,14 @@ func (e *EmbyService) ItemCounts(ctx context.Context, userID string) (map[string
|
|||||||
return nil, err
|
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{
|
return map[string]any{
|
||||||
"MovieCount": movieCount,
|
"MovieCount": movieCount,
|
||||||
"SeriesCount": seriesCount,
|
"SeriesCount": seriesCount,
|
||||||
|
|||||||
@@ -9,9 +9,10 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type embyItemsCacheValue struct {
|
type embyItemsCacheValue struct {
|
||||||
Items []map[string]any `json:"items"`
|
Items []map[string]any `json:"items"`
|
||||||
TotalRecordCount int64 `json:"total_record_count"`
|
TotalRecordCount int64 `json:"total_record_count"`
|
||||||
StartIndex int `json:"start_index"`
|
StartIndex int `json:"start_index"`
|
||||||
|
Artwork map[string]embyArtworkRef `json:"artwork,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type embyLatestCacheValue struct {
|
type embyLatestCacheValue struct {
|
||||||
@@ -19,6 +20,13 @@ type embyLatestCacheValue struct {
|
|||||||
Artwork map[string]embyArtworkRef `json:"artwork,omitempty"`
|
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 {
|
func (e *EmbyService) embyItemsCacheKey(kind string, p ItemsParams) string {
|
||||||
includeTypes := append([]string(nil), p.IncludeItemTypes...)
|
includeTypes := append([]string(nil), p.IncludeItemTypes...)
|
||||||
filters := append([]string(nil), p.Filters...)
|
filters := append([]string(nil), p.Filters...)
|
||||||
|
|||||||
@@ -538,6 +538,14 @@ func (e *EmbyService) resolveMediaPeople(ctx context.Context, m *model.Media) []
|
|||||||
if m == nil || strings.TrimSpace(m.Path) == "" {
|
if m == nil || strings.TrimSpace(m.Path) == "" {
|
||||||
return []map[string]any{}
|
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)
|
dir := filepath.Dir(m.Path)
|
||||||
candidates := make([]string, 0, 6)
|
candidates := make([]string, 0, 6)
|
||||||
seenPath := map[string]struct{}{}
|
seenPath := map[string]struct{}{}
|
||||||
@@ -610,9 +618,48 @@ func (e *EmbyService) resolveMediaPeople(ctx context.Context, m *model.Media) []
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
e.rememberMediaPeople(cacheKey, people)
|
||||||
return 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 {
|
func embyPersonID(name, roleType string) string {
|
||||||
sum := sha256.Sum256([]byte(strings.ToLower(strings.TrimSpace(name)) + ":" + strings.ToLower(strings.TrimSpace(roleType))))
|
sum := sha256.Sum256([]byte(strings.ToLower(strings.TrimSpace(name)) + ":" + strings.ToLower(strings.TrimSpace(roleType))))
|
||||||
return "person-" + hex.EncodeToString(sum[:8])
|
return "person-" + hex.EncodeToString(sum[:8])
|
||||||
|
|||||||
@@ -75,9 +75,9 @@ func (e *EmbyService) mediaItems(ctx context.Context, p ItemsParams) (map[string
|
|||||||
orderIncludesDirection = false
|
orderIncludesDirection = false
|
||||||
case "premieredate", "productionyear":
|
case "premieredate", "productionyear":
|
||||||
order = mediaReleaseOrderSQL(desc)
|
order = mediaReleaseOrderSQL(desc)
|
||||||
case "datecreated", "datelastmediaadded", "datelastcontentadded":
|
case "datecreated", "datelastmediaadded", "datelastcontentadded":
|
||||||
order = "media.created_at"
|
order = "media.created_at"
|
||||||
orderIncludesDirection = false
|
orderIncludesDirection = false
|
||||||
case "dateplayed":
|
case "dateplayed":
|
||||||
order = "resume.watched_at"
|
order = "resume.watched_at"
|
||||||
orderIncludesDirection = false
|
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) {
|
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.repo.DB.WithContext(ctx).Model(&model.Media{}).Where("season_num > 0 OR episode_num > 0")
|
||||||
q = e.applyUserMediaVisibility(ctx, q, p.UserID)
|
q = e.applyUserMediaVisibility(ctx, q, p.UserID)
|
||||||
if libraryID != "" {
|
if libraryID != "" {
|
||||||
@@ -246,9 +257,19 @@ func (e *EmbyService) seriesItemsForLibrary(ctx context.Context, libraryID strin
|
|||||||
groups := e.seriesGroupsFromMedia(ctx, rows)
|
groups := e.seriesGroupsFromMedia(ctx, rows)
|
||||||
sortSeriesGroups(groups, p)
|
sortSeriesGroups(groups, p)
|
||||||
total := len(groups)
|
total := len(groups)
|
||||||
items := make([]map[string]any, 0, minInt(p.Limit, len(groups)))
|
pageGroups := pageSlice(groups, p.StartIndex, p.Limit)
|
||||||
for _, group := range pageSlice(groups, p.StartIndex, p.Limit) {
|
items := make([]map[string]any, 0, len(pageGroups))
|
||||||
|
for _, group := range pageGroups {
|
||||||
items = append(items, e.seriesPayload(group))
|
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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -35,6 +35,17 @@ func (e *EmbyService) movieLibraryHasEpisodicContent(ctx context.Context, librar
|
|||||||
// 与 mediaItems 的区别: 后者会把剧集结构行当散装 Episode 漏出;这里改为聚合成
|
// 与 mediaItems 的区别: 后者会把剧集结构行当散装 Episode 漏出;这里改为聚合成
|
||||||
// Series,从根本上消除「电影库里整部剧被拆成单集」的现象。
|
// Series,从根本上消除「电影库里整部剧被拆成单集」的现象。
|
||||||
func (e *EmbyService) movieLibraryItems(ctx context.Context, p ItemsParams) (map[string]any, error) {
|
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)
|
libIDs := e.mergedLibraryIDs(ctx, p.ParentID)
|
||||||
apply := func(q *gorm.DB) *gorm.DB {
|
apply := func(q *gorm.DB) *gorm.DB {
|
||||||
q = e.applyUserMediaVisibility(ctx, q, p.UserID)
|
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 {
|
if err := movieQ.Find(&movieRows).Error; err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
movieItems, err := e.payloadsForMedia(ctx, movieRows, p.UserID)
|
// 先按版本去重,再参与排序。这里不立即构建 payload:大电影库可能有
|
||||||
if err != nil {
|
// 数万行,而客户端一页通常只要几十条,提前构建会触发大量 NFO / 数据
|
||||||
return nil, err
|
// 查询并把响应时间浪费在用户根本看不到的条目上。
|
||||||
}
|
movieRows = e.collapseMediaVersionRows(ctx, movieRows)
|
||||||
|
|
||||||
// 合并: Series 卡片 + Movie 项, 统一按首播/上映日期倒序。
|
// 合并: Series 卡片 + Movie 项, 统一按首播/上映日期倒序。
|
||||||
type entry struct {
|
type entry struct {
|
||||||
sortAt time.Time
|
sortAt time.Time
|
||||||
payload map[string]any
|
media *model.Media
|
||||||
|
group *embySeriesGroup
|
||||||
}
|
}
|
||||||
entries := make([]entry, 0, len(seriesGroups)+len(movieItems))
|
entries := make([]entry, 0, len(seriesGroups)+len(movieRows))
|
||||||
for _, g := range seriesGroups {
|
for i := range seriesGroups {
|
||||||
entries = append(entries, entry{sortAt: embySeriesReleaseSortTime(g), payload: e.seriesPayload(g)})
|
group := &seriesGroups[i]
|
||||||
|
entries = append(entries, entry{sortAt: embySeriesReleaseSortTime(*group), group: group})
|
||||||
}
|
}
|
||||||
for _, item := range movieItems {
|
for i := range movieRows {
|
||||||
entries = append(entries, entry{sortAt: embyPayloadReleaseSortTime(item), payload: item})
|
media := &movieRows[i]
|
||||||
|
entries = append(entries, entry{sortAt: embyMediaReleaseSortTime(*media), media: media})
|
||||||
}
|
}
|
||||||
sort.SliceStable(entries, func(i, j int) bool {
|
sort.SliceStable(entries, func(i, j int) bool {
|
||||||
return entries[i].sortAt.After(entries[j].sortAt)
|
return entries[i].sortAt.After(entries[j].sortAt)
|
||||||
})
|
})
|
||||||
total := len(entries)
|
total := len(entries)
|
||||||
paged := pageSlice(entries, p.StartIndex, p.Limit)
|
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 {
|
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),用于合并排序。
|
// embyPayloadCreatedAt 从 item payload 里取 DateCreated(time.Time),用于合并排序。
|
||||||
|
|||||||
@@ -578,3 +578,45 @@ func TestEmbySeriesSortByDateLastMediaAdded(t *testing.T) {
|
|||||||
t.Fatalf("DateLastMediaAdded = %v, want %v", items[0]["DateLastMediaAdded"], tNew)
|
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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -17,6 +17,7 @@ import (
|
|||||||
"net"
|
"net"
|
||||||
"net/http"
|
"net/http"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"runtime"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"syscall"
|
"syscall"
|
||||||
@@ -35,6 +36,12 @@ type ImageProxy struct {
|
|||||||
cacheDir string
|
cacheDir string
|
||||||
mu sync.Mutex
|
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
|
// libraryRootsFn returns the configured media library roots so that
|
||||||
// sidecar poster/artwork files stored alongside media (under arbitrary
|
// sidecar poster/artwork files stored alongside media (under arbitrary
|
||||||
// per-library paths) are allowed by isAllowedLocalPath. It is provided
|
// per-library paths) are allowed by isAllowedLocalPath. It is provided
|
||||||
@@ -55,6 +62,7 @@ type ImageProxy struct {
|
|||||||
const (
|
const (
|
||||||
imageBrowserCacheControl = "public, max-age=2592000, immutable"
|
imageBrowserCacheControl = "public, max-age=2592000, immutable"
|
||||||
imagePlaceholderCacheControl = "no-store"
|
imagePlaceholderCacheControl = "no-store"
|
||||||
|
imageMaxResizeConcurrency = 4
|
||||||
)
|
)
|
||||||
|
|
||||||
// NewImageProxy is the constructor.
|
// NewImageProxy is the constructor.
|
||||||
@@ -64,6 +72,7 @@ func NewImageProxy(cfg *config.Config, log *zap.Logger) *ImageProxy {
|
|||||||
log: log,
|
log: log,
|
||||||
cacheDir: filepath.Join(cfg.Cache.CacheDir, "images"),
|
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
|
// Honor HTTP(S)_PROXY env vars so deployments behind GFW can pull
|
||||||
// from image.tmdb.org via their HTTP proxy without extra config. On
|
// 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]
|
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.
|
// Prune removes oldest cached images until disk usage is within the configured limit.
|
||||||
func (p *ImageProxy) Prune() (PruneImageCacheResult, error) {
|
func (p *ImageProxy) Prune() (PruneImageCacheResult, error) {
|
||||||
if p.cfg == nil || p.cfg.Cache.ImagesMaxSizeMB <= 0 {
|
if p.cfg == nil || p.cfg.Cache.ImagesMaxSizeMB <= 0 {
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
|
"io"
|
||||||
"net/http"
|
"net/http"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
@@ -134,12 +135,38 @@ func (p *ImageProxy) serveCachedImage(w http.ResponseWriter, r *http.Request, ke
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (p *ImageProxy) removeUnusableImageCache(cachePath, failPath string) {
|
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 {
|
if err != nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
ctype := detectContentType(data)
|
stat, err := file.Stat()
|
||||||
if len(data) > 0 && isImageContentType(ctype) && !isTransparentPlaceholderData(data) {
|
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
|
return
|
||||||
}
|
}
|
||||||
_ = os.Remove(cachePath)
|
_ = os.Remove(cachePath)
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package service
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"crypto/sha256"
|
"crypto/sha256"
|
||||||
"encoding/hex"
|
"encoding/hex"
|
||||||
"errors"
|
"errors"
|
||||||
@@ -203,6 +204,28 @@ func (p *ImageProxy) resizeCachePath(key string) string {
|
|||||||
return filepath.Join(p.cacheDir, imageResizeCacheSubdir, key+".img")
|
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 读取图片,按选项缩放后写出,并把结果缓存
|
// serveResizedFromFile 从 srcPath 读取图片,按选项缩放后写出,并把结果缓存
|
||||||
// 到磁盘以免每次请求都重新解码。原图已满足目标尺寸时直接输出原文件。
|
// 到磁盘以免每次请求都重新解码。原图已满足目标尺寸时直接输出原文件。
|
||||||
// 返回 false 表示缩放不可用,调用方应回退到原图直出。
|
// 返回 false 表示缩放不可用,调用方应回退到原图直出。
|
||||||
@@ -211,6 +234,25 @@ func (p *ImageProxy) serveResizedFromFile(w http.ResponseWriter, r *http.Request
|
|||||||
if err != nil || stat.IsDir() || stat.Size() <= 0 {
|
if err != nil || stat.IsDir() || stat.Size() <= 0 {
|
||||||
return false
|
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.
|
data, err := os.ReadFile(srcPath) // #nosec G304 -- srcPath comes from an allowed local path or a SHA-derived cache path.
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return false
|
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)
|
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)
|
p.writeResizeCache(cachePath, out)
|
||||||
w.Header().Set("Content-Type", ctype)
|
w.Header().Set("Content-Type", ctype)
|
||||||
w.Header().Set("Cache-Control", imageBrowserCacheControl)
|
w.Header().Set("Cache-Control", imageBrowserCacheControl)
|
||||||
|
|||||||
@@ -230,3 +230,47 @@ func TestServeResizedFromFileCachesScaledResult(t *testing.T) {
|
|||||||
t.Fatal("expected a non-empty body")
|
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")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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)
|
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) {
|
func embyPremiereDate(value string) (time.Time, bool) {
|
||||||
value = normalizeReleaseDate(value)
|
value = normalizeReleaseDate(value)
|
||||||
if value == "" {
|
if value == "" {
|
||||||
|
|||||||
Reference in New Issue
Block a user