diff --git a/internal/handler/cloud.go b/internal/handler/cloud.go index 296e983..f100bd1 100644 --- a/internal/handler/cloud.go +++ b/internal/handler/cloud.go @@ -284,6 +284,11 @@ func cloudPlayHandler(svc *service.Container) gin.HandlerFunc { } func serveCloudResolvedLink(svc *service.Container, c *gin.Context, typ, ref string) { + if isCloudImageRef(ref) && svc != nil && svc.ImageProxy != nil { + if svc.ImageProxy.ServeCloudCached(c.Writer, c.Request, typ+":"+ref) { + return + } + } if svc == nil || svc.StorageCfg == nil { c.JSON(http.StatusServiceUnavailable, gin.H{"error": "cloud storage service unavailable"}) return diff --git a/internal/service/image_proxy.go b/internal/service/image_proxy.go index a909310..b95d1c4 100644 --- a/internal/service/image_proxy.go +++ b/internal/service/image_proxy.go @@ -236,6 +236,77 @@ func serveCachedPlaceholder(w http.ResponseWriter) { _, _ = w.Write(transparent1x1PNG) } +func (p *ImageProxy) cloudImageCachePaths(stableKey string) (string, string, string) { + stableKey = strings.TrimSpace(stableKey) + if stableKey == "" { + stableKey = "unknown" + } + sum := sha1.Sum([]byte("cloud-image:" + stableKey)) + key := "cloud-" + hex.EncodeToString(sum[:]) + cachePath := filepath.Join(p.cacheDir, key) + return key, cachePath, cachePath + ".fail" +} + +func serveCachedImageFile(w http.ResponseWriter, r *http.Request, key, cachePath string) bool { + data, err := os.ReadFile(cachePath) + if err != nil || len(data) == 0 { + return false + } + w.Header().Set("Content-Type", detectContentType(data)) + w.Header().Set("Cache-Control", imageBrowserCacheControl) + stat, _ := os.Stat(cachePath) + modTime := time.Now() + if stat != nil { + modTime = stat.ModTime() + } + http.ServeContent(w, r, key, modTime, bytes.NewReader(data)) + return true +} + +func freshNegativeImageCache(failPath string) bool { + stat, err := os.Stat(failPath) + if err != nil { + return false + } + if time.Since(stat.ModTime()) < imageNegativeCacheTTL { + return true + } + _ = os.Remove(failPath) + return false +} + +// CloudImageCached reports whether a stable cloud-image ref already has a +// usable positive or short-lived negative cache entry. Scanner pre-warm uses it +// to avoid repeatedly resolving the same cloud sidecar image. +func (p *ImageProxy) CloudImageCached(stableKey string) bool { + if p == nil { + return false + } + _, cachePath, failPath := p.cloudImageCachePaths(stableKey) + if stat, err := os.Stat(cachePath); err == nil && stat.Size() > 0 { + return true + } + return freshNegativeImageCache(failPath) +} + +// ServeCloudCached serves an already-local cloud sidecar image without asking +// the cloud provider for a fresh direct link. It returns true when it wrote a +// response, including a fresh negative-cache placeholder. +func (p *ImageProxy) ServeCloudCached(w http.ResponseWriter, r *http.Request, stableKey string) bool { + if p == nil { + return false + } + key, cachePath, failPath := p.cloudImageCachePaths(stableKey) + if serveCachedImageFile(w, r, key, cachePath) { + return true + } + if freshNegativeImageCache(failPath) { + serveCachedPlaceholder(w) + return true + } + return false +} + // Serve writes the requested image to w. Caller is expected to validate // the JWT before invoking it. func (p *ImageProxy) Serve(ctx context.Context, w http.ResponseWriter, r *http.Request, raw string) error { @@ -392,46 +463,60 @@ func (p *ImageProxy) ServeCloudResolved(ctx context.Context, w http.ResponseWrit if stableKey == "" { stableKey = link.URL } - sum := sha1.Sum([]byte("cloud-image:" + stableKey)) - key := "cloud-" + hex.EncodeToString(sum[:]) - cachePath := filepath.Join(p.cacheDir, key) - failPath := cachePath + ".fail" + key, cachePath, failPath := p.cloudImageCachePaths(stableKey) - if data, err := os.ReadFile(cachePath); err == nil && len(data) > 0 { - w.Header().Set("Content-Type", detectContentType(data)) - w.Header().Set("Cache-Control", imageBrowserCacheControl) - stat, _ := os.Stat(cachePath) - modTime := time.Now() - if stat != nil { - modTime = stat.ModTime() - } - http.ServeContent(w, r, key, modTime, bytes.NewReader(data)) + if serveCachedImageFile(w, r, key, cachePath) { return nil } - if stat, err := os.Stat(failPath); err == nil && time.Since(stat.ModTime()) < imageNegativeCacheTTL { + if freshNegativeImageCache(failPath) { serveCachedPlaceholder(w) return nil - } else if err == nil { - _ = os.Remove(failPath) } - if err := os.MkdirAll(p.cacheDir, 0o755); err != nil { - p.log.Warn("imageproxy: mkdir failed", zap.String("dir", p.cacheDir), zap.Error(err)) - servePlaceholder(w) + + data, ctype, err := p.fetchAndCacheCloudImage(ctx, stableKey, link, r.UserAgent()) + if err != nil { + p.log.Warn("imageproxy: cloud image fetch failed", zap.String("url", link.URL), zap.Error(err)) + serveCachedPlaceholder(w) return nil } + w.Header().Set("Content-Type", ctype) + w.Header().Set("Cache-Control", imageBrowserCacheControl) + http.ServeContent(w, r, key, time.Now(), bytes.NewReader(data)) + return nil +} + +// PrefetchCloudResolved downloads a cloud sidecar image into the local cache +// without writing an HTTP response. It is intentionally best-effort; callers +// should queue it with low concurrency so large cloud libraries do not overload +// small NAS devices. +func (p *ImageProxy) PrefetchCloudResolved(ctx context.Context, stableKey string, link *cloud.DirectLink) error { + if p == nil || link == nil || strings.TrimSpace(link.URL) == "" { + return nil + } + if p.CloudImageCached(stableKey) { + return nil + } + _, _, err := p.fetchAndCacheCloudImage(ctx, stableKey, link, "MediaStationGo/0.1") + return err +} + +func (p *ImageProxy) fetchAndCacheCloudImage(ctx context.Context, stableKey string, link *cloud.DirectLink, userAgent string) ([]byte, string, error) { + if err := os.MkdirAll(p.cacheDir, 0o755); err != nil { + return nil, "", err + } + _, cachePath, failPath := p.cloudImageCachePaths(stableKey) + req, err := http.NewRequestWithContext(ctx, http.MethodGet, link.URL, nil) if err != nil { - p.log.Warn("imageproxy: build cloud image request failed", zap.Error(err)) - servePlaceholder(w) - return nil + return nil, "", err } for k, v := range link.Headers { req.Header.Set(k, v) } if req.Header.Get("User-Agent") == "" { - if ua := r.UserAgent(); ua != "" { - req.Header.Set("User-Agent", ua) + if strings.TrimSpace(userAgent) != "" { + req.Header.Set("User-Agent", userAgent) } else { req.Header.Set("User-Agent", "MediaStationGo/0.1") } @@ -440,25 +525,22 @@ func (p *ImageProxy) ServeCloudResolved(ctx context.Context, w http.ResponseWrit resp, err := p.client.Do(req) if err != nil { - p.log.Warn("imageproxy: cloud image fetch failed", zap.String("url", link.URL), zap.Error(err)) p.markImageFetchFailed(failPath) - serveCachedPlaceholder(w) - return nil + return nil, "", err } defer resp.Body.Close() if resp.StatusCode >= 400 { - p.log.Warn("imageproxy: cloud image returned non-OK", - zap.String("url", link.URL), zap.String("status", resp.Status)) p.markImageFetchFailed(failPath) - serveCachedPlaceholder(w) - return nil + return nil, "", errors.New("cloud image returned " + resp.Status) } data, err := io.ReadAll(io.LimitReader(resp.Body, 32<<20)) - if err != nil || len(data) == 0 { - p.log.Warn("imageproxy: read cloud image body failed", zap.String("url", link.URL), zap.Error(err)) + if err != nil { p.markImageFetchFailed(failPath) - serveCachedPlaceholder(w) - return nil + return nil, "", err + } + if len(data) == 0 { + p.markImageFetchFailed(failPath) + return nil, "", errors.New("cloud image body is empty") } p.mu.Lock() @@ -482,10 +564,7 @@ func (p *ImageProxy) ServeCloudResolved(ctx context.Context, w http.ResponseWrit if ctype == "" { ctype = detectContentType(data) } - w.Header().Set("Content-Type", ctype) - w.Header().Set("Cache-Control", imageBrowserCacheControl) - http.ServeContent(w, r, key, time.Now(), bytes.NewReader(data)) - return nil + return data, ctype, nil } func (p *ImageProxy) markImageFetchFailed(failPath string) { diff --git a/internal/service/image_proxy_test.go b/internal/service/image_proxy_test.go index 97ca8ee..e7cc9fb 100644 --- a/internal/service/image_proxy_test.go +++ b/internal/service/image_proxy_test.go @@ -131,6 +131,9 @@ func TestImageProxyCachesCloudResolvedImage(t *testing.T) { })} link := &cloud.DirectLink{URL: "http://cloud-provider.invalid/poster.png"} + if proxy.ServeCloudCached(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "/api/cloud/play/openlist?ref=poster.png", nil), "openlist:poster.png") { + t.Fatal("ServeCloudCached returned true before the cloud image was cached") + } for i := 0; i < 2; i++ { rec := httptest.NewRecorder() if err := proxy.ServeCloudResolved(t.Context(), rec, httptest.NewRequest(http.MethodGet, "/api/cloud/play/openlist?ref=poster.png", nil), "openlist:poster.png", link); err != nil { @@ -146,6 +149,50 @@ func TestImageProxyCachesCloudResolvedImage(t *testing.T) { if got := atomic.LoadInt32(&calls); got != 1 { t.Fatalf("upstream calls = %d, want 1 due to cloud image cache", got) } + + rec := httptest.NewRecorder() + if !proxy.ServeCloudCached(rec, httptest.NewRequest(http.MethodGet, "/api/cloud/play/openlist?ref=poster.png", nil), "openlist:poster.png") { + t.Fatal("ServeCloudCached returned false after the cloud image was cached") + } + if got := atomic.LoadInt32(&calls); got != 1 { + t.Fatalf("upstream calls after ServeCloudCached = %d, want 1", got) + } + if rec.Body.Len() != len(transparent1x1PNG) { + t.Fatalf("cached body length = %d, want %d", rec.Body.Len(), len(transparent1x1PNG)) + } +} + +func TestImageProxyPrefetchCloudResolvedImage(t *testing.T) { + var calls int32 + proxy := NewImageProxy(&config.Config{Cache: config.CacheConfig{CacheDir: filepath.Join(t.TempDir(), "cache")}}, zap.NewNop()) + proxy.client = &http.Client{Transport: imageRoundTripFunc(func(req *http.Request) (*http.Response, error) { + atomic.AddInt32(&calls, 1) + return &http.Response{ + StatusCode: http.StatusOK, + Status: "200 OK", + Header: http.Header{"Content-Type": []string{"image/png"}}, + Body: io.NopCloser(strings.NewReader(string(transparent1x1PNG))), + Request: req, + }, nil + })} + + link := &cloud.DirectLink{URL: "http://cloud-provider.invalid/folder.png"} + if err := proxy.PrefetchCloudResolved(t.Context(), "openlist:folder.png", link); err != nil { + t.Fatal(err) + } + if err := proxy.PrefetchCloudResolved(t.Context(), "openlist:folder.png", link); err != nil { + t.Fatal(err) + } + if got := atomic.LoadInt32(&calls); got != 1 { + t.Fatalf("upstream calls = %d, want 1 after prefetch cache hit", got) + } + rec := httptest.NewRecorder() + if !proxy.ServeCloudCached(rec, httptest.NewRequest(http.MethodGet, "/api/cloud/play/openlist?ref=folder.png", nil), "openlist:folder.png") { + t.Fatal("prefetched cloud image was not served from cache") + } + if rec.Body.Len() != len(transparent1x1PNG) { + t.Fatalf("cached body length = %d, want %d", rec.Body.Len(), len(transparent1x1PNG)) + } } type imageRoundTripFunc func(*http.Request) (*http.Response, error) diff --git a/internal/service/scanner.go b/internal/service/scanner.go index 303bf71..9cbd07b 100644 --- a/internal/service/scanner.go +++ b/internal/service/scanner.go @@ -56,9 +56,15 @@ type ScannerService struct { scraper *ScraperService storage *StorageConfigService - cloudScanMu sync.Mutex - cloudScans map[string]*cloudScanEntry - cloudSlots chan struct{} + imageProxy *ImageProxy + + cloudScanMu sync.Mutex + cloudScans map[string]*cloudScanEntry + cloudSlots chan struct{} + cloudImagePrefetchOnce sync.Once + cloudImagePrefetchQueue chan cloudImagePrefetchTask + cloudImagePrefetchMu sync.Mutex + cloudImagePrefetching map[string]struct{} } // NewScannerService is the constructor. @@ -72,9 +78,12 @@ func NewScannerService( ) *ScannerService { return &ScannerService{ cfg: cfg, log: log, repo: repo, hub: hub, - probe: probe, scraper: scraper, - cloudScans: make(map[string]*cloudScanEntry), - cloudSlots: make(chan struct{}, 1), + probe: probe, + scraper: scraper, + cloudScans: make(map[string]*cloudScanEntry), + cloudSlots: make(chan struct{}, 1), + cloudImagePrefetchQueue: make(chan cloudImagePrefetchTask, 256), + cloudImagePrefetching: make(map[string]struct{}), } } @@ -85,6 +94,146 @@ func (s *ScannerService) SetStorageConfig(storage *StorageConfigService) { s.storage = storage } +// 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. +func (s *ScannerService) SetImageProxy(imageProxy *ImageProxy) { + s.imageProxy = imageProxy + if imageProxy != nil { + s.cloudImagePrefetchOnce.Do(func() { + go s.cloudImagePrefetchWorker() + }) + } +} + +func (s *ScannerService) cloudImagePrefetchWorker() { + for task := range s.cloudImagePrefetchQueue { + s.prefetchCloudImage(task) + } +} + +func (s *ScannerService) queueCloudArtworkPrefetch(raw string) { + if s == nil || s.storage == nil || s.imageProxy == nil { + return + } + typ, ref, ok := parseCloudImagePlaybackURL(raw) + if !ok { + return + } + stableKey := typ + ":" + ref + if s.imageProxy.CloudImageCached(stableKey) { + return + } + s.cloudImagePrefetchMu.Lock() + if _, ok := s.cloudImagePrefetching[stableKey]; ok { + s.cloudImagePrefetchMu.Unlock() + return + } + s.cloudImagePrefetching[stableKey] = struct{}{} + s.cloudImagePrefetchMu.Unlock() + + task := cloudImagePrefetchTask{typ: typ, ref: ref, stableKey: stableKey} + select { + case s.cloudImagePrefetchQueue <- task: + default: + s.cloudImagePrefetchMu.Lock() + delete(s.cloudImagePrefetching, stableKey) + s.cloudImagePrefetchMu.Unlock() + if s.log != nil { + s.log.Debug("cloud artwork prefetch queue full", zap.String("provider", typ), zap.String("ref", ref)) + } + } +} + +func (s *ScannerService) prefetchCloudImage(task cloudImagePrefetchTask) { + defer func() { + s.cloudImagePrefetchMu.Lock() + delete(s.cloudImagePrefetching, task.stableKey) + s.cloudImagePrefetchMu.Unlock() + }() + if s == nil || s.storage == nil || s.imageProxy == nil || s.imageProxy.CloudImageCached(task.stableKey) { + return + } + ctx, cancel := context.WithTimeout(context.Background(), 45*time.Second) + defer cancel() + link, err := s.storage.CloudResolve(ctx, task.typ, task.ref, "") + if err != nil { + if s.log != nil { + s.log.Debug("resolve cloud artwork for prefetch failed", zap.String("provider", task.typ), zap.String("ref", task.ref), zap.Error(err)) + } + return + } + if err := s.imageProxy.PrefetchCloudResolved(ctx, task.stableKey, link); err != nil && s.log != nil { + s.log.Debug("prefetch cloud artwork failed", zap.String("provider", task.typ), zap.String("ref", task.ref), zap.Error(err)) + } +} + +func (s *ScannerService) cacheCloudArtworkNow(ctx context.Context, raw string) { + if s == nil || s.storage == nil || s.imageProxy == nil { + return + } + typ, ref, ok := parseCloudImagePlaybackURL(raw) + if !ok { + return + } + stableKey := typ + ":" + ref + if s.imageProxy.CloudImageCached(stableKey) { + return + } + cacheCtx, cancel := context.WithTimeout(ctx, 20*time.Second) + defer cancel() + link, err := s.storage.CloudResolve(cacheCtx, typ, ref, "") + if err != nil { + if s.log != nil { + s.log.Debug("resolve cloud artwork for priority cache failed", zap.String("provider", typ), zap.String("ref", ref), zap.Error(err)) + } + s.queueCloudArtworkPrefetch(raw) + return + } + if err := s.imageProxy.PrefetchCloudResolved(cacheCtx, stableKey, link); err != nil { + if s.log != nil { + s.log.Debug("priority cache cloud artwork failed", zap.String("provider", typ), zap.String("ref", ref), zap.Error(err)) + } + s.queueCloudArtworkPrefetch(raw) + } +} + +func (s *ScannerService) cacheCloudMetadataArtworkNow(ctx context.Context, meta *LocalMetadata) { + if meta == nil { + return + } + s.cacheCloudArtworkNow(ctx, meta.PosterURL) + s.cacheCloudArtworkNow(ctx, meta.BackdropURL) +} + +func parseCloudImagePlaybackURL(raw string) (string, string, bool) { + u, err := url.Parse(strings.TrimSpace(raw)) + if err != nil { + return "", "", false + } + path := strings.Trim(u.Path, "/") + const prefix = "api/cloud/play/" + if !strings.HasPrefix(strings.ToLower(path), prefix) { + return "", "", false + } + typ := strings.TrimSpace(path[len(prefix):]) + ref := strings.TrimSpace(u.Query().Get("ref")) + if typ == "" || ref == "" || !isCloudArtworkRef(ref) { + return "", "", false + } + return typ, ref, true +} + +func isCloudArtworkRef(ref string) bool { + ref = strings.ToLower(strings.TrimSpace(ref)) + for _, suffix := range []string{".jpg", ".jpeg", ".png", ".webp", ".gif", ".bmp"} { + if strings.HasSuffix(ref, suffix) { + return true + } + } + return false +} + // ScanResult summarises a scan run. type ScanResult struct { LibraryID string `json:"library_id"` @@ -126,6 +275,12 @@ type cloudScanEntry struct { cancel context.CancelFunc } +type cloudImagePrefetchTask struct { + typ string + ref string + stableKey string +} + func (s *ScannerService) beginCloudScan(ctx context.Context, lib *model.Library, mount CloudMountInfo) (context.Context, func(*ScanResult, error), error) { if s == nil || lib == nil { return ctx, func(*ScanResult, error) {}, nil @@ -622,6 +777,7 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar publishProgress("listing", dirsVisited == 1 || dirsVisited%20 == 0) sidecars := newCloudSidecarSet(typ, entries) dirMeta := s.cloudDirectoryMetadata(ctx, typ, displayDir, sidecars, inheritedMeta) + s.cacheCloudMetadataArtworkNow(ctx, dirMeta) for _, entry := range entries { select { case <-ctx.Done(): @@ -654,12 +810,14 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar publishProgress("listing", filesDiscovered%100 == 0) displayPath := joinCloudDisplayPath(displayDir, entry.Name) path := cloudMediaPath(typ, displayPath) + localMeta := s.cloudFileMetadata(ctx, typ, displayPath, entry.Name, sidecars, dirMeta, librarySupportsSeasons(lib)) + s.cacheCloudMetadataArtworkNow(ctx, localMeta) candidate := cloudCandidate{ ref: ref, name: entry.Name, size: entry.Size, path: path, - localMeta: s.cloudFileMetadata(ctx, typ, displayPath, entry.Name, sidecars, dirMeta, librarySupportsSeasons(lib)), + localMeta: localMeta, } key := cloudMediaDedupeKey(lib, displayDir, entry.Name, entry.Size) if key != "" { @@ -800,6 +958,8 @@ func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library if localMeta != nil { applyLocalMetadata(m, localMeta) res.LocalMetadata++ + s.queueCloudArtworkPrefetch(localMeta.PosterURL) + s.queueCloudArtworkPrefetch(localMeta.BackdropURL) } if err := s.repo.Media.Upsert(ctx, m); err != nil { s.log.Warn("upsert cloud media failed", zap.String("path", path), zap.Error(err)) diff --git a/internal/service/scanner_cloud_test.go b/internal/service/scanner_cloud_test.go index 02823ed..ab71a40 100644 --- a/internal/service/scanner_cloud_test.go +++ b/internal/service/scanner_cloud_test.go @@ -422,6 +422,9 @@ func TestScanCloudLibraryReadsRemoteNFOAndArtwork(t *testing.T) { _, _ = w.Write([]byte(`剑来2024天地有剑气`)) case "/dav/Anime/JianLai/Season1/JianLai.S01E01.nfo": _, _ = w.Write([]byte(`剑来第一集11`)) + case "/dav/Anime/JianLai/poster.jpg": + w.Header().Set("Content-Type", "image/jpeg") + _, _ = w.Write([]byte("cloud-poster-bytes")) default: t.Fatalf("unexpected get path %s", r.URL.Path) } @@ -455,6 +458,8 @@ func TestScanCloudLibraryReadsRemoteNFOAndArtwork(t *testing.T) { } scanner := NewScannerService(&config.Config{}, log, repos, NewHub(log), nil, nil) scanner.SetStorageConfig(storage) + imageProxy := NewImageProxy(&config.Config{Cache: config.CacheConfig{CacheDir: t.TempDir()}}, log) + scanner.SetImageProxy(imageProxy) res, err := scanner.ScanLibrary(t.Context(), lib.ID) if err != nil { @@ -476,7 +481,27 @@ func TestScanCloudLibraryReadsRemoteNFOAndArtwork(t *testing.T) { if media.PosterURL != "/api/cloud/play/openlist?ref=%2FAnime%2FJianLai%2Fposter.jpg" { t.Fatalf("poster url = %q", media.PosterURL) } + rec := httptest.NewRecorder() + if !imageProxy.ServeCloudCached(rec, httptest.NewRequest(http.MethodGet, media.PosterURL, nil), "openlist:/Anime/JianLai/poster.jpg") { + t.Fatal("cloud poster should be cached locally during scan before media is exposed") + } + if got := rec.Body.String(); got != "cloud-poster-bytes" { + t.Fatalf("cached poster body = %q", got) + } if media.ScrapeStatus != "matched" { t.Fatalf("scrape status = %q", media.ScrapeStatus) } } + +func TestParseCloudImagePlaybackURL(t *testing.T) { + typ, ref, ok := parseCloudImagePlaybackURL("http://nas.local/api/cloud/play/openlist?ref=%2FAnime%2FJianLai%2Fposter.jpg") + if !ok || typ != "openlist" || ref != "/Anime/JianLai/poster.jpg" { + t.Fatalf("parse cloud image url = typ=%q ref=%q ok=%v", typ, ref, ok) + } + if _, _, ok := parseCloudImagePlaybackURL("/api/cloud/play/openlist?ref=%2FAnime%2FJianLai%2Fmovie.mkv"); ok { + t.Fatal("video cloud url should not be treated as artwork") + } + if _, _, ok := parseCloudImagePlaybackURL("https://image.tmdb.org/t/p/w500/poster.jpg"); ok { + t.Fatal("remote HTTP poster should not be treated as cloud artwork") + } +} diff --git a/internal/service/service.go b/internal/service/service.go index de1b60c..195b842 100644 --- a/internal/service/service.go +++ b/internal/service/service.go @@ -162,6 +162,7 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont } return roots }) + scanner.SetImageProxy(imageProxy) ctx, cancel := context.WithCancel(context.Background())