diff --git a/cmd/wire_gen.go b/cmd/wire_gen.go index fd04c44c5..6449e46fd 100644 --- a/cmd/wire_gen.go +++ b/cmd/wire_gen.go @@ -89,13 +89,13 @@ func CreateSubsonicAPIRouter(ctx context.Context) *subsonic.Router { fileCache := artwork.GetImageCache() imageStore := artwork.GetImageStore() fFmpeg := ffmpeg.New() - artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg) + broker := events.GetBroker() + artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg, broker) transcodingCache := stream.GetTranscodingCache() mediaStreamer := stream.NewMediaStreamer(dataStore, fFmpeg, transcodingCache) share := core.NewShare(dataStore) archiver := core.NewArchiver(mediaStreamer, dataStore, share) players := core.NewPlayers(dataStore) - broker := events.GetBroker() metricsMetrics := metrics.GetPrometheusInstance(dataStore) manager := plugins.GetManager(dataStore, broker, metricsMetrics) agentsAgents := agents.GetAgents(dataStore, manager) @@ -119,12 +119,12 @@ func CreateJellyfinAPIRouter(ctx context.Context) *jellyfin.Router { fileCache := artwork.GetImageCache() imageStore := artwork.GetImageStore() fFmpeg := ffmpeg.New() - artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg) + broker := events.GetBroker() + artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg, broker) transcodingCache := stream.GetTranscodingCache() mediaStreamer := stream.NewMediaStreamer(dataStore, fFmpeg, transcodingCache) transcodeDecider := stream.NewTranscodeDecider(dataStore, fFmpeg) players := core.NewPlayers(dataStore) - broker := events.GetBroker() metricsMetrics := metrics.GetPrometheusInstance(dataStore) manager := plugins.GetManager(dataStore, broker, metricsMetrics) playTracker := scrobbler.GetPlayTracker(dataStore, broker, manager) @@ -145,7 +145,8 @@ func CreatePublicRouter() *public.Router { fileCache := artwork.GetImageCache() imageStore := artwork.GetImageStore() fFmpeg := ffmpeg.New() - artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg) + broker := events.GetBroker() + artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg, broker) transcodingCache := stream.GetTranscodingCache() mediaStreamer := stream.NewMediaStreamer(dataStore, fFmpeg, transcodingCache) share := core.NewShare(dataStore) diff --git a/core/artwork/artwork.go b/core/artwork/artwork.go index 663d06d25..38fdef783 100644 --- a/core/artwork/artwork.go +++ b/core/artwork/artwork.go @@ -16,6 +16,7 @@ import ( "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/resources" + "github.com/navidrome/navidrome/server/events" "github.com/navidrome/navidrome/utils/cache" ) @@ -31,6 +32,7 @@ type Image struct { ETag string // representation validator; "" means Hash applies (full-size original) LastUpdated time.Time Placeholder bool + Transient bool // stand-in while the real representation is produced in the background } // representationTag varies with dimensions and encode settings, so a config change invalidates @@ -48,8 +50,8 @@ type Artwork interface { GetOrPlaceholder(ctx context.Context, id string, size int, square bool) (*Image, error) } -func NewArtwork(ds model.DataStore, cache cache.FileCache, store *ImageStore, ffm ffmpeg.FFmpeg) Artwork { - return &service{ds: ds, cache: cache, store: store, ffmpeg: ffm} +func NewArtwork(ds model.DataStore, cache cache.FileCache, store *ImageStore, ffm ffmpeg.FFmpeg, broker events.Broker) Artwork { + return &service{ds: ds, cache: cache, store: store, ffmpeg: ffm, broker: broker} } // entityExists reports whether the entity an artwork id points at is still there: state rows @@ -85,6 +87,7 @@ type service struct { cache cache.FileCache store *ImageStore ffmpeg ffmpeg.FFmpeg + broker events.Broker } func (s *service) GetOrPlaceholder(ctx context.Context, id string, size int, square bool) (*Image, error) { @@ -135,7 +138,7 @@ func (s *service) serveEntity(ctx context.Context, artID model.ArtworkID, size i // serveSource is the one place bytes become an Image. hash is the pixel identity ("" for disc art) // and doubles as the full-size validator, so an ETag is only needed when resized or hash is "". -func (s *service) serveSource(ctx context.Context, key, hash string, lastUpdate time.Time, +func (s *service) serveSource(ctx context.Context, artID model.ArtworkID, key, hash string, lastUpdate time.Time, size int, square bool, open func() (io.ReadCloser, error), ) (*Image, error) { if size == 0 && !square { @@ -152,15 +155,55 @@ func (s *service) serveSource(ctx context.Context, key, hash string, lastUpdate } return img, nil } - stream, err := s.cache.Get(ctx, &resizedItem{ - hash: key, size: size, square: square, ffmpeg: s.ffmpeg, open: open, - }) + // Without a cache a background conversion could never be served, so convert inline instead. + item := &resizedItem{hash: key, size: size, square: square, ffmpeg: s.ffmpeg, open: open, + deferAnimated: s.cache.Available(ctx)} + stream, err := s.cache.Get(ctx, item) + if deferred, ok := errors.AsType[*deferredAnimation](err); ok { + s.convertInBackground(artID, item, deferred.data) + return &Image{ReadCloser: io.NopCloser(deferred.standIn), Hash: hash, LastUpdated: lastUpdate, Transient: true}, nil + } if err != nil { return nil, err } return &Image{ReadCloser: stream, Hash: hash, ETag: representationTag(key, size, square), LastUpdated: lastUpdate}, nil } +// animConvertSlot bounds animated conversions process-wide: they outlive their request, so no +// request throttle limits them. +var animConvertSlot = make(chan struct{}, 1) + +const animConvertTimeout = time.Minute + +// convertInBackground caches item's conversion of data. When the slot is busy it does nothing, +// and a later request for the key tries again. +func (s *service) convertInBackground(artID model.ArtworkID, item *resizedItem, data []byte) { + select { + case animConvertSlot <- struct{}{}: + default: + return + } + key, size, square := item.Key(), item.size, item.square + go func() { + defer func() { <-animConvertSlot }() + ctx, cancel := context.WithTimeout(context.Background(), animConvertTimeout) + defer cancel() + start := time.Now() + resized, _, err := resizeImageData(ctx, s.ffmpeg, data, size, square) + out, _ := io.ReadAll(orOriginal(data, resized, err)) + // Converting before touching the cache keeps the entry from blocking readers for the whole conversion. + if err := warmCache(ctx, s.cache, &storedItem{key: key, data: out}); err != nil { + log.Warn(ctx, "Artwork: Could not cache animated image", "key", key, err) + return + } + log.Debug(ctx, "Artwork: Converted animated image", "key", key, "bytes", len(out), "elapsed", time.Since(start)) + // The image URL is unchanged, so UIs holding the stand-in only reload it when told. + if res, ok := artworkKindToResource[artID.Kind]; ok { + s.broker.SendBroadcastMessage(ctx, (&events.RefreshResource{}).With(res, artID.ID)) + } + }() +} + // serveHash serves the bytes of a found state row. A mismatch/open error is dangling, but a // cancelled request is not: it must not enqueue a re-resolution. func (s *service) serveHash(ctx context.Context, artID model.ArtworkID, ia *model.ItemArtwork, size int, square bool) (*Image, error) { @@ -175,7 +218,7 @@ func (s *service) serveHash(ctx context.Context, artID model.ArtworkID, ia *mode } return nil, err } - img, err := s.serveSource(ctx, ia.Hash, ia.Hash, ia.UpdatedAt, size, square, + img, err := s.serveSource(ctx, artID, ia.Hash, ia.Hash, ia.UpdatedAt, size, square, func() (io.ReadCloser, error) { return openOriginal(ia, art.Mime, s.store) }) if err != nil { if errors.Is(err, context.Canceled) { @@ -233,11 +276,11 @@ func (s *service) provisional(ctx context.Context, artID model.ArtworkID, size i s.enqueue(ctx, artID, model.ArtworkPriorityBump) log.Debug(ctx, "Artwork: Provisional read-through, no state row yet", "artID", artID, "source", res.source, "hit", res.reader != nil) - return s.serveResolution(ctx, res, size, square) + return s.serveResolution(ctx, artID, res, size, square) } // serveResolution turns a local resolution's bytes into a servable Image (byte-hash only, no decode). -func (s *service) serveResolution(ctx context.Context, res resolution, size int, square bool) (*Image, error) { +func (s *service) serveResolution(ctx context.Context, artID model.ArtworkID, res resolution, size int, square bool) (*Image, error) { if res.reader == nil { return nil, ErrUnavailable } @@ -251,7 +294,7 @@ func (s *service) serveResolution(ctx context.Context, res resolution, size int, return nil, ErrUnavailable } // Keyed by the byte-hash, so the entry lines up with the worker's eventual store entry. - return s.serveSource(ctx, hash, hash, unixMtime(res.refMtime), size, square, + return s.serveSource(ctx, artID, hash, hash, unixMtime(res.refMtime), size, square, func() (io.ReadCloser, error) { return io.NopCloser(bytes.NewReader(data)), nil }) } @@ -301,7 +344,7 @@ func (s *service) provisionalEmbedded(ctx context.Context, artID model.ArtworkID // Eligible but unextractable: fall back the way CoverArtID does, not to a placeholder. return s.Get(ctx, mf.DiscCoverArtID(), size, square) } - return s.serveResolution(ctx, res, size, square) + return s.serveResolution(ctx, artID, res, size, square) } // serveDisc reads disc art through with no state row and no enqueue, falling back to the album cover. @@ -319,7 +362,7 @@ func (s *service) serveDisc(ctx context.Context, artID model.ArtworkID, size int // Disc art has no state row, hence no content hash: keying on id, album mtime and // DiscArtPriority lets a warm cache answer without running the chain or touching the disk. key := fmt.Sprintf("%s|%d|%s", artID.ID, dr.cacheTime().UnixNano(), conf.Server.DiscArtPriority) - img, err := s.serveSource(ctx, key, "", dr.cacheTime(), size, square, selectImage) + img, err := s.serveSource(ctx, artID, key, "", dr.cacheTime(), size, square, selectImage) if err != nil { if errors.Is(err, context.Canceled) { return nil, err diff --git a/core/artwork/artwork_test.go b/core/artwork/artwork_test.go index 907b300de..af6e75214 100644 --- a/core/artwork/artwork_test.go +++ b/core/artwork/artwork_test.go @@ -3,10 +3,12 @@ package artwork import ( "bytes" "context" + "errors" "image" "io" "os" "path/filepath" + "sync" "time" "github.com/navidrome/navidrome/conf" @@ -14,6 +16,7 @@ import ( "github.com/navidrome/navidrome/consts" "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/resources" + "github.com/navidrome/navidrome/server/events" "github.com/navidrome/navidrome/tests" "github.com/navidrome/navidrome/utils/cache" . "github.com/onsi/ginkgo/v2" @@ -107,7 +110,7 @@ var _ = Describe("Artwork", func() { return arg.(artworkReader).Reader(ctx) }) Eventually(func() bool { return imgCache.Available(ctx) }, 10*time.Second).Should(BeTrue()) - svc = NewArtwork(ds, imgCache, store, ffm) + svc = NewArtwork(ds, imgCache, store, ffm, events.NoopBroker()) }) Describe("found state", func() { @@ -422,6 +425,95 @@ var _ = Describe("Artwork", func() { }) }) + Describe("animated GIF", func() { + var ( + gifBytes []byte + fake *animFFmpeg + broker *fakeEventBroker + ) + + get := func(id string) *Image { + GinkgoHelper() + img, err := svc.Get(ctx, model.MustParseArtworkID(id), 100, false) + Expect(err).ToNot(HaveOccurred()) + return img + } + waitForConversions := func() { + GinkgoHelper() + // Taking the slot, not len(), gives the race detector a happens-before with the conversion. + Eventually(func() bool { + select { + case animConvertSlot <- struct{}{}: + <-animConvertSlot + return true + default: + return false + } + }).Should(BeTrue()) + } + + BeforeEach(func() { + gifBytes = createAnimatedGIF(3) + fake = &animFFmpeg{MockFFmpeg: tests.NewMockFFmpeg(""), out: []byte("animated-webp")} + broker = &fakeEventBroker{} + svc = NewArtwork(ds, imgCache, store, fake, broker) + seedFoundStore("al", "al1", gifBytes) + DeferCleanup(waitForConversions) + }) + + It("serves a transient stand-in, then the conversion from the cache", func() { + img := get("al-al1") + Expect(img.Transient).To(BeTrue()) + Expect(readAll(img)).To(Equal(gifBytes)) + + waitForConversions() + img = get("al-al1") + Expect(img.Transient).To(BeFalse()) + Expect(readAll(img)).To(Equal([]byte("animated-webp"))) + Expect(fake.calls.Load()).To(Equal(int32(1))) + }) + + It("tells the UI to reload the item once its conversion is cached", func() { + readAll(get("al-al1")) + waitForConversions() + + sent := broker.getEvents() + Expect(sent).To(HaveLen(1)) + Expect(sent[0].Data(sent[0])).To(MatchJSON(`{"album":["al1"]}`)) + }) + + It("runs one conversion at a time and does not queue the others", func() { + fake.release = make(chan struct{}) + release := sync.OnceFunc(func() { close(fake.release) }) + DeferCleanup(release) + seedFoundStore("al", "al2", createAnimatedGIF(4)) + + readAll(get("al-al1")) + Eventually(fake.calls.Load).Should(Equal(int32(1))) + img := get("al-al2") + Expect(img.Transient).To(BeTrue()) + readAll(img) + Consistently(fake.calls.Load, "100ms").Should(Equal(int32(1))) + + release() + waitForConversions() + img = get("al-al2") + Expect(img.Transient).To(BeTrue()) + readAll(img) + }) + + It("caches the static fallback when the conversion fails, so it is not retried", func() { + fake.Error = errors.New("pipe:0: Input/output error") + readAll(get("al-al1")) + waitForConversions() + + img := get("al-al1") + Expect(img.Transient).To(BeFalse()) + Expect(readAll(img)).To(Equal(gifBytes)) + Expect(fake.calls.Load()).To(Equal(int32(1))) + }) + }) + Describe("GetOrPlaceholder", func() { It("accepts a raw entity id and serves its cover art", func() { albumRepo.SetData(model.Albums{{ID: "rawal", Name: "Album"}}) diff --git a/core/artwork/e2e/acquire_serve_test.go b/core/artwork/e2e/acquire_serve_test.go index 34dfb2cac..ea23fa08c 100644 --- a/core/artwork/e2e/acquire_serve_test.go +++ b/core/artwork/e2e/acquire_serve_test.go @@ -108,7 +108,7 @@ var _ = Describe("Acquisition → serve loop", func() { }) Eventually(func() bool { return imgCache.Available(ctx) }, 10*time.Second).Should(BeTrue()) - svc = artwork.NewArtwork(ds, imgCache, store, ffm) + svc = artwork.NewArtwork(ds, imgCache, store, ffm, events.NoopBroker()) worker = artwork.NewWorker(ds, store, agents.GetAgents(ds, nil), ffm, events.NoopBroker(), imgCache) }) diff --git a/core/artwork/e2e/resolution_harness_test.go b/core/artwork/e2e/resolution_harness_test.go index fff62ce69..6427d7fce 100644 --- a/core/artwork/e2e/resolution_harness_test.go +++ b/core/artwork/e2e/resolution_harness_test.go @@ -119,7 +119,7 @@ func setupResolutionHarness() { }) Eventually(func() bool { return imgCache.Available(rctx) }, 10*time.Second).Should(BeTrue()) - rsvc = artwork.NewArtwork(rds, imgCache, rstore, ffm) + rsvc = artwork.NewArtwork(rds, imgCache, rstore, ffm, events.NoopBroker()) rworker = artwork.NewWorker(rds, rstore, agents.GetAgents(rds, nil), ffm, events.NoopBroker(), imgCache) } diff --git a/core/artwork/image_cache.go b/core/artwork/image_cache.go index b1970d21d..7d58b56f1 100644 --- a/core/artwork/image_cache.go +++ b/core/artwork/image_cache.go @@ -42,6 +42,8 @@ type resizedItem struct { square bool ffmpeg ffmpeg.FFmpeg open func() (io.ReadCloser, error) + // deferAnimated makes Reader return a *deferredAnimation for animated GIFs instead of converting inline. + deferAnimated bool } // Key is the ETag namespaced for the cache, so the validator a client holds and the entry it @@ -64,13 +66,50 @@ func (r *resizedItem) Reader(ctx context.Context) (io.ReadCloser, error) { if err != nil { return nil, err } + if r.deferAnimated && isAnimatedGIF(data) && r.ffmpeg.IsAvailable() { + static, _, err := resizeStaticImage(data, r.size, r.square) + return nil, &deferredAnimation{data: data, standIn: orOriginal(data, static, err)} + } resized, _, err := resizeImageData(ctx, r.ffmpeg, data, r.size, r.square) - if err != nil || resized == nil { - // Resize failed or image already within bounds: serve the original bytes. - return io.NopCloser(bytes.NewReader(data)), nil - } - if rc, ok := resized.(io.ReadCloser); ok { - return rc, nil - } - return io.NopCloser(resized), nil + return io.NopCloser(orOriginal(data, resized, err)), nil +} + +// orOriginal serves the original bytes when the resize failed or the image was already within bounds. +func orOriginal(data []byte, resized io.Reader, err error) io.Reader { + if err != nil || resized == nil { + return bytes.NewReader(data) + } + return resized +} + +// deferredAnimation is returned as an error so the cache stores nothing: it carries the GIF bytes +// to convert in the background and a static stand-in to serve meanwhile. +type deferredAnimation struct { + data []byte + standIn io.Reader +} + +func (*deferredAnimation) Error() string { return "animated image conversion deferred" } + +// warmCache stores item in c, returning once the entry is written. +func warmCache(ctx context.Context, c cache.FileCache, item cache.Item) error { + stream, err := c.Get(ctx, item) + if err != nil { + return err + } + defer stream.Close() + _, err = io.Copy(io.Discard, stream) + return err +} + +// storedItem caches bytes produced outside the cache. +type storedItem struct { + key string + data []byte +} + +func (i *storedItem) Key() string { return i.key } + +func (i *storedItem) Reader(context.Context) (io.ReadCloser, error) { + return io.NopCloser(bytes.NewReader(i.data)), nil } diff --git a/core/artwork/resize.go b/core/artwork/resize.go index 9b3ff46e5..f233f0586 100644 --- a/core/artwork/resize.go +++ b/core/artwork/resize.go @@ -3,6 +3,7 @@ package artwork import ( "bytes" "context" + "errors" "fmt" "image" "image/draw" @@ -48,9 +49,9 @@ func resizeImageData(ctx context.Context, ffm ffmpeg.FFmpeg, data []byte, size i if isAnimatedGIF(data) { if ffm.IsAvailable() { // Animated GIF: convert to animated WebP via ffmpeg (with optional resize) - r, err := ffm.ConvertAnimatedImage(ctx, bytes.NewReader(data), size, conf.Server.CoverArtQuality) + out, err := convertAnimatedGIF(ctx, ffm, data, size) if err == nil { - return r, 0, nil + return bytes.NewReader(out), 0, nil } log.Warn(ctx, "Artwork: Could not convert animated GIF, falling back to static", err) } @@ -62,6 +63,24 @@ func resizeImageData(ctx context.Context, ffm ffmpeg.FFmpeg, data []byte, size i return resizeStaticImage(data, size, square) } +// convertAnimatedGIF buffers the output because ffmpeg can fail after it starts streaming, +// too late for a caller holding the stream to fall back. +func convertAnimatedGIF(ctx context.Context, ffm ffmpeg.FFmpeg, data []byte, size int) ([]byte, error) { + r, err := ffm.ConvertAnimatedImage(ctx, bytes.NewReader(data), size, conf.Server.CoverArtQuality) + if err != nil { + return nil, err + } + defer r.Close() + out, err := io.ReadAll(r) + if err != nil { + return nil, err + } + if len(out) == 0 { + return nil, errors.New("ffmpeg produced no output") + } + return out, nil +} + // toFastScaleType converts types x/image/draw has no optimized scaler for (e.g. *image.NYCbCrA, // *image.Paletted) to *image.RGBA, avoiding CatmullRom.Scale's generic per-pixel fallback. func toFastScaleType(img image.Image) image.Image { diff --git a/core/artwork/resize_test.go b/core/artwork/resize_test.go index 76b95974d..73463f901 100644 --- a/core/artwork/resize_test.go +++ b/core/artwork/resize_test.go @@ -1,6 +1,17 @@ package artwork import ( + "bytes" + "context" + "errors" + "image" + "io" + "sync/atomic" + "testing/iotest" + + "github.com/navidrome/navidrome/conf" + "github.com/navidrome/navidrome/conf/configtest" + "github.com/navidrome/navidrome/tests" . "github.com/onsi/ginkgo/v2" . "github.com/onsi/gomega" ) @@ -11,3 +22,62 @@ var _ = Describe("resizeStaticImage", func() { Expect(err).To(MatchError(ContainSubstring("exceed pixel cap"))) }) }) + +var _ = Describe("resizeImageData", func() { + var gifBytes []byte + + BeforeEach(func() { + DeferCleanup(configtest.SetupConfig()) + conf.Server.EnableWebPEncoding = false + gifBytes = createAnimatedGIF(3) + }) + + It("converts an animated GIF with ffmpeg", func() { + fake := &animFFmpeg{MockFFmpeg: tests.NewMockFFmpeg(""), out: []byte("animated-webp")} + r, _, err := resizeImageData(context.Background(), fake, gifBytes, 1, false) + Expect(err).ToNot(HaveOccurred()) + Expect(io.ReadAll(r)).To(Equal([]byte("animated-webp"))) + }) + + DescribeTable("falls back to a static resize when ffmpeg fails", + func(fake *animFFmpeg) { + r, _, err := resizeImageData(context.Background(), fake, gifBytes, 1, false) + Expect(err).ToNot(HaveOccurred()) + data, err := io.ReadAll(r) + Expect(err).ToNot(HaveOccurred()) + _, format, err := image.DecodeConfig(bytes.NewReader(data)) + Expect(err).ToNot(HaveOccurred()) + Expect(format).To(Equal("jpeg")) + }, + Entry("before producing output", &animFFmpeg{MockFFmpeg: &tests.MockFFmpeg{Error: errors.New("no libwebp_anim")}}), + Entry("mid-stream, like ffmpeg 5.1 reading a GIF from a pipe", &animFFmpeg{MockFFmpeg: tests.NewMockFFmpeg(""), streamErr: errors.New("pipe:0: Input/output error")}), + Entry("with empty output", &animFFmpeg{MockFFmpeg: tests.NewMockFFmpeg("")}), + ) +}) + +// animFFmpeg fakes animated conversion; a non-nil release blocks it until closed. +type animFFmpeg struct { + *tests.MockFFmpeg + calls atomic.Int32 + release chan struct{} + out []byte + streamErr error +} + +func (f *animFFmpeg) ConvertAnimatedImage(ctx context.Context, _ io.Reader, _ int, _ int) (io.ReadCloser, error) { + f.calls.Add(1) + if f.release != nil { + select { + case <-f.release: + case <-ctx.Done(): + return nil, ctx.Err() + } + } + if f.Error != nil { + return nil, f.Error + } + if f.streamErr != nil { + return io.NopCloser(iotest.ErrReader(f.streamErr)), nil //nolint:nilerr // the stream fails, not the call + } + return io.NopCloser(bytes.NewReader(f.out)), nil +} diff --git a/core/artwork/worker.go b/core/artwork/worker.go index bb09be55e..3c76e3c72 100644 --- a/core/artwork/worker.go +++ b/core/artwork/worker.go @@ -310,13 +310,10 @@ func (w *Worker) precache(ctx context.Context, got *acquired) { ffmpeg: w.ffmpeg, open: func() (io.ReadCloser, error) { return io.NopCloser(bytes.NewReader(got.data)), nil }, } - stream, err := w.cache.Get(ctx, item) - if err != nil { + if err := warmCache(ctx, w.cache, item); err != nil { log.Debug(ctx, "Artwork: Precache failed", "kind", got.ia.ItemKind, "id", got.ia.ItemID, err) return } - _, _ = io.Copy(io.Discard, stream) - _ = stream.Close() log.Trace(ctx, "Artwork: Precached UI size", "kind", got.ia.ItemKind, "id", got.ia.ItemID, "size", conf.Server.UICoverArtSize, "elapsed", time.Since(precacheStart)) } diff --git a/server/imghttp/headers.go b/server/imghttp/headers.go index 9308bcdab..561da34b2 100644 --- a/server/imghttp/headers.go +++ b/server/imghttp/headers.go @@ -13,8 +13,8 @@ import ( // (the caller must then not write a body). requestedHash is the hash the client asserted, or "". func WriteImageHeaders(w http.ResponseWriter, r *http.Request, img *artwork.Image, requestedHash string) (wrote304 bool) { h := w.Header() - // Placeholders are transient stand-ins for unresolved art: never cached, no validators. - if img.Placeholder { + // Placeholders and transient stand-ins must not outlive what they stand in for: never cached, no validators. + if img.Placeholder || img.Transient { h.Set("Cache-Control", "no-store") return false } diff --git a/server/imghttp/headers_test.go b/server/imghttp/headers_test.go index 2e8313faf..ec06f84d6 100644 --- a/server/imghttp/headers_test.go +++ b/server/imghttp/headers_test.go @@ -41,6 +41,12 @@ func resized() *artwork.Image { } } +func transient() *artwork.Image { + img := resized() + img.Transient = true + return img +} + func unvalidated() *artwork.Image { return &artwork.Image{ReadCloser: io.NopCloser(strings.NewReader("IMG")), LastUpdated: lastMod} } @@ -105,6 +111,9 @@ var _ = Describe("WriteImageHeaders", func() { testCase{img: found(), ifNoneMatch: `"deadbeefdeadbeef"`, want304: false, wantCache: "public, no-cache", wantETag: `"` + testHash + `"`, wantLastMod: true}), Entry("placeholder ignores If-None-Match and never 304s", testCase{img: placeholder(), ifNoneMatch: "*", want304: false, wantCache: "no-store"}), + // A stand-in shares the final representation's ETag inputs, so a validator would pin it. + Entry("transient stand-in is never cached and carries no validators", + testCase{img: transient(), ifNoneMatch: `"` + testRepTag + `"`, want304: false, wantCache: "no-store"}), // With no validator, an emitted ETag would be the same empty tag on every such response, // and matching it would 304 changed bytes. diff --git a/server/subsonic/e2e/subsonic_artwork_test.go b/server/subsonic/e2e/subsonic_artwork_test.go index 9324ea9e3..57d916f64 100644 --- a/server/subsonic/e2e/subsonic_artwork_test.go +++ b/server/subsonic/e2e/subsonic_artwork_test.go @@ -132,7 +132,7 @@ var _ = Describe("Artwork Serving", Ordered, func() { store := artwork.NewImageStore(GinkgoT().TempDir()) imgCache := newDummyImageCache(ctx) ffm := harness.NoopFFmpeg{} - artSvc = artwork.NewArtwork(ds, imgCache, store, ffm) + artSvc = artwork.NewArtwork(ds, imgCache, store, ffm, events.NoopBroker()) worker = artwork.NewWorker(ds, store, agents.GetAgents(ds, nil), ffm, events.NoopBroker(), imgCache) artRouter = buildArtworkRouter(artSvc) diff --git a/ui/src/common/Artwork.jsx b/ui/src/common/Artwork.jsx index 21ac5e07b..3f39bc94a 100644 --- a/ui/src/common/Artwork.jsx +++ b/ui/src/common/Artwork.jsx @@ -5,6 +5,7 @@ import { makeStyles } from '@material-ui/core/styles' import config from '../config' import subsonic from '../subsonic' import { useImageUrl } from './useImageUrl' +import { useArtworkRefresh } from './useArtworkRefresh' import { ThumbHashCanvas } from './ThumbHashCanvas' // Drives both the CSS transition and the timer that retires the placeholder, so they cannot drift. @@ -47,7 +48,8 @@ export const Artwork = ({ }) => { const classes = useStyles() const url = record ? subsonic.getCoverArtUrl(record, size, square) : '' - const { imgUrl, fromCache } = useImageUrl(url) + const version = useArtworkRefresh(record?.id) + const { imgUrl, fromCache } = useImageUrl(url, version) const [decoded, setDecoded] = useState(false) const [faded, setFaded] = useState(false) diff --git a/ui/src/common/Artwork.test.jsx b/ui/src/common/Artwork.test.jsx index 18b3c916f..bbf744057 100644 --- a/ui/src/common/Artwork.test.jsx +++ b/ui/src/common/Artwork.test.jsx @@ -2,6 +2,7 @@ import { render, fireEvent, act } from '@testing-library/react' import { describe, it, expect, vi, beforeEach } from 'vitest' vi.mock('./useImageUrl', () => ({ useImageUrl: vi.fn() })) +vi.mock('./useArtworkRefresh', () => ({ useArtworkRefresh: vi.fn() })) vi.mock('../subsonic', () => ({ default: { // Mirrors the real hash suffix, so a refreshed record yields a different URL. @@ -13,6 +14,7 @@ vi.mock('../subsonic', () => ({ vi.mock('../config', () => ({ default: { uiCoverArtSize: 300 } })) import { useImageUrl } from './useImageUrl' +import { useArtworkRefresh } from './useArtworkRefresh' import { Artwork } from './Artwork' const withArt = { @@ -28,6 +30,14 @@ describe('Artwork', () => { HTMLCanvasElement.prototype.getContext = vi.fn(() => null) }) + it('loads the image under the version of the last refresh naming its record', () => { + useImageUrl.mockReturnValue({ imgUrl: null, loading: true }) + useArtworkRefresh.mockReturnValue(1234) + render() + expect(useArtworkRefresh).toHaveBeenCalledWith('al-1') + expect(useImageUrl).toHaveBeenCalledWith('/rest/getCoverArt?id=al-1', 1234) + }) + it('renders nothing without a record', () => { useImageUrl.mockReturnValue({ imgUrl: null, loading: false }) const { container } = render() diff --git a/ui/src/common/useArtworkRefresh.js b/ui/src/common/useArtworkRefresh.js new file mode 100644 index 000000000..7fb92586b --- /dev/null +++ b/ui/src/common/useArtworkRefresh.js @@ -0,0 +1,23 @@ +import { useState } from 'react' +import { useSelector } from 'react-redux' + +const namesId = (resources, id) => + Object.values(resources || {}).some( + (ids) => Array.isArray(ids) && ids.includes(id), + ) + +// Artwork can change behind an unchanged URL (a background conversion replacing its stand-in), so +// this returns when the server last announced new artwork for id, for useImageUrl to version by. +export const useArtworkRefresh = (id) => { + // A primitive, so an event naming other items does not re-render this cover. + const named = useSelector(({ activity }) => + id && namesId(activity?.refresh?.resources, id) + ? activity.refresh.lastReceived + : undefined, + ) + const [refreshedAt, setRefreshedAt] = useState() + if (named !== undefined && named !== refreshedAt) { + setRefreshedAt(named) + } + return refreshedAt +} diff --git a/ui/src/common/useArtworkRefresh.test.js b/ui/src/common/useArtworkRefresh.test.js new file mode 100644 index 000000000..89e6f9a1d --- /dev/null +++ b/ui/src/common/useArtworkRefresh.test.js @@ -0,0 +1,44 @@ +import { renderHook } from '@testing-library/react-hooks' +import { vi, describe, it, expect, beforeEach } from 'vitest' +import { useSelector } from 'react-redux' +import { useArtworkRefresh } from './useArtworkRefresh' + +vi.mock('react-redux', () => ({ useSelector: vi.fn() })) + +const withRefresh = (refresh) => + useSelector.mockImplementation((select) => select({ activity: { refresh } })) + +describe('useArtworkRefresh', () => { + beforeEach(() => { + vi.clearAllMocks() + }) + + it('has no version before any refresh event', () => { + withRefresh(undefined) + const { result } = renderHook(() => useArtworkRefresh('ar-1')) + expect(result.current).toBeUndefined() + }) + + it('versions the artwork with the time of an event naming its id', () => { + withRefresh({ lastReceived: 100, resources: { artist: ['ar-1'] } }) + const { result } = renderHook(() => useArtworkRefresh('ar-1')) + expect(result.current).toBe(100) + }) + + it('keeps its version when a later event names something else', () => { + withRefresh({ lastReceived: 100, resources: { artist: ['ar-1'] } }) + const { result, rerender } = renderHook(() => useArtworkRefresh('ar-1')) + + withRefresh({ lastReceived: 200, resources: { album: ['al-9'] } }) + rerender() + + expect(result.current).toBe(100) + }) + + // Honoring a wildcard would refetch every cover on screen after each scan. + it('ignores wildcard events', () => { + withRefresh({ lastReceived: 100, resources: { '*': '*', artist: ['*'] } }) + const { result } = renderHook(() => useArtworkRefresh('ar-1')) + expect(result.current).toBeUndefined() + }) +}) diff --git a/ui/src/common/useImageUrl.js b/ui/src/common/useImageUrl.js index 8a7f2bcc7..71811dbed 100644 --- a/ui/src/common/useImageUrl.js +++ b/ui/src/common/useImageUrl.js @@ -34,9 +34,10 @@ const evictIfNeeded = () => { /** * Loads an image via fetch() with AbortController so that in-flight requests * are canceled on unmount (e.g., during pagination). Uses a module-level cache - * so remounting returns the cached blob URL instantly. + * so remounting returns the cached blob URL instantly. A no-store response is a stand-in: it is + * re-fetched on remount and on a new version, and painted until its replacement arrives. */ -export const useImageUrl = (url) => { +export const useImageUrl = (url, version) => { const cached = url ? cache.get(url) : null const [imgUrl, setImgUrl] = useState(cached?.blobUrl || null) const [loading, setLoading] = useState(!!url && !cached) @@ -68,7 +69,7 @@ export const useImageUrl = (url) => { // Re-check: another component's effect may have populated the cache // between this component's render and effect execution. const entry = cache.get(url) - if (entry) { + if (entry && !entry.transient) { entry.refCount++ setImgUrl(entry.blobUrl) setLoading(false) @@ -77,12 +78,15 @@ export const useImageUrl = (url) => { entry.refCount-- } } + const standIn = !!entry const controller = new AbortController() let queued = true - setImgUrl(null) - setLoading(true) - setError(false) + if (!standIn) { + setImgUrl(null) + setLoading(true) + setError(false) + } const doFetch = () => { queued = false @@ -92,9 +96,12 @@ export const useImageUrl = (url) => { if (!res.ok) { throw new Error(`HTTP ${res.status}`) } - return res.blob() + const transient = !!res.headers + ?.get('Cache-Control') + ?.includes('no-store') + return res.blob().then((blob) => ({ blob, transient })) }) - .then((blob) => { + .then(({ blob, transient }) => { activeFetches-- processQueue() // Guard against late resolution after abort @@ -105,12 +112,15 @@ export const useImageUrl = (url) => { // Handle concurrent fetches: if another component already cached // this URL, use its entry and discard our blob. const existing = cache.get(url) - if (existing && existing.blobUrl) { + if (existing?.blobUrl && !existing.transient) { existing.refCount++ URL.revokeObjectURL(objectUrl) setImgUrl(existing.blobUrl) } else { - cache.set(url, { blobUrl: objectUrl, refCount: 1 }) + if (existing?.blobUrl && existing.refCount <= 0) { + URL.revokeObjectURL(existing.blobUrl) + } + cache.set(url, { blobUrl: objectUrl, refCount: 1, transient }) evictIfNeeded() setImgUrl(objectUrl) } @@ -122,6 +132,9 @@ export const useImageUrl = (url) => { if (err.name === 'AbortError') { return // Expected on unmount or URL change } + if (standIn) { + return // Keep painting the stand-in + } // Cache the error so repeated mounts don't re-fetch broken URLs cache.set(url, { blobUrl: null, error: true, refCount: 0 }) setError(true) @@ -150,7 +163,7 @@ export const useImageUrl = (url) => { entry.refCount-- } } - }, [url]) + }, [url, version]) return { imgUrl, loading, error, fromCache } } diff --git a/ui/src/common/useImageUrl.test.js b/ui/src/common/useImageUrl.test.js index 818732d06..fc7b91ab8 100644 --- a/ui/src/common/useImageUrl.test.js +++ b/ui/src/common/useImageUrl.test.js @@ -234,4 +234,98 @@ describe('useImageUrl', () => { // A remembered failure has no blob, so callers must not treat it as instantly painted. expect(result2.current.fromCache).toBe(false) }) + + describe('no-store stand-ins', () => { + const image = (data, cacheControl = 'public, no-cache') => + Promise.resolve({ + ok: true, + headers: { + get: (name) => (name === 'Cache-Control' ? cacheControl : null), + }, + blob: () => Promise.resolve(new Blob([data])), + }) + const standIn = (data) => image(data, 'no-store') + const renderVersioned = (version) => + renderHook( + ({ version }) => useImageUrl('http://example.com/img.jpg', version), + { initialProps: { version } }, + ) + const settle = () => + act(async () => { + await flushPromises() + }) + + it('re-fetches a stand-in when its version changes, painting it meanwhile', async () => { + let resolveFinal + global.URL.createObjectURL = vi + .fn() + .mockReturnValueOnce('blob:stand-in') + .mockReturnValueOnce('blob:final') + global.fetch = vi + .fn() + .mockImplementationOnce(() => standIn('still')) + .mockImplementationOnce( + () => new Promise((resolve) => (resolveFinal = resolve)), + ) + + const { result, rerender } = renderVersioned(undefined) + await settle() + expect(result.current.imgUrl).toBe('blob:stand-in') + + rerender({ version: 1 }) + await settle() + expect(global.fetch).toHaveBeenCalledTimes(2) + expect(result.current.imgUrl).toBe('blob:stand-in') + expect(result.current.loading).toBe(false) + + await act(async () => { + resolveFinal(image('animated')) + await flushPromises() + }) + expect(result.current.imgUrl).toBe('blob:final') + }) + + it('keeps the stand-in when the re-fetch fails', async () => { + global.fetch = vi + .fn() + .mockImplementationOnce(() => standIn('still')) + .mockImplementationOnce(() => + Promise.resolve({ ok: false, status: 500 }), + ) + + const { result, rerender } = renderVersioned(undefined) + await settle() + rerender({ version: 1 }) + await settle() + + expect(global.fetch).toHaveBeenCalledTimes(2) + expect(result.current.imgUrl).toBe('blob:mock-url') + expect(result.current.error).toBe(false) + }) + + // Plays, stars and ratings also name items, and must not re-download final images. + it('does not re-fetch a final image when its version changes', async () => { + global.fetch = vi.fn(() => image('final')) + + const { rerender } = renderVersioned(undefined) + await settle() + rerender({ version: 1 }) + await settle() + + expect(global.fetch).toHaveBeenCalledTimes(1) + }) + + it('re-fetches a stand-in on remount', async () => { + global.fetch = vi.fn(() => standIn('still')) + + const { unmount } = renderVersioned(undefined) + await settle() + unmount() + const { result } = renderVersioned(undefined) + await settle() + + expect(global.fetch).toHaveBeenCalledTimes(2) + expect(result.current.imgUrl).toBe('blob:mock-url') + }) + }) }) diff --git a/utils/cache/file_caches.go b/utils/cache/file_caches.go index 48cd135cb..91cf87f10 100644 --- a/utils/cache/file_caches.go +++ b/utils/cache/file_caches.go @@ -178,7 +178,12 @@ func (fc *fileCache) Get(ctx context.Context, arg Item) (*CachedStream, error) { reader, err := fc.getReader(ctx, arg) if err != nil { _ = r.Close() - _ = w.Close() + // Cancelling fails readers that joined meanwhile with err, and keeps the removal from blocking on them. + if cw, ok := w.(interface{ CloseWithError(error) error }); ok { + _ = cw.CloseWithError(err) + } else { + _ = w.Close() + } _ = fc.invalidate(ctx, key) return nil, err } diff --git a/utils/cache/file_caches_test.go b/utils/cache/file_caches_test.go index 3189de6b2..5547db736 100644 --- a/utils/cache/file_caches_test.go +++ b/utils/cache/file_caches_test.go @@ -106,6 +106,30 @@ var _ = Describe("File Caches", func() { Expect(called).To(BeTrue()) }) + It("fails a reader that joined a miss whose loader failed, without blocking the removal", func() { + entered, release := make(chan struct{}), make(chan struct{}) + fc := callNewFileCache("test", "1KB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) { + close(entered) + <-release + return nil, errors.New("loader failed") + }) + firstErr := make(chan error, 1) + go func() { + _, err := fc.Get(context.Background(), &testArg{"k"}) + firstErr <- err + }() + <-entered + + joined, err := fc.Get(context.Background(), &testArg{"k"}) + Expect(err).ToNot(HaveOccurred()) + close(release) + Eventually(firstErr).Should(Receive(MatchError("loader failed"))) + + _, err = io.ReadAll(joined) + Expect(err).To(MatchError("loader failed")) + _ = joined.Close() + }) + It("writes a completion marker after a successful cache write", func() { fc := callNewFileCache("test", "1KB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) { return strings.NewReader("complete-data"), nil