diff --git a/core/artwork/artwork.go b/core/artwork/artwork.go index 1da419dff..fb67917b5 100644 --- a/core/artwork/artwork.go +++ b/core/artwork/artwork.go @@ -153,20 +153,17 @@ func (s *service) serveSource(ctx context.Context, key, hash string, lastUpdate } return img, nil } - 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 } - } + 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(item, deferred.data) + return &Image{ReadCloser: io.NopCloser(deferred.standIn), Hash: hash, LastUpdated: lastUpdate, Transient: true}, nil + } 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 } @@ -184,37 +181,19 @@ func (s *service) convertInBackground(item *resizedItem, data []byte) { default: return } - conv := *item - conv.deferAnimated = nil - conv.open = func() (io.ReadCloser, error) { return io.NopCloser(bytes.NewReader(data)), nil } - key := conv.Key() + 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() - 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 - } + 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. - stream, err := s.cache.Get(ctx, &storedItem{key: key, data: out}) - if err != nil { + 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 } - 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)) }() } diff --git a/core/artwork/artwork_test.go b/core/artwork/artwork_test.go index 707eb3e89..bbbcdf43d 100644 --- a/core/artwork/artwork_test.go +++ b/core/artwork/artwork_test.go @@ -491,7 +491,7 @@ var _ = Describe("Artwork", func() { }) It("caches the static fallback when the conversion fails, so it is not retried", func() { - fake.err = errors.New("pipe:0: Input/output error") + fake.Error = errors.New("pipe:0: Input/output error") readAll(get("al-al1")) waitForConversions() diff --git a/core/artwork/image_cache.go b/core/artwork/image_cache.go index 53eeb26f5..7d58b56f1 100644 --- a/core/artwork/image_cache.go +++ b/core/artwork/image_cache.go @@ -42,8 +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) + // 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 @@ -66,23 +66,40 @@ 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) + if r.deferAnimated && isAnimatedGIF(data) && r.ffmpeg.IsAvailable() { static, _, err := resizeStaticImage(data, r.size, r.square) - if err != nil || static == nil { - static = bytes.NewReader(data) - } - return cache.Transient(io.NopCloser(static)), nil + return nil, &deferredAnimation{data: data, standIn: orOriginal(data, static, err)} } resized, _, err := resizeImageData(ctx, r.ffmpeg, data, r.size, r.square) + 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 { - // Resize failed or image already within bounds: serve the original bytes. - return io.NopCloser(bytes.NewReader(data)), nil + return bytes.NewReader(data) } - if rc, ok := resized.(io.ReadCloser); ok { - return rc, nil + 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 } - return io.NopCloser(resized), nil + defer stream.Close() + _, err = io.Copy(io.Discard, stream) + return err } // storedItem caches bytes produced outside the cache. diff --git a/core/artwork/resize_test.go b/core/artwork/resize_test.go index d89ee0f3b..73463f901 100644 --- a/core/artwork/resize_test.go +++ b/core/artwork/resize_test.go @@ -41,7 +41,6 @@ var _ = Describe("resizeImageData", func() { 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) @@ -50,9 +49,9 @@ var _ = Describe("resizeImageData", func() { 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{}), + 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("")}), ) }) @@ -62,7 +61,6 @@ type animFFmpeg struct { calls atomic.Int32 release chan struct{} out []byte - err error streamErr error } @@ -75,8 +73,8 @@ func (f *animFFmpeg) ConvertAnimatedImage(ctx context.Context, _ io.Reader, _ in return nil, ctx.Err() } } - if f.err != nil { - return nil, f.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 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/utils/cache/file_caches.go b/utils/cache/file_caches.go index 90764f326..91cf87f10 100644 --- a/utils/cache/file_caches.go +++ b/utils/cache/file_caches.go @@ -25,18 +25,6 @@ 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) @@ -167,8 +155,7 @@ func (fc *fileCache) Get(ctx context.Context, arg Item) (*CachedStream, error) { if err != nil { return nil, err } - reader, transient := unwrapTransient(reader) - return &CachedStream{Reader: reader, Transient: transient}, nil + return &CachedStream{Reader: reader}, nil } key := arg.Key() @@ -191,14 +178,15 @@ 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 } - 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) @@ -230,27 +218,12 @@ 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 - Transient bool + Cached bool } func (s *CachedStream) Close() error { diff --git a/utils/cache/file_caches_test.go b/utils/cache/file_caches_test.go index 70cd2a0dc..5547db736 100644 --- a/utils/cache/file_caches_test.go +++ b/utils/cache/file_caches_test.go @@ -106,65 +106,28 @@ var _ = Describe("File Caches", func() { Expect(called).To(BeTrue()) }) - It("serves a transient result without storing it", func() { - var calls atomic.Int32 + 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) { - calls.Add(1) - return Transient(io.NopCloser(strings.NewReader("stand-in"))), nil + 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 - 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"}) + joined, err := fc.Get(context.Background(), &testArg{"k"}) Expect(err).ToNot(HaveOccurred()) - Expect(io.ReadAll(s)).To(Equal([]byte("stand-in"))) - _ = s.Close() + close(release) + Eventually(firstErr).Should(Receive(MatchError("loader failed"))) - // 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"))) + _, err = io.ReadAll(joined) + Expect(err).To(MatchError("loader failed")) + _ = joined.Close() }) It("writes a completion marker after a successful cache write", func() {