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