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.
This commit is contained in:
Deluan 2026-09-14 19:36:50 -04:00
commit bc16cafb91
9 changed files with 362 additions and 9 deletions

View file

@ -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) {

View file

@ -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"}})

View file

@ -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
}

View file

@ -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 {

View file

@ -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
}

View file

@ -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
}

View file

@ -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.

View file

@ -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 {

View file

@ -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