mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-07 05:46:38 +08:00
fix: cache cloud artwork before library import
This commit is contained in:
@@ -284,6 +284,11 @@ func cloudPlayHandler(svc *service.Container) gin.HandlerFunc {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func serveCloudResolvedLink(svc *service.Container, c *gin.Context, typ, ref string) {
|
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 {
|
if svc == nil || svc.StorageCfg == nil {
|
||||||
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "cloud storage service unavailable"})
|
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "cloud storage service unavailable"})
|
||||||
return
|
return
|
||||||
|
|||||||
+118
-39
@@ -236,6 +236,77 @@ func serveCachedPlaceholder(w http.ResponseWriter) {
|
|||||||
_, _ = w.Write(transparent1x1PNG)
|
_, _ = 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
|
// Serve writes the requested image to w. Caller is expected to validate
|
||||||
// the JWT before invoking it.
|
// the JWT before invoking it.
|
||||||
func (p *ImageProxy) Serve(ctx context.Context, w http.ResponseWriter, r *http.Request, raw string) error {
|
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 == "" {
|
if stableKey == "" {
|
||||||
stableKey = link.URL
|
stableKey = link.URL
|
||||||
}
|
}
|
||||||
sum := sha1.Sum([]byte("cloud-image:" + stableKey))
|
key, cachePath, failPath := p.cloudImageCachePaths(stableKey)
|
||||||
key := "cloud-" + hex.EncodeToString(sum[:])
|
|
||||||
cachePath := filepath.Join(p.cacheDir, key)
|
|
||||||
failPath := cachePath + ".fail"
|
|
||||||
|
|
||||||
if data, err := os.ReadFile(cachePath); err == nil && len(data) > 0 {
|
if serveCachedImageFile(w, r, key, cachePath) {
|
||||||
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 nil
|
return nil
|
||||||
}
|
}
|
||||||
if stat, err := os.Stat(failPath); err == nil && time.Since(stat.ModTime()) < imageNegativeCacheTTL {
|
if freshNegativeImageCache(failPath) {
|
||||||
serveCachedPlaceholder(w)
|
serveCachedPlaceholder(w)
|
||||||
return nil
|
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))
|
data, ctype, err := p.fetchAndCacheCloudImage(ctx, stableKey, link, r.UserAgent())
|
||||||
servePlaceholder(w)
|
if err != nil {
|
||||||
|
p.log.Warn("imageproxy: cloud image fetch failed", zap.String("url", link.URL), zap.Error(err))
|
||||||
|
serveCachedPlaceholder(w)
|
||||||
return nil
|
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)
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, link.URL, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
p.log.Warn("imageproxy: build cloud image request failed", zap.Error(err))
|
return nil, "", err
|
||||||
servePlaceholder(w)
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
for k, v := range link.Headers {
|
for k, v := range link.Headers {
|
||||||
req.Header.Set(k, v)
|
req.Header.Set(k, v)
|
||||||
}
|
}
|
||||||
if req.Header.Get("User-Agent") == "" {
|
if req.Header.Get("User-Agent") == "" {
|
||||||
if ua := r.UserAgent(); ua != "" {
|
if strings.TrimSpace(userAgent) != "" {
|
||||||
req.Header.Set("User-Agent", ua)
|
req.Header.Set("User-Agent", userAgent)
|
||||||
} else {
|
} else {
|
||||||
req.Header.Set("User-Agent", "MediaStationGo/0.1")
|
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)
|
resp, err := p.client.Do(req)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
p.log.Warn("imageproxy: cloud image fetch failed", zap.String("url", link.URL), zap.Error(err))
|
|
||||||
p.markImageFetchFailed(failPath)
|
p.markImageFetchFailed(failPath)
|
||||||
serveCachedPlaceholder(w)
|
return nil, "", err
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
defer resp.Body.Close()
|
defer resp.Body.Close()
|
||||||
if resp.StatusCode >= 400 {
|
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)
|
p.markImageFetchFailed(failPath)
|
||||||
serveCachedPlaceholder(w)
|
return nil, "", errors.New("cloud image returned " + resp.Status)
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
data, err := io.ReadAll(io.LimitReader(resp.Body, 32<<20))
|
data, err := io.ReadAll(io.LimitReader(resp.Body, 32<<20))
|
||||||
if err != nil || len(data) == 0 {
|
if err != nil {
|
||||||
p.log.Warn("imageproxy: read cloud image body failed", zap.String("url", link.URL), zap.Error(err))
|
|
||||||
p.markImageFetchFailed(failPath)
|
p.markImageFetchFailed(failPath)
|
||||||
serveCachedPlaceholder(w)
|
return nil, "", err
|
||||||
return nil
|
}
|
||||||
|
if len(data) == 0 {
|
||||||
|
p.markImageFetchFailed(failPath)
|
||||||
|
return nil, "", errors.New("cloud image body is empty")
|
||||||
}
|
}
|
||||||
|
|
||||||
p.mu.Lock()
|
p.mu.Lock()
|
||||||
@@ -482,10 +564,7 @@ func (p *ImageProxy) ServeCloudResolved(ctx context.Context, w http.ResponseWrit
|
|||||||
if ctype == "" {
|
if ctype == "" {
|
||||||
ctype = detectContentType(data)
|
ctype = detectContentType(data)
|
||||||
}
|
}
|
||||||
w.Header().Set("Content-Type", ctype)
|
return data, ctype, nil
|
||||||
w.Header().Set("Cache-Control", imageBrowserCacheControl)
|
|
||||||
http.ServeContent(w, r, key, time.Now(), bytes.NewReader(data))
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (p *ImageProxy) markImageFetchFailed(failPath string) {
|
func (p *ImageProxy) markImageFetchFailed(failPath string) {
|
||||||
|
|||||||
@@ -131,6 +131,9 @@ func TestImageProxyCachesCloudResolvedImage(t *testing.T) {
|
|||||||
})}
|
})}
|
||||||
|
|
||||||
link := &cloud.DirectLink{URL: "http://cloud-provider.invalid/poster.png"}
|
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++ {
|
for i := 0; i < 2; i++ {
|
||||||
rec := httptest.NewRecorder()
|
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 {
|
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 {
|
if got := atomic.LoadInt32(&calls); got != 1 {
|
||||||
t.Fatalf("upstream calls = %d, want 1 due to cloud image cache", got)
|
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)
|
type imageRoundTripFunc func(*http.Request) (*http.Response, error)
|
||||||
|
|||||||
+167
-7
@@ -56,9 +56,15 @@ type ScannerService struct {
|
|||||||
scraper *ScraperService
|
scraper *ScraperService
|
||||||
storage *StorageConfigService
|
storage *StorageConfigService
|
||||||
|
|
||||||
cloudScanMu sync.Mutex
|
imageProxy *ImageProxy
|
||||||
cloudScans map[string]*cloudScanEntry
|
|
||||||
cloudSlots chan struct{}
|
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.
|
// NewScannerService is the constructor.
|
||||||
@@ -72,9 +78,12 @@ func NewScannerService(
|
|||||||
) *ScannerService {
|
) *ScannerService {
|
||||||
return &ScannerService{
|
return &ScannerService{
|
||||||
cfg: cfg, log: log, repo: repo, hub: hub,
|
cfg: cfg, log: log, repo: repo, hub: hub,
|
||||||
probe: probe, scraper: scraper,
|
probe: probe,
|
||||||
cloudScans: make(map[string]*cloudScanEntry),
|
scraper: scraper,
|
||||||
cloudSlots: make(chan struct{}, 1),
|
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
|
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.
|
// ScanResult summarises a scan run.
|
||||||
type ScanResult struct {
|
type ScanResult struct {
|
||||||
LibraryID string `json:"library_id"`
|
LibraryID string `json:"library_id"`
|
||||||
@@ -126,6 +275,12 @@ type cloudScanEntry struct {
|
|||||||
cancel context.CancelFunc
|
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) {
|
func (s *ScannerService) beginCloudScan(ctx context.Context, lib *model.Library, mount CloudMountInfo) (context.Context, func(*ScanResult, error), error) {
|
||||||
if s == nil || lib == nil {
|
if s == nil || lib == nil {
|
||||||
return ctx, func(*ScanResult, error) {}, 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)
|
publishProgress("listing", dirsVisited == 1 || dirsVisited%20 == 0)
|
||||||
sidecars := newCloudSidecarSet(typ, entries)
|
sidecars := newCloudSidecarSet(typ, entries)
|
||||||
dirMeta := s.cloudDirectoryMetadata(ctx, typ, displayDir, sidecars, inheritedMeta)
|
dirMeta := s.cloudDirectoryMetadata(ctx, typ, displayDir, sidecars, inheritedMeta)
|
||||||
|
s.cacheCloudMetadataArtworkNow(ctx, dirMeta)
|
||||||
for _, entry := range entries {
|
for _, entry := range entries {
|
||||||
select {
|
select {
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
@@ -654,12 +810,14 @@ func (s *ScannerService) scanCloudLibrary(ctx context.Context, lib *model.Librar
|
|||||||
publishProgress("listing", filesDiscovered%100 == 0)
|
publishProgress("listing", filesDiscovered%100 == 0)
|
||||||
displayPath := joinCloudDisplayPath(displayDir, entry.Name)
|
displayPath := joinCloudDisplayPath(displayDir, entry.Name)
|
||||||
path := cloudMediaPath(typ, displayPath)
|
path := cloudMediaPath(typ, displayPath)
|
||||||
|
localMeta := s.cloudFileMetadata(ctx, typ, displayPath, entry.Name, sidecars, dirMeta, librarySupportsSeasons(lib))
|
||||||
|
s.cacheCloudMetadataArtworkNow(ctx, localMeta)
|
||||||
candidate := cloudCandidate{
|
candidate := cloudCandidate{
|
||||||
ref: ref,
|
ref: ref,
|
||||||
name: entry.Name,
|
name: entry.Name,
|
||||||
size: entry.Size,
|
size: entry.Size,
|
||||||
path: path,
|
path: path,
|
||||||
localMeta: s.cloudFileMetadata(ctx, typ, displayPath, entry.Name, sidecars, dirMeta, librarySupportsSeasons(lib)),
|
localMeta: localMeta,
|
||||||
}
|
}
|
||||||
key := cloudMediaDedupeKey(lib, displayDir, entry.Name, entry.Size)
|
key := cloudMediaDedupeKey(lib, displayDir, entry.Name, entry.Size)
|
||||||
if key != "" {
|
if key != "" {
|
||||||
@@ -800,6 +958,8 @@ func (s *ScannerService) ingestCloudFile(ctx context.Context, lib *model.Library
|
|||||||
if localMeta != nil {
|
if localMeta != nil {
|
||||||
applyLocalMetadata(m, localMeta)
|
applyLocalMetadata(m, localMeta)
|
||||||
res.LocalMetadata++
|
res.LocalMetadata++
|
||||||
|
s.queueCloudArtworkPrefetch(localMeta.PosterURL)
|
||||||
|
s.queueCloudArtworkPrefetch(localMeta.BackdropURL)
|
||||||
}
|
}
|
||||||
if err := s.repo.Media.Upsert(ctx, m); err != nil {
|
if err := s.repo.Media.Upsert(ctx, m); err != nil {
|
||||||
s.log.Warn("upsert cloud media failed", zap.String("path", path), zap.Error(err))
|
s.log.Warn("upsert cloud media failed", zap.String("path", path), zap.Error(err))
|
||||||
|
|||||||
@@ -422,6 +422,9 @@ func TestScanCloudLibraryReadsRemoteNFOAndArtwork(t *testing.T) {
|
|||||||
_, _ = w.Write([]byte(`<tvshow><title>剑来</title><year>2024</year><plot>天地有剑气</plot></tvshow>`))
|
_, _ = w.Write([]byte(`<tvshow><title>剑来</title><year>2024</year><plot>天地有剑气</plot></tvshow>`))
|
||||||
case "/dav/Anime/JianLai/Season1/JianLai.S01E01.nfo":
|
case "/dav/Anime/JianLai/Season1/JianLai.S01E01.nfo":
|
||||||
_, _ = w.Write([]byte(`<episodedetails><showtitle>剑来</showtitle><title>第一集</title><season>1</season><episode>1</episode></episodedetails>`))
|
_, _ = w.Write([]byte(`<episodedetails><showtitle>剑来</showtitle><title>第一集</title><season>1</season><episode>1</episode></episodedetails>`))
|
||||||
|
case "/dav/Anime/JianLai/poster.jpg":
|
||||||
|
w.Header().Set("Content-Type", "image/jpeg")
|
||||||
|
_, _ = w.Write([]byte("cloud-poster-bytes"))
|
||||||
default:
|
default:
|
||||||
t.Fatalf("unexpected get path %s", r.URL.Path)
|
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 := NewScannerService(&config.Config{}, log, repos, NewHub(log), nil, nil)
|
||||||
scanner.SetStorageConfig(storage)
|
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)
|
res, err := scanner.ScanLibrary(t.Context(), lib.ID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -476,7 +481,27 @@ func TestScanCloudLibraryReadsRemoteNFOAndArtwork(t *testing.T) {
|
|||||||
if media.PosterURL != "/api/cloud/play/openlist?ref=%2FAnime%2FJianLai%2Fposter.jpg" {
|
if media.PosterURL != "/api/cloud/play/openlist?ref=%2FAnime%2FJianLai%2Fposter.jpg" {
|
||||||
t.Fatalf("poster url = %q", media.PosterURL)
|
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" {
|
if media.ScrapeStatus != "matched" {
|
||||||
t.Fatalf("scrape status = %q", media.ScrapeStatus)
|
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")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -162,6 +162,7 @@ func New(cfg *config.Config, log *zap.Logger, repos *repository.Container) *Cont
|
|||||||
}
|
}
|
||||||
return roots
|
return roots
|
||||||
})
|
})
|
||||||
|
scanner.SetImageProxy(imageProxy)
|
||||||
|
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user