mirror of
https://github.com/navidrome/navidrome.git
synced 2026-10-08 18:37:09 +02:00
refactor(artwork): defer animated GIF conversion via a loader error
Replace cache.Transient with an artwork-local *deferredAnimation error returned by resizedItem.Reader. The file cache already stores nothing when its loader fails, so the generic cache no longer needs a transient mode. The error carries the GIF bytes and the static stand-in, replacing the deferAnimated callback. The cache's loader-error path now cancels the entry writer, so readers that joined the miss fail with the loader's error instead of reading an empty entry, and the removal no longer blocks on them. The background job resizes the bytes it already holds instead of re-reading them, orOriginal centralizes the serve-the-original fallback, and warmCache shares the store-and-drain step with the worker's precache.
This commit is contained in:
parent
bc16cafb91
commit
c0da9fcd93
7 changed files with 73 additions and 146 deletions
|
|
@ -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))
|
||||
}()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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))
|
||||
}
|
||||
|
|
|
|||
43
utils/cache/file_caches.go
vendored
43
utils/cache/file_caches.go
vendored
|
|
@ -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 {
|
||||
|
|
|
|||
71
utils/cache/file_caches_test.go
vendored
71
utils/cache/file_caches_test.go
vendored
|
|
@ -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() {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue