Improve large library storage backends and scans

This commit is contained in:
ShukeBta
2026-06-15 18:16:34 +08:00
parent c9e6adf041
commit a5bf4bfdd4
36 changed files with 2340 additions and 93 deletions
+1 -1
View File
@@ -51,7 +51,7 @@ func (c *Container) runBootCloudLibraryScanQueue(cloudLibs []model.Library) {
for _, lib := range cloudLibs {
libID := lib.ID
libName := lib.Name
scanCtx, cancel := context.WithTimeout(context.Background(), 2*time.Hour)
scanCtx, cancel := cloudScanContext(context.Background(), cloudScanTimeout(context.Background(), c.Repo, 24*time.Hour))
c.Log.Info("boot: scanning cloud library", zap.String("id", libID), zap.String("name", libName))
if _, err := c.Scan.ScanLibraryWithoutAutoScrape(scanCtx, libID); err != nil {
c.Log.Warn("boot: cloud library scan failed", zap.String("id", libID), zap.String("name", libName), zap.Error(err))
+78 -3
View File
@@ -61,6 +61,7 @@ type EmbyService struct {
repo *repository.Container
storage cloudPlaybackResolver
probe cloudPlaybackProber
cache *RuntimeCacheService
virtualMu sync.RWMutex
virtualSeries map[string]embySeriesCacheEntry
@@ -87,6 +88,13 @@ func NewEmbyService(cfg *config.Config, log *zap.Logger, repo *repository.Contai
return &EmbyService{cfg: cfg, log: log, repo: repo}
}
func (e *EmbyService) SetRuntimeCache(cache *RuntimeCacheService) *EmbyService {
if e != nil {
e.cache = cache
}
return e
}
func (e *EmbyService) SetCloudProbe(storage cloudPlaybackResolver, probe cloudPlaybackProber) {
if e == nil {
return
@@ -439,6 +447,11 @@ func (e *EmbyService) Items(ctx context.Context, p ItemsParams) (map[string]any,
}
func (e *EmbyService) mediaItems(ctx context.Context, p ItemsParams) (map[string]any, error) {
cacheKey := e.embyItemsCacheKey("items", p)
var cached embyItemsCacheValue
if e.cache != nil && e.cache.GetJSON(ctx, cacheKey, &cached) {
return map[string]any{"Items": cached.Items, "TotalRecordCount": cached.TotalRecordCount, "StartIndex": cached.StartIndex}, nil
}
q := e.repo.DB.WithContext(ctx).Model(&model.Media{})
q = e.applyUserMediaVisibility(ctx, q, p.UserID)
if p.ParentID != "" {
@@ -503,7 +516,57 @@ func (e *EmbyService) mediaItems(ctx context.Context, p ItemsParams) (map[string
if err != nil {
return nil, err
}
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: total, StartIndex: p.StartIndex}, time.Duration(e.mediaCacheTTLSeconds())*time.Second)
}
return out, nil
}
type embyItemsCacheValue struct {
Items []map[string]any `json:"items"`
TotalRecordCount int64 `json:"total_record_count"`
StartIndex int `json:"start_index"`
}
type embyLatestCacheValue struct {
Items []map[string]any `json:"items"`
}
func (e *EmbyService) embyItemsCacheKey(kind string, p ItemsParams) string {
includeTypes := append([]string(nil), p.IncludeItemTypes...)
filters := append([]string(nil), p.Filters...)
ids := append([]string(nil), p.IDs...)
sort.Strings(includeTypes)
sort.Strings(filters)
sort.Strings(ids)
sum := sha256.Sum256([]byte(strings.Join([]string{
kind,
p.UserID,
p.ParentID,
strings.Join(ids, ","),
p.SearchTerm,
strings.Join(includeTypes, ","),
strings.Join(filters, ","),
strconv.FormatBool(p.Recursive),
p.SortBy,
p.SortOrder,
strconv.Itoa(p.StartIndex),
strconv.Itoa(p.Limit),
}, "|")))
return "media:emby:" + hex.EncodeToString(sum[:])
}
func (e *EmbyService) embyLatestCacheKey(userID, parentID string, limit int) string {
sum := sha256.Sum256([]byte(strings.Join([]string{"latest", userID, parentID, strconv.Itoa(limit)}, "|")))
return "media:emby:" + hex.EncodeToString(sum[:])
}
func (e *EmbyService) mediaCacheTTLSeconds() int {
if e == nil || e.cfg == nil || e.cfg.Cache.MediaTTLSeconds < 1 {
return 15
}
return e.cfg.Cache.MediaTTLSeconds
}
func (e *EmbyService) episodeItems(ctx context.Context, rows []model.Media, p ItemsParams) (map[string]any, error) {
@@ -636,13 +699,22 @@ func (e *EmbyService) LatestItems(ctx context.Context, userID, parentID string,
if limit <= 0 || limit > 100 {
limit = 20
}
cacheKey := e.embyLatestCacheKey(userID, parentID, limit)
var cached embyLatestCacheValue
if e.cache != nil && e.cache.GetJSON(ctx, cacheKey, &cached) {
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 {
return e.latestSeriesItemsForLibrary(ctx, userID, parentID, limit)
out, err := e.latestSeriesItemsForLibrary(ctx, userID, parentID, limit)
if err == nil && e.cache != nil {
e.cache.SetJSON(ctx, cacheKey, embyLatestCacheValue{Items: out}, time.Duration(e.mediaCacheTTLSeconds())*time.Second)
}
return out, err
}
q = q.Where("library_id = ?", parentID)
q = q.Where("library_id IN ?", e.mergedLibraryIDs(ctx, parentID))
}
var rows []model.Media
if err := q.Order("media.created_at desc").Limit(limit).Find(&rows).Error; err != nil {
@@ -669,6 +741,9 @@ func (e *EmbyService) LatestItems(ctx context.Context, userID, parentID string,
for _, m := range rows {
out = append(out, e.itemPayload(ctx, &m, favs[m.ID], 0))
}
if e.cache != nil {
e.cache.SetJSON(ctx, cacheKey, embyLatestCacheValue{Items: out}, time.Duration(e.mediaCacheTTLSeconds())*time.Second)
}
return out, nil
}
+40
View File
@@ -103,6 +103,46 @@ func TestEmbyItemsExposeSeriesSeasonEpisodeHierarchy(t *testing.T) {
}
}
func TestEmbyLatestItemsIncludesMergedCloudMovieLibrary(t *testing.T) {
svc := newTestEmbyService(t)
local := model.Library{Name: "国产电影", Path: `/media/国产电影`, Type: "movie", Enabled: true}
cloud := model.Library{Name: "OpenList · 国产电影", Path: BuildCloudLibraryPath("openlist", "/国产电影", "/国产电影"), Type: "movie", Enabled: true}
for _, lib := range []*model.Library{&local, &cloud} {
if err := svc.repo.Library.Create(t.Context(), lib); err != nil {
t.Fatalf("create library: %v", err)
}
}
for _, media := range []model.Media{
{
Base: model.Base{ID: "local-movie", CreatedAt: time.Now().Add(-time.Minute)},
LibraryID: local.ID,
Title: "本地版本",
Path: `/media/国产电影/local.mkv`,
},
{
Base: model.Base{ID: "cloud-movie", CreatedAt: time.Now()},
LibraryID: cloud.ID,
Title: "云盘版本",
Path: `cloud://openlist/国产电影/cloud.mkv`,
},
} {
if err := svc.repo.DB.Create(&media).Error; err != nil {
t.Fatalf("create media: %v", err)
}
}
latest, err := svc.LatestItems(t.Context(), "user-1", local.ID, 10)
if err != nil {
t.Fatalf("latest items: %v", err)
}
if len(latest) != 2 {
t.Fatalf("latest items = %#v, want local and merged cloud media", latest)
}
if latest[0]["Id"] != "cloud-movie" || latest[1]["Id"] != "local-movie" {
t.Fatalf("latest order/items = %#v, want cloud then local", latest)
}
}
func TestEmbyVirtualSeriesArtworkUsesListCache(t *testing.T) {
svc := newTestEmbyService(t)
lib := model.Library{Name: "番剧", Path: `/media/anime`, Type: "anime", Enabled: true}
+252 -10
View File
@@ -3,11 +3,15 @@ package service
import (
"context"
"crypto/sha1"
"encoding/hex"
"errors"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"time"
"go.uber.org/zap"
@@ -18,9 +22,10 @@ import (
// MediaService offers high-level CRUD over libraries and media items.
type MediaService struct {
cfg *config.Config
log *zap.Logger
repo *repository.Container
cfg *config.Config
log *zap.Logger
repo *repository.Container
cache *RuntimeCacheService
}
type MediaVisibility struct {
@@ -29,6 +34,11 @@ type MediaVisibility struct {
HiddenLibraryIDs []string
}
type MediaItem struct {
model.Media
Versions []model.Media `json:"versions,omitempty"`
}
const maxMediaSearchLimit = 50000
const maxMediaSearchPageSize = 2000
@@ -60,6 +70,13 @@ func NewMediaService(cfg *config.Config, log *zap.Logger, repo *repository.Conta
return &MediaService{cfg: cfg, log: log, repo: repo}
}
func (s *MediaService) SetRuntimeCache(cache *RuntimeCacheService) *MediaService {
if s != nil {
s.cache = cache
}
return s
}
// CreateLibrary persists a library after validating that its path exists.
func (s *MediaService) CreateLibrary(ctx context.Context, name, path, kind string) (*model.Library, error) {
if name == "" || path == "" {
@@ -74,6 +91,7 @@ func (s *MediaService) CreateLibrary(ctx context.Context, name, path, kind strin
if err := s.repo.Library.Create(ctx, lib); err != nil {
return nil, err
}
s.invalidateMediaCache(ctx)
return lib, nil
}
@@ -306,13 +324,21 @@ func (s *MediaService) DeleteLibrary(ctx context.Context, id string) error {
if err := s.repo.Media.PurgeByLibrary(ctx, id); err != nil {
return err
}
return s.repo.DB.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.Library{}).Error
err := s.repo.DB.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.Library{}).Error
if err == nil {
s.invalidateMediaCache(ctx)
}
return err
}
}
if err := s.repo.Media.DeleteByLibrary(ctx, id); err != nil {
return err
}
return s.repo.Library.Delete(ctx, id)
err = s.repo.Library.Delete(ctx, id)
if err == nil {
s.invalidateMediaCache(ctx)
}
return err
}
// ListMedia paginates media items inside a library.
@@ -335,11 +361,194 @@ func (s *MediaService) ListMediaVisible(ctx context.Context, libraryID string, p
if err != nil {
return nil, 0, err
}
return s.repo.Media.ListByLibrariesFiltered(ctx, libraryIDs, (page-1)*pageSize, pageSize, repository.MediaQueryFilter{
filter := repository.MediaQueryFilter{
IncludeNSFW: visibility.IncludeNSFW,
AllowedLibraryIDs: visibility.AllowedLibraryIDs,
HiddenLibraryIDs: visibility.HiddenLibraryIDs,
}
cacheKey := s.mediaListCacheKey(libraryID, libraryIDs, page, pageSize, filter)
var cached mediaListCacheValue
if s.cache != nil && s.cache.GetJSON(ctx, cacheKey, &cached) {
return cached.Items, cached.Total, nil
}
items, total, err := s.repo.Media.ListByLibrariesFiltered(ctx, libraryIDs, (page-1)*pageSize, pageSize, filter)
if err != nil {
return nil, 0, err
}
if s.cache != nil {
s.cache.SetJSON(ctx, cacheKey, mediaListCacheValue{Items: items, Total: total}, time.Duration(s.mediaCacheTTLSeconds())*time.Second)
}
return items, total, nil
}
func (s *MediaService) ListMediaVisibleGrouped(ctx context.Context, libraryID string, page, pageSize int, visibility MediaVisibility) ([]MediaItem, int64, error) {
items, _, err := s.ListMediaVisible(ctx, libraryID, page, pageSize, visibility)
if err != nil {
return nil, 0, err
}
grouped := groupMediaVersions(items)
return grouped, int64(len(grouped)), nil
}
type mediaListCacheValue struct {
Items []model.Media `json:"items"`
Total int64 `json:"total"`
}
func (s *MediaService) mediaListCacheKey(libraryID string, libraryIDs []string, page, pageSize int, filter repository.MediaQueryFilter) string {
allowed := append([]string(nil), filter.AllowedLibraryIDs...)
hidden := append([]string(nil), filter.HiddenLibraryIDs...)
libs := append([]string(nil), libraryIDs...)
sort.Strings(allowed)
sort.Strings(hidden)
sort.Strings(libs)
sum := sha1.Sum([]byte(strings.Join([]string{
libraryID,
strings.Join(libs, ","),
fmt.Sprintf("%d:%d:%t", page, pageSize, filter.IncludeNSFW),
strings.Join(allowed, ","),
strings.Join(hidden, ","),
}, "|")))
return "media:list:" + hex.EncodeToString(sum[:])
}
func (s *MediaService) mediaCacheTTLSeconds() int {
if s == nil || s.cfg == nil || s.cfg.Cache.MediaTTLSeconds < 1 {
return 15
}
return s.cfg.Cache.MediaTTLSeconds
}
func (s *MediaService) invalidateMediaCache(ctx context.Context) {
if s != nil && s.cache != nil {
s.cache.DeletePrefix(ctx, "media:")
s.cache.DeletePrefix(ctx, "stats:")
}
}
func groupMediaVersions(items []model.Media) []MediaItem {
if len(items) == 0 {
return nil
}
type group struct {
key string
primary model.Media
rows []model.Media
}
groups := make([]group, 0, len(items))
byKey := make(map[string]int, len(items))
for _, item := range items {
key := mediaVersionGroupKey(item)
if key == "" {
groups = append(groups, group{primary: item, rows: []model.Media{item}})
continue
}
if idx, ok := byKey[key]; ok {
groups[idx].rows = append(groups[idx].rows, item)
if betterMediaVersion(item, groups[idx].primary) {
groups[idx].primary = item
}
continue
}
byKey[key] = len(groups)
groups = append(groups, group{key: key, primary: item, rows: []model.Media{item}})
}
out := make([]MediaItem, 0, len(groups))
for _, g := range groups {
sort.SliceStable(g.rows, func(i, j int) bool {
return betterMediaVersion(g.rows[i], g.rows[j])
})
item := MediaItem{Media: g.primary}
if len(g.rows) > 1 {
item.Versions = g.rows
}
out = append(out, item)
}
sort.SliceStable(out, func(i, j int) bool {
return out[i].CreatedAt.After(out[j].CreatedAt)
})
return out
}
func mediaVersionGroupKey(m model.Media) string {
if m.SeasonNum > 0 || m.EpisodeNum > 0 {
title := firstNonEmpty(m.OriginalName, m.Title)
if title == "" {
title, _ = CleanQuery(m.Path)
}
title = normalizeMediaVersionText(title)
if title == "" {
return ""
}
return strings.Join([]string{
"episode",
strings.ToLower(strings.TrimSpace(m.LibraryID)),
title,
fmt.Sprintf("%d:%d", m.SeasonNum, m.EpisodeNum),
}, "|")
}
switch {
case m.TMDbID > 0:
return fmt.Sprintf("tmdb:%d", m.TMDbID)
case m.BangumiID > 0:
return fmt.Sprintf("bangumi:%d", m.BangumiID)
case strings.TrimSpace(m.DoubanID) != "":
return "douban:" + strings.ToLower(strings.TrimSpace(m.DoubanID))
case strings.TrimSpace(m.TheTVDBID) != "":
return "thetvdb:" + strings.ToLower(strings.TrimSpace(m.TheTVDBID))
}
title := firstNonEmpty(m.OriginalName, m.Title)
if title == "" {
title, _ = CleanQuery(m.Path)
}
title = normalizeMediaVersionText(title)
if title == "" {
return ""
}
year := m.Year
if year <= 0 {
_, year = CleanQuery(m.Path)
}
return fmt.Sprintf("movie:%s:%d", title, year)
}
func normalizeMediaVersionText(value string) string {
value = strings.ToLower(strings.TrimSpace(value))
if value == "" {
return ""
}
fields := strings.FieldsFunc(value, func(r rune) bool {
switch r {
case '.', '_', '-', ' ', '\t', '/', '\\', '[', ']', '(', ')', '(', ')', '【', '】':
return true
default:
return false
}
})
out := fields[:0]
for _, field := range fields {
field = strings.TrimSpace(field)
if field == "" {
continue
}
if _, noise := noiseTokenSet[field]; noise {
continue
}
out = append(out, field)
}
return strings.Join(out, " ")
}
func betterMediaVersion(candidate, current model.Media) bool {
candidatePixels := candidate.Width * candidate.Height
currentPixels := current.Width * current.Height
if candidatePixels != currentPixels {
return candidatePixels > currentPixels
}
if candidate.SizeBytes != current.SizeBytes {
return candidate.SizeBytes > current.SizeBytes
}
return candidate.CreatedAt.After(current.CreatedAt)
}
// SearchMedia performs a simple LIKE search across titles.
@@ -361,6 +570,14 @@ func (s *MediaService) SearchMediaVisible(ctx context.Context, query string, lim
})
}
func (s *MediaService) SearchMediaVisibleGrouped(ctx context.Context, query string, limit int, visibility MediaVisibility) ([]MediaItem, error) {
items, err := s.SearchMediaVisible(ctx, query, limit, visibility)
if err != nil {
return nil, err
}
return groupMediaVersions(items), nil
}
func (s *MediaService) SearchMediaVisiblePage(ctx context.Context, query string, page, pageSize int, visibility MediaVisibility) ([]model.Media, int64, error) {
if pageSize <= 0 {
pageSize = 50
@@ -379,6 +596,15 @@ func (s *MediaService) SearchMediaVisiblePage(ctx context.Context, query string,
})
}
func (s *MediaService) SearchMediaVisiblePageGrouped(ctx context.Context, query string, page, pageSize int, visibility MediaVisibility) ([]MediaItem, int64, error) {
items, _, err := s.SearchMediaVisiblePage(ctx, query, page, pageSize, visibility)
if err != nil {
return nil, 0, err
}
grouped := groupMediaVersions(items)
return grouped, int64(len(grouped)), nil
}
// GetMedia returns a single media row.
func (s *MediaService) GetMedia(ctx context.Context, id string) (*model.Media, error) {
return s.repo.Media.FindByID(ctx, id)
@@ -392,15 +618,27 @@ func (s *MediaService) SoftDelete(ctx context.Context, id string) error {
return err
}
if media != nil && isCloudMediaPath(media.Path) {
return s.repo.DB.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.Media{}).Error
err := s.repo.DB.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.Media{}).Error
if err == nil {
s.invalidateMediaCache(ctx)
}
return err
}
return s.repo.DB.Where("id = ?", id).Delete(&model.Media{}).Error
err = s.repo.DB.WithContext(ctx).Where("id = ?", id).Delete(&model.Media{}).Error
if err == nil {
s.invalidateMediaCache(ctx)
}
return err
}
// RestoreDeleted unsets DeletedAt for a single media row.
func (s *MediaService) RestoreDeleted(ctx context.Context, id string) error {
return s.repo.DB.Unscoped().Model(&model.Media{}).
err := s.repo.DB.WithContext(ctx).Unscoped().Model(&model.Media{}).
Where("id = ?", id).Update("deleted_at", nil).Error
if err == nil {
s.invalidateMediaCache(ctx)
}
return err
}
// ListRecycleBin returns every soft-deleted row, newest first.
@@ -419,5 +657,9 @@ func (s *MediaService) ListRecycleBin(ctx context.Context, limit int) ([]model.M
// PurgeDeleted permanently removes a soft-deleted row from the database.
func (s *MediaService) PurgeDeleted(ctx context.Context, id string) error {
return s.repo.DB.Unscoped().Where("id = ?", id).Delete(&model.Media{}).Error
err := s.repo.DB.WithContext(ctx).Unscoped().Where("id = ?", id).Delete(&model.Media{}).Error
if err == nil {
s.invalidateMediaCache(ctx)
}
return err
}
+37
View File
@@ -4,6 +4,7 @@ import (
"os"
"path/filepath"
"testing"
"time"
"github.com/glebarez/sqlite"
"go.uber.org/zap"
@@ -264,3 +265,39 @@ func TestSoftDeleteCloudMediaPurgesRecordWithoutRecycleBin(t *testing.T) {
t.Fatalf("cloud media removal must not populate recycle bin: %#v", recycle)
}
}
func TestSoftDeleteInvalidatesMediaAndStatsCache(t *testing.T) {
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err := db.AutoMigrate(&model.Media{}); err != nil {
t.Fatal(err)
}
repos := repository.New(db)
media := model.Media{
Base: model.Base{ID: "local-media"},
Title: "Cached Movie",
Path: filepath.Join(t.TempDir(), "Cached Movie.mkv"),
}
if err := repos.DB.Create(&media).Error; err != nil {
t.Fatal(err)
}
cache := NewRuntimeCacheService(&config.Config{}, zap.NewNop())
cache.SetJSON(t.Context(), "media:list:test", map[string]string{"state": "stale"}, time.Minute)
cache.SetJSON(t.Context(), "stats:snapshot:base", map[string]int{"media": 1}, time.Minute)
svc := NewMediaService(&config.Config{}, zap.NewNop(), repos).SetRuntimeCache(cache)
if err := svc.SoftDelete(t.Context(), media.ID); err != nil {
t.Fatal(err)
}
var mediaCache map[string]string
if cache.GetJSON(t.Context(), "media:list:test", &mediaCache) {
t.Fatal("soft delete should invalidate media cache")
}
var statsCache map[string]int
if cache.GetJSON(t.Context(), "stats:snapshot:base", &statsCache) {
t.Fatal("soft delete should invalidate stats cache")
}
}
+2 -1
View File
@@ -283,7 +283,8 @@ func (o *OrganizerService) resolveTransferMode(ctx context.Context, override Tra
}
if mode == TransferMove && o.keepSeedingEnabled(ctx) {
// 移动会删除源文件导致 qBittorrent 停止做种;保种开启时改用硬链接
//(跨盘自动退化为复制),既规范命名又保留源文件继续做种上传。
// 既规范命名又保留源文件继续做种上传。硬链接失败时会报错,避免静默
// 退化复制后占用双份磁盘空间。
return TransferHardlink
}
return mode
+193
View File
@@ -0,0 +1,193 @@
package service
import (
"context"
"encoding/json"
"strings"
"sync"
"time"
"github.com/redis/go-redis/v9"
"go.uber.org/zap"
"github.com/ShukeBta/MediaStationGo/internal/config"
)
type RuntimeCacheService struct {
log *zap.Logger
client *redis.Client
prefix string
mu sync.RWMutex
memory map[string]runtimeCacheItem
limit int
}
type runtimeCacheItem struct {
raw []byte
expiresAt time.Time
}
func NewRuntimeCacheService(cfg *config.Config, log *zap.Logger) *RuntimeCacheService {
c := &RuntimeCacheService{log: log, memory: map[string]runtimeCacheItem{}, limit: 2048}
if cfg == nil {
return c
}
c.prefix = strings.Trim(strings.TrimSpace(cfg.Cache.RedisPrefix), ":")
if c.prefix == "" {
c.prefix = "mediastationgo"
}
rawURL := strings.TrimSpace(cfg.Cache.RedisURL)
if rawURL == "" {
return c
}
opts, err := redis.ParseURL(rawURL)
if err != nil {
if log != nil {
log.Warn("redis cache disabled: invalid redis url", zap.Error(err))
}
return c
}
client := redis.NewClient(opts)
pingCtx, cancel := context.WithTimeout(context.Background(), 1200*time.Millisecond)
defer cancel()
if err := client.Ping(pingCtx).Err(); err != nil {
if log != nil {
log.Warn("redis cache unavailable; using in-process cache", zap.Error(err))
}
_ = client.Close()
return c
}
c.client = client
if log != nil {
log.Info("redis runtime cache enabled with in-process L1", zap.String("addr", opts.Addr), zap.String("prefix", c.prefix))
}
return c
}
func (c *RuntimeCacheService) Enabled() bool {
return c != nil
}
func (c *RuntimeCacheService) Close() error {
if c == nil || c.client == nil {
return nil
}
return c.client.Close()
}
func (c *RuntimeCacheService) GetJSON(ctx context.Context, key string, out any) bool {
if !c.Enabled() || strings.TrimSpace(key) == "" || out == nil {
return false
}
fullKey := c.key(key)
if raw, ok := c.getMemory(fullKey); ok {
return json.Unmarshal(raw, out) == nil
}
if c.client != nil {
raw, err := c.client.Get(ctx, fullKey).Bytes()
if err == nil {
if json.Unmarshal(raw, out) != nil {
return false
}
c.setMemory(fullKey, raw, 2*time.Second)
return true
}
}
return false
}
func (c *RuntimeCacheService) SetJSON(ctx context.Context, key string, value any, ttl time.Duration) {
if !c.Enabled() || strings.TrimSpace(key) == "" || value == nil || ttl <= 0 {
return
}
raw, err := json.Marshal(value)
if err != nil {
return
}
fullKey := c.key(key)
c.setMemory(fullKey, raw, ttl)
if c.client != nil {
_ = c.client.Set(ctx, fullKey, raw, ttl).Err()
}
}
func (c *RuntimeCacheService) DeletePrefix(ctx context.Context, prefix string) {
if !c.Enabled() || strings.TrimSpace(prefix) == "" {
return
}
fullPrefix := c.key(prefix)
c.deleteMemoryPrefix(fullPrefix)
if c.client != nil {
pattern := fullPrefix + "*"
var cursor uint64
for {
keys, next, err := c.client.Scan(ctx, cursor, pattern, 200).Result()
if err != nil {
return
}
if len(keys) > 0 {
_ = c.client.Del(ctx, keys...).Err()
}
cursor = next
if cursor == 0 {
return
}
}
}
}
func (c *RuntimeCacheService) key(key string) string {
key = strings.TrimLeft(strings.TrimSpace(key), ":")
if c.prefix == "" {
return key
}
return c.prefix + ":" + key
}
func (c *RuntimeCacheService) getMemory(key string) ([]byte, bool) {
now := time.Now()
c.mu.RLock()
item, ok := c.memory[key]
c.mu.RUnlock()
if !ok {
return nil, false
}
if now.After(item.expiresAt) {
c.mu.Lock()
delete(c.memory, key)
c.mu.Unlock()
return nil, false
}
return item.raw, true
}
func (c *RuntimeCacheService) setMemory(key string, raw []byte, ttl time.Duration) {
if ttl <= 0 || len(raw) == 0 {
return
}
c.mu.Lock()
defer c.mu.Unlock()
if len(c.memory) >= c.limit {
now := time.Now()
for k, item := range c.memory {
if now.After(item.expiresAt) || len(c.memory) >= c.limit {
delete(c.memory, k)
}
if len(c.memory) < c.limit {
break
}
}
}
c.memory[key] = runtimeCacheItem{raw: append([]byte(nil), raw...), expiresAt: time.Now().Add(ttl)}
}
func (c *RuntimeCacheService) deleteMemoryPrefix(prefix string) {
c.mu.Lock()
defer c.mu.Unlock()
for key := range c.memory {
if strings.HasPrefix(key, prefix) {
delete(c.memory, key)
}
}
}
+302 -41
View File
@@ -56,6 +56,7 @@ type ScannerService struct {
probe *FFprobeService
scraper *ScraperService
storage *StorageConfigService
cache *RuntimeCacheService
imageProxy *ImageProxy
@@ -113,6 +114,12 @@ func (s *ScannerService) SetStorageConfig(storage *StorageConfigService) {
}
}
func (s *ScannerService) SetRuntimeCache(cache *RuntimeCacheService) {
if s != nil {
s.cache = cache
}
}
// SetImageProxy lets cloud scans warm sidecar poster/backdrop files into the
// local image cache. This keeps library opening fast without forcing the UI or
// Emby clients to resolve/download every cloud poster on demand.
@@ -255,18 +262,38 @@ func isCloudArtworkRef(ref string) bool {
// ScanResult summarises a scan run.
type ScanResult struct {
LibraryID string `json:"library_id"`
Visited int `json:"visited"`
Added int `json:"added"`
Updated int `json:"updated"`
Skipped int `json:"skipped"`
Probed int `json:"probed"`
LocalMetadata int `json:"local_metadata"`
Removed int64 `json:"removed"`
LibraryID string `json:"library_id"`
Visited int `json:"visited"`
Added int `json:"added"`
Updated int `json:"updated"`
Skipped int `json:"skipped"`
Probed int `json:"probed"`
LocalMetadata int `json:"local_metadata"`
Removed int64 `json:"removed"`
ErrorCount int `json:"error_count,omitempty"`
Errors []string `json:"errors,omitempty"`
}
var ErrCloudScanAlreadyRunning = errors.New("cloud scan already running")
const maxScanErrorDetails = 20
func addScanError(res *ScanResult, path string, err error) {
if res == nil || err == nil {
return
}
res.ErrorCount++
if len(res.Errors) >= maxScanErrorDetails {
return
}
path = strings.TrimSpace(path)
msg := strings.TrimSpace(err.Error())
if path != "" {
msg = path + ": " + msg
}
res.Errors = append(res.Errors, msg)
}
const maxCloudMediaProbeQueuePerScan = 32
const cloudMediaProbeFailureBackoff = 6 * time.Hour
@@ -291,6 +318,8 @@ type CloudScanStatus struct {
Updated int `json:"updated"`
Skipped int `json:"skipped"`
Removed int64 `json:"removed"`
ErrorCount int `json:"error_count,omitempty"`
Errors []string `json:"errors,omitempty"`
Error string `json:"error,omitempty"`
ResumeHint string `json:"resume_hint,omitempty"`
Estimate string `json:"estimate_message,omitempty"`
@@ -327,6 +356,17 @@ type existingCloudMedia struct {
STRMURL string
}
type existingLocalMedia struct {
SizeBytes int64
DurationSec int
Width int
Height int
VideoCodec string
AudioCodec string
Container string
STRMURL string
}
func (s *ScannerService) cloudMediaProbeWorker() {
for task := range s.cloudMediaProbeQueue {
s.probeCloudMediaAsync(task)
@@ -451,6 +491,8 @@ func (s *ScannerService) beginCloudScan(ctx context.Context, lib *model.Library,
current.status.Updated = res.Updated
current.status.Skipped = res.Skipped
current.status.Removed = res.Removed
current.status.ErrorCount = res.ErrorCount
current.status.Errors = append([]string(nil), res.Errors...)
}
current.status.UpdatedAt = now
current.status.FinishedAt = now
@@ -471,22 +513,28 @@ func (s *ScannerService) beginCloudScan(ctx context.Context, lib *model.Library,
default:
current.status.State = "finished"
current.status.Stage = "finished"
current.status.Error = ""
if current.status.ErrorCount > 0 {
current.status.Error = fmt.Sprintf("部分文件入库失败:%d 个,详情见 errors", current.status.ErrorCount)
} else {
current.status.Error = ""
}
}
if s.hub != nil {
s.hub.Publish("scan", map[string]any{
"library_id": lib.ID,
"provider": mount.Provider,
"cloud": true,
"finished": true,
"state": current.status.State,
"stage": current.status.Stage,
"error": current.status.Error,
"visited": current.status.Visited,
"added": current.status.Added,
"updated": current.status.Updated,
"skipped": current.status.Skipped,
"removed": current.status.Removed,
"library_id": lib.ID,
"provider": mount.Provider,
"cloud": true,
"finished": true,
"state": current.status.State,
"stage": current.status.Stage,
"error": current.status.Error,
"visited": current.status.Visited,
"added": current.status.Added,
"updated": current.status.Updated,
"skipped": current.status.Skipped,
"removed": current.status.Removed,
"error_count": current.status.ErrorCount,
"errors": current.status.Errors,
})
}
}
@@ -643,7 +691,7 @@ func (s *ScannerService) StartCloudLibraryScan(libraryID string, autoScrape bool
s.cloudScanMu.Unlock()
go func() {
ctx, cancel := context.WithTimeout(context.Background(), 6*time.Hour)
ctx, cancel := cloudScanContext(context.Background(), cloudScanTimeout(context.Background(), s.repo, 24*time.Hour))
defer cancel()
if autoScrape {
_, err = s.ScanLibrary(ctx, libraryID)
@@ -667,6 +715,28 @@ func (s *ScannerService) StartCloudLibraryScan(libraryID string, autoScrape bool
return status, true, nil
}
func cloudScanContext(parent context.Context, timeout time.Duration) (context.Context, context.CancelFunc) {
if timeout <= 0 {
return context.WithCancel(parent)
}
return context.WithTimeout(parent, timeout)
}
func cloudScanTimeout(ctx context.Context, repo *repository.Container, fallback time.Duration) time.Duration {
if repo == nil || repo.Setting == nil {
return fallback
}
value, err := repo.Setting.Get(ctx, "cloud.scan_timeout_hours")
if err != nil || strings.TrimSpace(value) == "" {
return fallback
}
hours := parseIntSettingDefault(strings.TrimSpace(value), int(fallback/time.Hour))
if hours <= 0 {
return 0
}
return time.Duration(hours) * time.Hour
}
func (s *ScannerService) StartAllCloudLibraryScans() ([]CloudScanStatus, error) {
if s == nil {
return nil, errors.New("scanner unavailable")
@@ -754,8 +824,19 @@ func (s *ScannerService) scanLibrary(ctx context.Context, libraryID string, auto
res := &ScanResult{LibraryID: lib.ID}
seen := make(map[string]struct{})
seenInodes := make(map[string]string)
writeBatch := newLocalMediaWriteBatch(s, ctx, res, 100)
existingMedia, err := s.existingLocalMediaSnapshot(ctx, lib.ID)
if err != nil {
s.log.Warn("load existing local media snapshot failed", zap.String("library_id", lib.ID), zap.Error(err))
existingMedia = nil
}
walkFn := func(path string, info walkInfo) error {
select {
case <-ctx.Done():
return ctx.Err()
default:
}
if info.isDir {
return nil
}
@@ -764,12 +845,18 @@ func (s *ScannerService) scanLibrary(ctx context.Context, libraryID string, auto
return nil
}
seen[filepath.Clean(path)] = struct{}{}
s.ingestFile(ctx, lib, path, info.size, seenInodes, res)
s.ingestFile(ctx, lib, path, info.size, seenInodes, existingMedia, writeBatch, res)
return nil
}
if err := walk(lib.Path, walkFn); err != nil {
return res, err
walkErr := walk(lib.Path, walkFn)
writeBatch.Flush()
if walkErr != nil {
addScanError(res, lib.Path, walkErr)
if res.Added+res.Updated > 0 {
s.invalidateMediaCache(ctx)
}
return res, walkErr
}
removed, err := s.pruneMissingMedia(ctx, lib.ID, seen)
if err != nil {
@@ -779,15 +866,18 @@ func (s *ScannerService) scanLibrary(ctx context.Context, libraryID string, auto
}
s.hub.Publish("scan", map[string]any{
"library_id": lib.ID,
"finished": true,
"visited": res.Visited,
"added": res.Added,
"updated": res.Updated,
"probed": res.Probed,
"local_meta": res.LocalMetadata,
"removed": res.Removed,
"library_id": lib.ID,
"finished": true,
"visited": res.Visited,
"added": res.Added,
"updated": res.Updated,
"probed": res.Probed,
"local_meta": res.LocalMetadata,
"removed": res.Removed,
"error_count": res.ErrorCount,
"errors": res.Errors,
})
s.invalidateMediaCache(ctx)
s.maybeGenerateSTRMAfterScan(lib.ID)
// Online enrichment is opt-in. Local NFO is always consumed first during
@@ -820,7 +910,10 @@ func (s *ScannerService) IngestPath(ctx context.Context, libraryID, path string)
return false, nil
}
res := &ScanResult{LibraryID: lib.ID}
s.ingestFile(ctx, lib, path, fi.Size(), make(map[string]string), res)
s.ingestFile(ctx, lib, path, fi.Size(), make(map[string]string), nil, nil, res)
if res.Added+res.Updated > 0 {
s.invalidateMediaCache(ctx)
}
return res.Added+res.Updated > 0, nil
}
@@ -1054,12 +1147,16 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar
"visited": res.Visited,
"added": res.Added,
"updated": res.Updated,
"skipped": res.Skipped,
"removed": res.Removed,
"error_count": res.ErrorCount,
"errors": res.Errors,
"discovered": filesDiscovered,
"dirs": dirsVisited,
"elapsed_seconds": int(time.Since(startedAt).Seconds()),
"cloud": true,
})
s.invalidateMediaCache(ctx)
s.maybeGenerateSTRMAfterScan(lib.ID)
if autoScrape && s.scraper != nil && s.scraper.AnyEnabled() && s.autoScrapeEnabled(ctx) {
s.startAutoScrape(ctx, lib.ID)
@@ -1067,6 +1164,13 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar
return res, nil
}
func (s *ScannerService) invalidateMediaCache(ctx context.Context) {
if s != nil && s.cache != nil {
s.cache.DeletePrefix(ctx, "media:")
s.cache.DeletePrefix(ctx, "stats:")
}
}
func (s *ScannerService) startAutoScrape(ctx context.Context, libraryID string) {
scrapeCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Minute)
go func() {
@@ -1118,6 +1222,43 @@ func (s *ScannerService) existingCloudMediaSnapshot(ctx context.Context, library
return out, nil
}
func (s *ScannerService) existingLocalMediaSnapshot(ctx context.Context, libraryID string) (map[string]existingLocalMedia, error) {
var rows []struct {
Path string
SizeBytes int64
DurationSec int
Width int
Height int
VideoCodec string
AudioCodec string
Container string
STRMURL string
}
if err := s.repo.DB.WithContext(ctx).
Model(&model.Media{}).
Select("path, size_bytes, duration_sec, width, height, video_codec, audio_codec, container, strm_url").
Where("library_id = ? AND path NOT LIKE ?", libraryID, "cloud://%").
Find(&rows).Error; err != nil {
return nil, err
}
out := make(map[string]existingLocalMedia, len(rows))
for _, row := range rows {
if row.Path != "" {
out[filepath.Clean(row.Path)] = existingLocalMedia{
SizeBytes: row.SizeBytes,
DurationSec: row.DurationSec,
Width: row.Width,
Height: row.Height,
VideoCodec: row.VideoCodec,
AudioCodec: row.AudioCodec,
Container: row.Container,
STRMURL: row.STRMURL,
}
}
}
return out, nil
}
func (s *ScannerService) shadowedCloudLibrary(ctx context.Context, lib *model.Library) *CloudMountConflict {
libs, err := s.repo.Library.List(ctx)
if err != nil {
@@ -1213,6 +1354,7 @@ func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library
s.queueCloudArtworkPrefetch(localMeta.BackdropURL)
}
if err := s.repo.Media.Upsert(ctx, m); err != nil {
addScanError(res, path, err)
s.log.Warn("upsert cloud media failed", zap.String("path", path), zap.Error(err))
return
}
@@ -1343,6 +1485,18 @@ func cloudTrackMetadataMissing(existing existingCloudMedia) bool {
strings.TrimSpace(existing.AudioCodec) == ""
}
func localTrackMetadataMissing(existing existingLocalMedia) bool {
return existing.DurationSec <= 0 ||
existing.Width <= 0 ||
existing.Height <= 0 ||
strings.TrimSpace(existing.VideoCodec) == "" ||
strings.TrimSpace(existing.AudioCodec) == ""
}
func localMetadataNeedsRefresh(local *LocalMetadata) bool {
return local != nil && (local.HasNFO || local.HasArtwork || localHasDescriptiveMetadata(local))
}
func cloudSeriesTitleFromMediaPath(mediaPath string) (string, int) {
displayPath := strings.TrimSpace(mediaPath)
if strings.HasPrefix(strings.ToLower(displayPath), "cloud://") {
@@ -1394,12 +1548,15 @@ func (s *ScannerService) RemovePath(ctx context.Context, path string) (int64, er
res := s.repo.DB.WithContext(ctx).
Where("path = ?", path).
Delete(&model.Media{})
if res.Error == nil && res.RowsAffected > 0 {
s.invalidateMediaCache(ctx)
}
return res.RowsAffected, res.Error
}
// ingestFile upserts a single media file. seenInodes dedups hardlinks within a
// single scan; pass a fresh map for one-off ingests. It mutates res counters.
func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, path string, size int64, seenInodes map[string]string, res *ScanResult) {
func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, path string, size int64, seenInodes map[string]string, existingMedia map[string]existingLocalMedia, writeBatch *localMediaWriteBatch, res *ScanResult) {
res.Visited++
ext := strings.ToLower(filepath.Ext(path))
@@ -1424,7 +1581,27 @@ func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, pat
seenInodes[fileID] = path
}
isNewMedia := !s.mediaPathExists(ctx, path)
cleanPath := filepath.Clean(path)
parsedSeason, parsedEpisode := ParseEpisode(path)
localMeta, localMetaErr := ReadLocalMetadata(path, lib.Path, librarySupportsSeasons(lib) || parsedSeason > 0 || parsedEpisode > 0)
if localMetaErr != nil {
s.log.Warn("read local metadata failed", zap.String("path", path), zap.Error(localMetaErr))
}
isNewMedia := false
if existingMedia != nil {
existing, exists := existingMedia[cleanPath]
isNewMedia = !exists
if exists &&
ext != ".strm" &&
existing.SizeBytes == size &&
(s.probe == nil || !localTrackMetadataMissing(existing)) &&
!localMetadataNeedsRefresh(localMeta) {
res.Skipped++
return
}
} else {
isNewMedia = !s.mediaPathExists(ctx, path)
}
title, year := CleanQuery(path)
if title == "" {
@@ -1449,15 +1626,12 @@ func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, pat
}
}
parsedSeason, parsedEpisode := ParseEpisode(path)
m.SeasonNum = parsedSeason
m.EpisodeNum = parsedEpisode
if local, err := ReadLocalMetadata(path, lib.Path, librarySupportsSeasons(lib) || parsedSeason > 0 || parsedEpisode > 0); err == nil && local != nil {
applyLocalMetadata(m, local)
if localMeta != nil {
applyLocalMetadata(m, localMeta)
res.LocalMetadata++
} else if err != nil {
s.log.Warn("read local metadata failed", zap.String("path", path), zap.Error(err))
}
// Best-effort ffprobe; failure does not abort the file.
@@ -1477,7 +1651,12 @@ func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, pat
}
}
if isNewMedia && writeBatch != nil {
writeBatch.Add(path, m)
return
}
if err := s.repo.Media.Upsert(ctx, m); err != nil {
addScanError(res, path, err)
s.log.Warn("upsert media failed", zap.String("path", path), zap.Error(err))
return
}
@@ -1497,6 +1676,88 @@ func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, pat
})
}
type localMediaWriteBatch struct {
scanner *ScannerService
ctx context.Context
res *ScanResult
limit int
items []localMediaWriteItem
}
type localMediaWriteItem struct {
path string
media *model.Media
}
func newLocalMediaWriteBatch(scanner *ScannerService, ctx context.Context, res *ScanResult, limit int) *localMediaWriteBatch {
if limit <= 0 {
limit = 100
}
return &localMediaWriteBatch{scanner: scanner, ctx: ctx, res: res, limit: limit}
}
func (b *localMediaWriteBatch) Add(path string, media *model.Media) {
if b == nil || b.scanner == nil || media == nil {
return
}
if media.ScrapeStatus == "" {
media.ScrapeStatus = "pending"
}
b.items = append(b.items, localMediaWriteItem{path: path, media: media})
if len(b.items) >= b.limit {
b.Flush()
}
}
func (b *localMediaWriteBatch) Flush() {
if b == nil || len(b.items) == 0 || b.scanner == nil || b.scanner.repo == nil || b.scanner.repo.DB == nil {
return
}
items := b.items
b.items = nil
media := make([]model.Media, 0, len(items))
for _, item := range items {
if item.media != nil {
media = append(media, *item.media)
}
}
if len(media) == 0 {
return
}
if err := b.scanner.repo.DB.WithContext(b.ctx).CreateInBatches(&media, b.limit).Error; err == nil {
b.res.Added += len(media)
b.publish()
return
}
for _, item := range items {
if item.media == nil {
continue
}
if err := b.scanner.repo.Media.Upsert(b.ctx, item.media); err != nil {
addScanError(b.res, item.path, err)
b.scanner.log.Warn("upsert media failed", zap.String("path", item.path), zap.Error(err))
continue
}
b.res.Added++
}
b.publish()
}
func (b *localMediaWriteBatch) publish() {
if b == nil || b.scanner == nil || b.scanner.hub == nil || b.res == nil {
return
}
b.scanner.hub.Publish("scan", map[string]any{
"library_id": b.res.LibraryID,
"visited": b.res.Visited,
"added": b.res.Added,
"updated": b.res.Updated,
"probed": b.res.Probed,
"local_meta": b.res.LocalMetadata,
"batched": true,
})
}
// duplicateByFileID reports an existing media path that shares the given inode
// identity but lives at a different path and still exists on disk.
func (s *ScannerService) duplicateByFileID(ctx context.Context, fileID, path string) (string, bool) {
@@ -4,6 +4,7 @@ import (
"os"
"path/filepath"
"testing"
"time"
"github.com/glebarez/sqlite"
"go.uber.org/zap"
@@ -95,6 +96,62 @@ func TestScanLibraryReadsLocalSTRMTarget(t *testing.T) {
}
}
func TestScanLibrarySkipsUnchangedExistingLocalMedia(t *testing.T) {
sc, repos := newScannerTestEnv(t)
root := t.TempDir()
lib := model.Library{Name: "Movies", Path: root, Type: "movie", Enabled: true}
if err := repos.Library.Create(t.Context(), &lib); err != nil {
t.Fatal(err)
}
file := filepath.Join(root, "Already In Library (2024).mkv")
if err := os.WriteFile(file, []byte("same-size"), 0o644); err != nil {
t.Fatal(err)
}
first, err := sc.ScanLibrary(t.Context(), lib.ID)
if err != nil {
t.Fatalf("first scan: %v", err)
}
if first.Added != 1 || first.Skipped != 0 {
t.Fatalf("first scan = %#v, want added=1 skipped=0", first)
}
second, err := sc.ScanLibrary(t.Context(), lib.ID)
if err != nil {
t.Fatalf("second scan: %v", err)
}
if second.Added != 0 || second.Updated != 0 || second.Skipped != 1 {
t.Fatalf("second scan = %#v, want unchanged file skipped", second)
}
if got := countMedia(t, repos); got != 1 {
t.Fatalf("media count = %d, want 1", got)
}
}
func TestScanLibraryReportsPerFileUpsertErrors(t *testing.T) {
sc, repos := newScannerTestEnv(t)
root := t.TempDir()
lib := model.Library{Name: "Broken DB", Path: root, Type: "movie", Enabled: true}
if err := repos.Library.Create(t.Context(), &lib); err != nil {
t.Fatal(err)
}
file := filepath.Join(root, "Cannot Insert (2024).mkv")
if err := os.WriteFile(file, []byte("data"), 0o644); err != nil {
t.Fatal(err)
}
if err := repos.DB.Exec("DROP TABLE media").Error; err != nil {
t.Fatal(err)
}
res, err := sc.ScanLibrary(t.Context(), lib.ID)
if err != nil {
t.Fatalf("scan should continue and report file errors, got top-level error: %v", err)
}
if res.ErrorCount != 1 || len(res.Errors) != 1 {
t.Fatalf("scan errors = count %d details %#v, want one detailed error", res.ErrorCount, res.Errors)
}
if res.Visited != 1 {
t.Fatalf("visited = %d, want 1", res.Visited)
}
}
func TestScanLibraryMapsPersistedHostLibraryPath(t *testing.T) {
sc, repos := newScannerTestEnv(t)
root := t.TempDir()
@@ -132,6 +189,8 @@ func TestScanLibraryMapsPersistedHostLibraryPath(t *testing.T) {
func TestRemovePathDeletesVanishedMedia(t *testing.T) {
sc, repos := newScannerTestEnv(t)
cache := NewRuntimeCacheService(&config.Config{}, zap.NewNop())
sc.SetRuntimeCache(cache)
root := t.TempDir()
lib := model.Library{Name: "Movies", Path: root, Type: "movie", Enabled: true}
if err := repos.Library.Create(t.Context(), &lib); err != nil {
@@ -147,10 +206,18 @@ func TestRemovePathDeletesVanishedMedia(t *testing.T) {
if countMedia(t, repos) != 1 {
t.Fatal("expected 1 media before removal")
}
cache.SetJSON(t.Context(), "media:list:test", map[string]string{"state": "stale"}, time.Minute)
var cached map[string]string
if !cache.GetJSON(t.Context(), "media:list:test", &cached) {
t.Fatal("expected media cache to be primed")
}
// A still-present file is not removed.
if removed, _ := sc.RemovePath(t.Context(), file); removed != 0 {
t.Fatalf("present file should not be removed, got %d", removed)
}
if !cache.GetJSON(t.Context(), "media:list:test", &cached) {
t.Fatal("present file should not invalidate media cache")
}
if err := os.Remove(file); err != nil {
t.Fatal(err)
}
@@ -164,6 +231,9 @@ func TestRemovePathDeletesVanishedMedia(t *testing.T) {
if countMedia(t, repos) != 0 {
t.Fatal("expected 0 media after removal")
}
if cache.GetJSON(t.Context(), "media:list:test", &cached) {
t.Fatal("vanished media removal should invalidate media cache")
}
}
// TestScanSkipsHardlinkDuplicate verifies that a hardlink (same inode) kept
+16 -2
View File
@@ -74,6 +74,7 @@ type Container struct {
Notify *NotifyService
Site *SiteService
Device *DeviceService
Cache *RuntimeCacheService
stopCtx context.Context
stopCancel context.CancelFunc
@@ -92,6 +93,13 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont
go sseHub.Run()
probe := NewFFprobeService(cfg, log)
runtimeCache := NewRuntimeCacheService(cfg, log)
if searchBackend := repository.NewOpenSearchMediaBackend(cfg.Search); searchBackend != nil && repos != nil && repos.Media != nil {
repos.Media.SetSearchBackend(searchBackend)
if log != nil {
log.Info("opensearch media search enabled", zap.String("index", cfg.Search.Index), zap.String("url", cfg.Search.OpenSearchURL))
}
}
crypto := NewCryptoService(cfg.Secrets.JWTSecret, log)
apiConfig := NewAPIConfigService(log, repos, crypto)
tmdb := NewTMDbProvider(cfg, log, apiConfig)
@@ -108,6 +116,7 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont
discover := NewDiscoverService(log, tmdb)
transcoder := NewTranscoderService(cfg, log, repos, hub)
scanner := NewScannerService(cfg, log, repos, hub, probe, scraper)
scanner.SetRuntimeCache(runtimeCache)
organizePipeline := NewOrganizePipelineService(log, repos, organizer, scanner, tasks)
watcher := NewWatcherService(log, repos, scanner)
nfo := NewNFOService(log, repos)
@@ -125,6 +134,7 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont
storageCfg := NewStorageConfigService(log, repos, crypto)
strmSvc := NewSTRMService(log, repos, cfg)
scanner.SetStorageConfig(storageCfg)
emby.SetRuntimeCache(runtimeCache)
emby.SetCloudProbe(storageCfg, probe)
downloadClients := NewDownloadClientService(log, repos)
assistant := NewAssistantService(log, repos, ai)
@@ -188,7 +198,7 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont
SSEHub: sseHub,
Tasks: tasks,
Auth: authSvc,
Media: NewMediaService(cfg, log, repos),
Media: NewMediaService(cfg, log, repos).SetRuntimeCache(runtimeCache),
Scan: scanner,
Stream: NewStreamService(cfg, log, repos, transcoder),
Transcoder: transcoder,
@@ -205,7 +215,7 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont
Downloads: downloads,
Subscription: subscription,
Subtitle: NewSubtitleService(log, repos),
Stats: NewStatsService(log, repos),
Stats: NewStatsService(log, repos).SetRuntimeCache(runtimeCache),
Profile: NewProfileService(log, repos),
Audit: NewAuditService(log, repos),
NFO: nfo,
@@ -237,6 +247,7 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont
Notify: notifySvc,
Site: siteSvc,
Device: deviceSvc,
Cache: runtimeCache,
stopCtx: ctx,
stopCancel: cancel,
}
@@ -444,6 +455,9 @@ func (c *Container) Close() {
if c.Transcoder != nil {
c.Transcoder.StopAll()
}
if c.Cache != nil {
_ = c.Cache.Close()
}
if c.WSHub != nil {
c.WSHub.Stop()
}
+24 -2
View File
@@ -23,8 +23,9 @@ import (
// StatsService computes aggregate stats.
type StatsService struct {
log *zap.Logger
repo *repository.Container
log *zap.Logger
repo *repository.Container
cache *RuntimeCacheService
}
// NewStatsService is the constructor.
@@ -32,6 +33,13 @@ func NewStatsService(log *zap.Logger, repo *repository.Container) *StatsService
return &StatsService{log: log, repo: repo}
}
func (s *StatsService) SetRuntimeCache(cache *RuntimeCacheService) *StatsService {
if s != nil {
s.cache = cache
}
return s
}
// Snapshot is the JSON returned by /api/stats.
type Snapshot struct {
Libraries int64 `json:"libraries"`
@@ -57,6 +65,15 @@ type Hardware struct {
// Compute builds a fresh snapshot.
func (s *StatsService) Compute(ctx context.Context, dataDir string) (*Snapshot, error) {
const cacheKey = "stats:snapshot:base"
if s.cache != nil {
var cached Snapshot
if s.cache.GetJSON(ctx, cacheKey, &cached) {
cached.GeneratedAt = time.Now()
cached.Hardware = readHardware(dataDir)
return &cached, nil
}
}
snap := &Snapshot{GeneratedAt: time.Now()}
libs, err := s.repo.Library.List(ctx)
if err != nil {
@@ -114,6 +131,11 @@ func (s *StatsService) Compute(ctx context.Context, dataDir string) (*Snapshot,
return nil, err
}
if s.cache != nil {
cacheCopy := *snap
cacheCopy.Hardware = Hardware{}
s.cache.SetJSON(ctx, cacheKey, cacheCopy, 10*time.Second)
}
snap.Hardware = readHardware(dataDir)
return snap, nil
}
+5 -6
View File
@@ -19,6 +19,8 @@ import (
"strings"
)
var linkFile = os.Link
// TransferMode 表示整理时文件的转移方式。
type TransferMode string
@@ -62,14 +64,11 @@ func transferFile(src, dst string, mode TransferMode) error {
case TransferCopy:
return copyFile(src, dst)
case TransferHardlink:
if err := os.Link(src, dst); err != nil {
if err := linkFile(src, dst); err != nil {
// Docker 部署里下载目录和媒体目录往往是两个独立的 bind mount,
// 即使在宿主机上同属一块盘,容器内 os.Link 也会因跨文件系统
// (EXDEV) 失败。此前直接报错导致 PT 下载完成后整理静默中断;
// 现在自动降级为复制(保留源文件继续做种,语义一致)。
if copyErr := copyFile(src, dst); copyErr == nil {
return nil
}
// (EXDEV) 失败。hardlink 模式必须保持零额外数据占用语义,不能
// 自动降级为复制;需要复制时请显式选择 copy。
return fmt.Errorf("hardlink failed: %w; source and target must be on the same filesystem, choose copy if you want to duplicate data", err)
}
return nil
+30
View File
@@ -1,8 +1,10 @@
package service
import (
"errors"
"os"
"path/filepath"
"strings"
"testing"
)
@@ -74,6 +76,34 @@ func TestTransferFileHardlinkSharesInodeAndKeepsSource(t *testing.T) {
}
}
func TestTransferFileHardlinkDoesNotFallBackToCopy(t *testing.T) {
dir := t.TempDir()
src := writeTemp(t, dir, "src.mkv", "payload")
dst := filepath.Join(dir, "dst.mkv")
origLinkFile := linkFile
linkFile = func(_, _ string) error {
return errors.New("simulated cross-device link")
}
t.Cleanup(func() {
linkFile = origLinkFile
})
err := transferFile(src, dst, TransferHardlink)
if err == nil {
t.Fatal("hardlink failure should be reported instead of falling back to copy")
}
if !strings.Contains(err.Error(), "hardlink failed") {
t.Fatalf("hardlink error = %q, want hardlink failure context", err.Error())
}
if _, statErr := os.Stat(dst); !os.IsNotExist(statErr) {
t.Fatalf("hardlink failure should not create copied dst, stat err = %v", statErr)
}
if b, readErr := os.ReadFile(src); readErr != nil || string(b) != "payload" {
t.Fatalf("hardlink failure should keep source unchanged, content=%q err=%v", b, readErr)
}
}
func TestTransferFileSymlinkKeepsSource(t *testing.T) {
dir := t.TempDir()
src := writeTemp(t, dir, "src.mkv", "payload")