From bc16cafb91a7795f545866fd4854937e5f05bb8d Mon Sep 17 00:00:00 2001 From: Deluan Date: Mon, 14 Sep 2026 19:36:50 -0400 Subject: [PATCH] fix(artwork): convert animated GIF covers in the background Resized animated GIF covers were converted to animated WebP by ffmpeg inside the image request, bound to its 10s timeout. On slow CPUs large GIFs never finished: ffmpeg was killed, nothing was cached, and every request repeated the work. With ffmpeg 5.1 the conversion also failed mid-stream when reading the GIF from a pipe, after the stream was already handed to the client, so there was no fallback. Requests now get a static stand-in right away, served with no-store and no validators, while the conversion runs in the background. A single process-wide slot bounds these conversions, since they outlive the request throttle; when it is busy nothing is queued and a later request retries. The result, or a static fallback if ffmpeg fails, is cached under the same key. ffmpeg output is now buffered so mid-stream failures fall back to a static resize. The file cache gains cache.Transient for loader results that must be served but not stored. Without an available cache, conversion stays inline. --- core/artwork/artwork.go | 64 ++++++++++++++++++++++++-- core/artwork/artwork_test.go | 80 +++++++++++++++++++++++++++++++++ core/artwork/image_cache.go | 22 +++++++++ core/artwork/resize.go | 23 +++++++++- core/artwork/resize_test.go | 72 +++++++++++++++++++++++++++++ server/imghttp/headers.go | 4 +- server/imghttp/headers_test.go | 9 ++++ utils/cache/file_caches.go | 36 ++++++++++++++- utils/cache/file_caches_test.go | 61 +++++++++++++++++++++++++ 9 files changed, 362 insertions(+), 9 deletions(-) diff --git a/core/artwork/artwork.go b/core/artwork/artwork.go index 663d06d25..1da419dff 100644 --- a/core/artwork/artwork.go +++ b/core/artwork/artwork.go @@ -31,6 +31,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 @@ -152,15 +153,72 @@ 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, - }) + item := &resizedItem{hash: key, size: size, square: square, ffmpeg: s.ffmpeg, open: open} + var animated []byte + // Without a cache a background conversion could never be served, so convert inline instead. + if s.cache.Available(ctx) { + item.deferAnimated = func(data []byte) { animated = data } + } + stream, err := s.cache.Get(ctx, item) if err != nil { return nil, err } + if stream.Transient { + s.convertInBackground(item, animated) + return &Image{ReadCloser: stream, Hash: hash, LastUpdated: lastUpdate, Transient: true}, nil + } 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(item *resizedItem, data []byte) { + select { + case animConvertSlot <- struct{}{}: + default: + return + } + conv := *item + conv.deferAnimated = nil + conv.open = func() (io.ReadCloser, error) { return io.NopCloser(bytes.NewReader(data)), nil } + key := conv.Key() + go func() { + defer func() { <-animConvertSlot }() + ctx, cancel := context.WithTimeout(context.Background(), animConvertTimeout) + defer cancel() + start := time.Now() + rc, err := conv.Reader(ctx) + if err != nil { + log.Warn(ctx, "Artwork: Could not convert animated image", "key", key, err) + return + } + out, err := io.ReadAll(rc) + _ = rc.Close() + if err != nil { + log.Warn(ctx, "Artwork: Could not convert animated image", "key", key, err) + return + } + // Converting before touching the cache keeps the entry from blocking readers for the whole conversion. + stream, err := s.cache.Get(ctx, &storedItem{key: key, data: out}) + if err != nil { + log.Warn(ctx, "Artwork: Could not cache animated image", "key", key, err) + return + } + defer stream.Close() + if _, err := io.Copy(io.Discard, stream); err != nil { + log.Debug(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)) + }() +} + // 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) { diff --git a/core/artwork/artwork_test.go b/core/artwork/artwork_test.go index 907b300de..707eb3e89 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" @@ -422,6 +424,84 @@ var _ = Describe("Artwork", func() { }) }) + Describe("animated GIF", func() { + var ( + gifBytes []byte + fake *animFFmpeg + ) + + 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")} + svc = NewArtwork(ds, imgCache, store, fake) + 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("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.err = 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/image_cache.go b/core/artwork/image_cache.go index b1970d21d..53eeb26f5 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, when set, receives an animated GIF's bytes, and a transient static stand-in is served. + deferAnimated func(data []byte) } // Key is the ETag namespaced for the cache, so the validator a client holds and the entry it @@ -64,6 +66,14 @@ func (r *resizedItem) Reader(ctx context.Context) (io.ReadCloser, error) { if err != nil { return nil, err } + if r.deferAnimated != nil && isAnimatedGIF(data) && r.ffmpeg.IsAvailable() { + r.deferAnimated(data) + static, _, err := resizeStaticImage(data, r.size, r.square) + if err != nil || static == nil { + static = bytes.NewReader(data) + } + return cache.Transient(io.NopCloser(static)), nil + } 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. @@ -74,3 +84,15 @@ func (r *resizedItem) Reader(ctx context.Context) (io.ReadCloser, error) { } return io.NopCloser(resized), nil } + +// 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..d89ee0f3b 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,64 @@ 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) { + fake.MockFFmpeg = tests.NewMockFFmpeg("") + 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{err: errors.New("no libwebp_anim")}), + Entry("mid-stream, like ffmpeg 5.1 reading a GIF from a pipe", &animFFmpeg{streamErr: errors.New("pipe:0: Input/output error")}), + Entry("with empty output", &animFFmpeg{}), + ) +}) + +// 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 + err error + 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.err != nil { + return nil, f.err + } + 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/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/utils/cache/file_caches.go b/utils/cache/file_caches.go index 48cd135cb..90764f326 100644 --- a/utils/cache/file_caches.go +++ b/utils/cache/file_caches.go @@ -25,6 +25,18 @@ type Item interface { Key() string } +// Transient wraps a ReadFunc result that Get must serve but never store. +func Transient(rc io.ReadCloser) io.ReadCloser { return transientReader{rc} } + +type transientReader struct{ io.ReadCloser } + +func unwrapTransient(r io.Reader) (io.Reader, bool) { + if t, ok := r.(transientReader); ok { + return t.ReadCloser, true + } + return r, false +} + // ReadFunc is a function that retrieves the data to be cached. It receives the Item to be cached and returns // an io.Reader with the data and an error. type ReadFunc func(ctx context.Context, item Item) (io.Reader, error) @@ -155,7 +167,8 @@ func (fc *fileCache) Get(ctx context.Context, arg Item) (*CachedStream, error) { if err != nil { return nil, err } - return &CachedStream{Reader: reader}, nil + reader, transient := unwrapTransient(reader) + return &CachedStream{Reader: reader, Transient: transient}, nil } key := arg.Key() @@ -182,6 +195,10 @@ func (fc *fileCache) Get(ctx context.Context, arg Item) (*CachedStream, error) { _ = fc.invalidate(ctx, key) return nil, err } + if reader, transient := unwrapTransient(reader); transient { + fc.discard(ctx, key, r, w) + return &CachedStream{Reader: reader, Transient: true}, nil + } go func() { if err := fc.copyAndClose(ctx, key, w, reader); err != nil { log.Debug(ctx, "Error storing file in cache", "cache", fc.name, "key", key, err) @@ -213,12 +230,27 @@ func (fc *fileCache) Get(ctx context.Context, arg Item) (*CachedStream, error) { return &CachedStream{Reader: r, Cached: cached}, nil } +// discard drops a miss entry that will never be written. Cancelling first makes the removal +// non-blocking, so the key is free when Get returns; readers that joined meanwhile fail. +func (fc *fileCache) discard(ctx context.Context, key string, r io.Closer, w io.WriteCloser) { + _ = r.Close() + if cw, ok := w.(interface{ CloseWithError(error) error }); ok { + _ = cw.CloseWithError(errTransientEntry) + } else { + _ = w.Close() + } + _ = fc.invalidate(ctx, key) +} + +var errTransientEntry = errors.New("cache entry was served without being stored") + // CachedStream is a wrapper around an io.ReadCloser that allows reading from a cache. type CachedStream struct { io.Reader io.Seeker io.Closer - Cached bool + Cached bool + Transient bool } func (s *CachedStream) Close() error { diff --git a/utils/cache/file_caches_test.go b/utils/cache/file_caches_test.go index 3189de6b2..70cd2a0dc 100644 --- a/utils/cache/file_caches_test.go +++ b/utils/cache/file_caches_test.go @@ -106,6 +106,67 @@ var _ = Describe("File Caches", func() { Expect(called).To(BeTrue()) }) + It("serves a transient result without storing it", func() { + var calls atomic.Int32 + fc := callNewFileCache("test", "1KB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) { + calls.Add(1) + return Transient(io.NopCloser(strings.NewReader("stand-in"))), nil + }) + + for range 2 { + s, err := fc.Get(context.Background(), &testArg{"transient"}) + Expect(err).ToNot(HaveOccurred()) + Expect(s.Transient).To(BeTrue()) + Expect(s.Cached).To(BeFalse()) + Expect(io.ReadAll(s)).To(Equal([]byte("stand-in"))) + Expect(s.Close()).To(Succeed()) + } + Expect(calls.Load()).To(Equal(int32(2))) + + dataPath := fcSpreadFS(fc).KeyMapper((&testArg{"transient"}).Key()) + _, statErr := os.Stat(dataPath + ".complete") + Expect(os.IsNotExist(statErr)).To(BeTrue()) + }) + + It("stores the next result once a transient one was served", func() { + var calls atomic.Int32 + fc := callNewFileCache("test", "1KB", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) { + if calls.Add(1) == 1 { + return Transient(io.NopCloser(strings.NewReader("stand-in"))), nil + } + return strings.NewReader("real"), nil + }) + + s, err := fc.Get(context.Background(), &testArg{"k"}) + Expect(err).ToNot(HaveOccurred()) + Expect(io.ReadAll(s)).To(Equal([]byte("stand-in"))) + _ = s.Close() + + // The transient entry is gone when Get returns, so this is a MISS that stores. + s, err = fc.Get(context.Background(), &testArg{"k"}) + Expect(err).ToNot(HaveOccurred()) + Expect(s.Transient).To(BeFalse()) + Expect(io.ReadAll(s)).To(Equal([]byte("real"))) + _ = s.Close() + + Eventually(func() bool { + s, _ := fc.Get(context.Background(), &testArg{"k"}) + defer s.Close() + return s.Cached + }).Should(BeTrue()) + Expect(calls.Load()).To(Equal(int32(2))) + }) + + It("reports a transient result when the cache is disabled", func() { + fc := callNewFileCache("test", "0", "test", 0, func(ctx context.Context, arg Item) (io.Reader, error) { + return Transient(io.NopCloser(strings.NewReader("stand-in"))), nil + }) + s, err := fc.Get(context.Background(), &testArg{"k"}) + Expect(err).ToNot(HaveOccurred()) + Expect(s.Transient).To(BeTrue()) + Expect(io.ReadAll(s)).To(Equal([]byte("stand-in"))) + }) + 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