mirror of
https://github.com/navidrome/navidrome.git
synced 2026-10-08 18:37:09 +02:00
Compare commits
4 commits
master
...
animated-g
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1bfe5accc8 | ||
|
|
2cc2ab1160 | ||
|
|
c0da9fcd93 | ||
|
|
bc16cafb91 |
20 changed files with 532 additions and 47 deletions
|
|
@ -89,13 +89,13 @@ func CreateSubsonicAPIRouter(ctx context.Context) *subsonic.Router {
|
|||
fileCache := artwork.GetImageCache()
|
||||
imageStore := artwork.GetImageStore()
|
||||
fFmpeg := ffmpeg.New()
|
||||
artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg)
|
||||
broker := events.GetBroker()
|
||||
artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg, broker)
|
||||
transcodingCache := stream.GetTranscodingCache()
|
||||
mediaStreamer := stream.NewMediaStreamer(dataStore, fFmpeg, transcodingCache)
|
||||
share := core.NewShare(dataStore)
|
||||
archiver := core.NewArchiver(mediaStreamer, dataStore, share)
|
||||
players := core.NewPlayers(dataStore)
|
||||
broker := events.GetBroker()
|
||||
metricsMetrics := metrics.GetPrometheusInstance(dataStore)
|
||||
manager := plugins.GetManager(dataStore, broker, metricsMetrics)
|
||||
agentsAgents := agents.GetAgents(dataStore, manager)
|
||||
|
|
@ -119,12 +119,12 @@ func CreateJellyfinAPIRouter(ctx context.Context) *jellyfin.Router {
|
|||
fileCache := artwork.GetImageCache()
|
||||
imageStore := artwork.GetImageStore()
|
||||
fFmpeg := ffmpeg.New()
|
||||
artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg)
|
||||
broker := events.GetBroker()
|
||||
artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg, broker)
|
||||
transcodingCache := stream.GetTranscodingCache()
|
||||
mediaStreamer := stream.NewMediaStreamer(dataStore, fFmpeg, transcodingCache)
|
||||
transcodeDecider := stream.NewTranscodeDecider(dataStore, fFmpeg)
|
||||
players := core.NewPlayers(dataStore)
|
||||
broker := events.GetBroker()
|
||||
metricsMetrics := metrics.GetPrometheusInstance(dataStore)
|
||||
manager := plugins.GetManager(dataStore, broker, metricsMetrics)
|
||||
playTracker := scrobbler.GetPlayTracker(dataStore, broker, manager)
|
||||
|
|
@ -145,7 +145,8 @@ func CreatePublicRouter() *public.Router {
|
|||
fileCache := artwork.GetImageCache()
|
||||
imageStore := artwork.GetImageStore()
|
||||
fFmpeg := ffmpeg.New()
|
||||
artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg)
|
||||
broker := events.GetBroker()
|
||||
artworkArtwork := artwork.NewArtwork(dataStore, fileCache, imageStore, fFmpeg, broker)
|
||||
transcodingCache := stream.GetTranscodingCache()
|
||||
mediaStreamer := stream.NewMediaStreamer(dataStore, fFmpeg, transcodingCache)
|
||||
share := core.NewShare(dataStore)
|
||||
|
|
|
|||
|
|
@ -16,6 +16,7 @@ import (
|
|||
"github.com/navidrome/navidrome/log"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
"github.com/navidrome/navidrome/resources"
|
||||
"github.com/navidrome/navidrome/server/events"
|
||||
"github.com/navidrome/navidrome/utils/cache"
|
||||
)
|
||||
|
||||
|
|
@ -31,6 +32,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
|
||||
|
|
@ -48,8 +50,8 @@ type Artwork interface {
|
|||
GetOrPlaceholder(ctx context.Context, id string, size int, square bool) (*Image, error)
|
||||
}
|
||||
|
||||
func NewArtwork(ds model.DataStore, cache cache.FileCache, store *ImageStore, ffm ffmpeg.FFmpeg) Artwork {
|
||||
return &service{ds: ds, cache: cache, store: store, ffmpeg: ffm}
|
||||
func NewArtwork(ds model.DataStore, cache cache.FileCache, store *ImageStore, ffm ffmpeg.FFmpeg, broker events.Broker) Artwork {
|
||||
return &service{ds: ds, cache: cache, store: store, ffmpeg: ffm, broker: broker}
|
||||
}
|
||||
|
||||
// entityExists reports whether the entity an artwork id points at is still there: state rows
|
||||
|
|
@ -85,6 +87,7 @@ type service struct {
|
|||
cache cache.FileCache
|
||||
store *ImageStore
|
||||
ffmpeg ffmpeg.FFmpeg
|
||||
broker events.Broker
|
||||
}
|
||||
|
||||
func (s *service) GetOrPlaceholder(ctx context.Context, id string, size int, square bool) (*Image, error) {
|
||||
|
|
@ -135,7 +138,7 @@ func (s *service) serveEntity(ctx context.Context, artID model.ArtworkID, size i
|
|||
|
||||
// serveSource is the one place bytes become an Image. hash is the pixel identity ("" for disc art)
|
||||
// and doubles as the full-size validator, so an ETag is only needed when resized or hash is "".
|
||||
func (s *service) serveSource(ctx context.Context, key, hash string, lastUpdate time.Time,
|
||||
func (s *service) serveSource(ctx context.Context, artID model.ArtworkID, key, hash string, lastUpdate time.Time,
|
||||
size int, square bool, open func() (io.ReadCloser, error),
|
||||
) (*Image, error) {
|
||||
if size == 0 && !square {
|
||||
|
|
@ -152,15 +155,55 @@ 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,
|
||||
})
|
||||
// Without a cache a background conversion could never be served, so convert inline instead.
|
||||
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(artID, item, deferred.data)
|
||||
return &Image{ReadCloser: io.NopCloser(deferred.standIn), Hash: hash, LastUpdated: lastUpdate, Transient: true}, nil
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
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(artID model.ArtworkID, item *resizedItem, data []byte) {
|
||||
select {
|
||||
case animConvertSlot <- struct{}{}:
|
||||
default:
|
||||
return
|
||||
}
|
||||
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()
|
||||
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.
|
||||
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
|
||||
}
|
||||
log.Debug(ctx, "Artwork: Converted animated image", "key", key, "bytes", len(out), "elapsed", time.Since(start))
|
||||
// The image URL is unchanged, so UIs holding the stand-in only reload it when told.
|
||||
if res, ok := artworkKindToResource[artID.Kind]; ok {
|
||||
s.broker.SendBroadcastMessage(ctx, (&events.RefreshResource{}).With(res, artID.ID))
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
// 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) {
|
||||
|
|
@ -175,7 +218,7 @@ func (s *service) serveHash(ctx context.Context, artID model.ArtworkID, ia *mode
|
|||
}
|
||||
return nil, err
|
||||
}
|
||||
img, err := s.serveSource(ctx, ia.Hash, ia.Hash, ia.UpdatedAt, size, square,
|
||||
img, err := s.serveSource(ctx, artID, ia.Hash, ia.Hash, ia.UpdatedAt, size, square,
|
||||
func() (io.ReadCloser, error) { return openOriginal(ia, art.Mime, s.store) })
|
||||
if err != nil {
|
||||
if errors.Is(err, context.Canceled) {
|
||||
|
|
@ -233,11 +276,11 @@ func (s *service) provisional(ctx context.Context, artID model.ArtworkID, size i
|
|||
s.enqueue(ctx, artID, model.ArtworkPriorityBump)
|
||||
log.Debug(ctx, "Artwork: Provisional read-through, no state row yet", "artID", artID,
|
||||
"source", res.source, "hit", res.reader != nil)
|
||||
return s.serveResolution(ctx, res, size, square)
|
||||
return s.serveResolution(ctx, artID, res, size, square)
|
||||
}
|
||||
|
||||
// serveResolution turns a local resolution's bytes into a servable Image (byte-hash only, no decode).
|
||||
func (s *service) serveResolution(ctx context.Context, res resolution, size int, square bool) (*Image, error) {
|
||||
func (s *service) serveResolution(ctx context.Context, artID model.ArtworkID, res resolution, size int, square bool) (*Image, error) {
|
||||
if res.reader == nil {
|
||||
return nil, ErrUnavailable
|
||||
}
|
||||
|
|
@ -251,7 +294,7 @@ func (s *service) serveResolution(ctx context.Context, res resolution, size int,
|
|||
return nil, ErrUnavailable
|
||||
}
|
||||
// Keyed by the byte-hash, so the entry lines up with the worker's eventual store entry.
|
||||
return s.serveSource(ctx, hash, hash, unixMtime(res.refMtime), size, square,
|
||||
return s.serveSource(ctx, artID, hash, hash, unixMtime(res.refMtime), size, square,
|
||||
func() (io.ReadCloser, error) { return io.NopCloser(bytes.NewReader(data)), nil })
|
||||
}
|
||||
|
||||
|
|
@ -301,7 +344,7 @@ func (s *service) provisionalEmbedded(ctx context.Context, artID model.ArtworkID
|
|||
// Eligible but unextractable: fall back the way CoverArtID does, not to a placeholder.
|
||||
return s.Get(ctx, mf.DiscCoverArtID(), size, square)
|
||||
}
|
||||
return s.serveResolution(ctx, res, size, square)
|
||||
return s.serveResolution(ctx, artID, res, size, square)
|
||||
}
|
||||
|
||||
// serveDisc reads disc art through with no state row and no enqueue, falling back to the album cover.
|
||||
|
|
@ -319,7 +362,7 @@ func (s *service) serveDisc(ctx context.Context, artID model.ArtworkID, size int
|
|||
// Disc art has no state row, hence no content hash: keying on id, album mtime and
|
||||
// DiscArtPriority lets a warm cache answer without running the chain or touching the disk.
|
||||
key := fmt.Sprintf("%s|%d|%s", artID.ID, dr.cacheTime().UnixNano(), conf.Server.DiscArtPriority)
|
||||
img, err := s.serveSource(ctx, key, "", dr.cacheTime(), size, square, selectImage)
|
||||
img, err := s.serveSource(ctx, artID, key, "", dr.cacheTime(), size, square, selectImage)
|
||||
if err != nil {
|
||||
if errors.Is(err, context.Canceled) {
|
||||
return nil, err
|
||||
|
|
|
|||
|
|
@ -3,10 +3,12 @@ package artwork
|
|||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"image"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/navidrome/navidrome/conf"
|
||||
|
|
@ -14,6 +16,7 @@ import (
|
|||
"github.com/navidrome/navidrome/consts"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
"github.com/navidrome/navidrome/resources"
|
||||
"github.com/navidrome/navidrome/server/events"
|
||||
"github.com/navidrome/navidrome/tests"
|
||||
"github.com/navidrome/navidrome/utils/cache"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
|
|
@ -107,7 +110,7 @@ var _ = Describe("Artwork", func() {
|
|||
return arg.(artworkReader).Reader(ctx)
|
||||
})
|
||||
Eventually(func() bool { return imgCache.Available(ctx) }, 10*time.Second).Should(BeTrue())
|
||||
svc = NewArtwork(ds, imgCache, store, ffm)
|
||||
svc = NewArtwork(ds, imgCache, store, ffm, events.NoopBroker())
|
||||
})
|
||||
|
||||
Describe("found state", func() {
|
||||
|
|
@ -422,6 +425,95 @@ var _ = Describe("Artwork", func() {
|
|||
})
|
||||
})
|
||||
|
||||
Describe("animated GIF", func() {
|
||||
var (
|
||||
gifBytes []byte
|
||||
fake *animFFmpeg
|
||||
broker *fakeEventBroker
|
||||
)
|
||||
|
||||
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")}
|
||||
broker = &fakeEventBroker{}
|
||||
svc = NewArtwork(ds, imgCache, store, fake, broker)
|
||||
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("tells the UI to reload the item once its conversion is cached", func() {
|
||||
readAll(get("al-al1"))
|
||||
waitForConversions()
|
||||
|
||||
sent := broker.getEvents()
|
||||
Expect(sent).To(HaveLen(1))
|
||||
Expect(sent[0].Data(sent[0])).To(MatchJSON(`{"album":["al1"]}`))
|
||||
})
|
||||
|
||||
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.Error = 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"}})
|
||||
|
|
|
|||
|
|
@ -108,7 +108,7 @@ var _ = Describe("Acquisition → serve loop", func() {
|
|||
})
|
||||
Eventually(func() bool { return imgCache.Available(ctx) }, 10*time.Second).Should(BeTrue())
|
||||
|
||||
svc = artwork.NewArtwork(ds, imgCache, store, ffm)
|
||||
svc = artwork.NewArtwork(ds, imgCache, store, ffm, events.NoopBroker())
|
||||
worker = artwork.NewWorker(ds, store, agents.GetAgents(ds, nil), ffm, events.NoopBroker(), imgCache)
|
||||
})
|
||||
|
||||
|
|
|
|||
|
|
@ -119,7 +119,7 @@ func setupResolutionHarness() {
|
|||
})
|
||||
Eventually(func() bool { return imgCache.Available(rctx) }, 10*time.Second).Should(BeTrue())
|
||||
|
||||
rsvc = artwork.NewArtwork(rds, imgCache, rstore, ffm)
|
||||
rsvc = artwork.NewArtwork(rds, imgCache, rstore, ffm, events.NoopBroker())
|
||||
rworker = artwork.NewWorker(rds, rstore, agents.GetAgents(rds, nil), ffm, events.NoopBroker(), imgCache)
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -42,6 +42,8 @@ type resizedItem struct {
|
|||
square bool
|
||||
ffmpeg ffmpeg.FFmpeg
|
||||
open func() (io.ReadCloser, error)
|
||||
// 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
|
||||
|
|
@ -64,13 +66,50 @@ func (r *resizedItem) Reader(ctx context.Context) (io.ReadCloser, error) {
|
|||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if r.deferAnimated && isAnimatedGIF(data) && r.ffmpeg.IsAvailable() {
|
||||
static, _, err := resizeStaticImage(data, r.size, r.square)
|
||||
return nil, &deferredAnimation{data: data, standIn: orOriginal(data, static, err)}
|
||||
}
|
||||
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.
|
||||
return io.NopCloser(bytes.NewReader(data)), nil
|
||||
}
|
||||
if rc, ok := resized.(io.ReadCloser); ok {
|
||||
return rc, nil
|
||||
}
|
||||
return io.NopCloser(resized), nil
|
||||
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 {
|
||||
return bytes.NewReader(data)
|
||||
}
|
||||
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
|
||||
}
|
||||
defer stream.Close()
|
||||
_, err = io.Copy(io.Discard, stream)
|
||||
return err
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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,62 @@ 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) {
|
||||
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{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("")}),
|
||||
)
|
||||
})
|
||||
|
||||
// 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
|
||||
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.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
|
||||
}
|
||||
return io.NopCloser(bytes.NewReader(f.out)), nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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))
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -132,7 +132,7 @@ var _ = Describe("Artwork Serving", Ordered, func() {
|
|||
store := artwork.NewImageStore(GinkgoT().TempDir())
|
||||
imgCache := newDummyImageCache(ctx)
|
||||
ffm := harness.NoopFFmpeg{}
|
||||
artSvc = artwork.NewArtwork(ds, imgCache, store, ffm)
|
||||
artSvc = artwork.NewArtwork(ds, imgCache, store, ffm, events.NoopBroker())
|
||||
worker = artwork.NewWorker(ds, store, agents.GetAgents(ds, nil), ffm, events.NoopBroker(), imgCache)
|
||||
|
||||
artRouter = buildArtworkRouter(artSvc)
|
||||
|
|
|
|||
|
|
@ -5,6 +5,7 @@ import { makeStyles } from '@material-ui/core/styles'
|
|||
import config from '../config'
|
||||
import subsonic from '../subsonic'
|
||||
import { useImageUrl } from './useImageUrl'
|
||||
import { useArtworkRefresh } from './useArtworkRefresh'
|
||||
import { ThumbHashCanvas } from './ThumbHashCanvas'
|
||||
|
||||
// Drives both the CSS transition and the timer that retires the placeholder, so they cannot drift.
|
||||
|
|
@ -47,7 +48,8 @@ export const Artwork = ({
|
|||
}) => {
|
||||
const classes = useStyles()
|
||||
const url = record ? subsonic.getCoverArtUrl(record, size, square) : ''
|
||||
const { imgUrl, fromCache } = useImageUrl(url)
|
||||
const version = useArtworkRefresh(record?.id)
|
||||
const { imgUrl, fromCache } = useImageUrl(url, version)
|
||||
|
||||
const [decoded, setDecoded] = useState(false)
|
||||
const [faded, setFaded] = useState(false)
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ import { render, fireEvent, act } from '@testing-library/react'
|
|||
import { describe, it, expect, vi, beforeEach } from 'vitest'
|
||||
|
||||
vi.mock('./useImageUrl', () => ({ useImageUrl: vi.fn() }))
|
||||
vi.mock('./useArtworkRefresh', () => ({ useArtworkRefresh: vi.fn() }))
|
||||
vi.mock('../subsonic', () => ({
|
||||
default: {
|
||||
// Mirrors the real hash suffix, so a refreshed record yields a different URL.
|
||||
|
|
@ -13,6 +14,7 @@ vi.mock('../subsonic', () => ({
|
|||
vi.mock('../config', () => ({ default: { uiCoverArtSize: 300 } }))
|
||||
|
||||
import { useImageUrl } from './useImageUrl'
|
||||
import { useArtworkRefresh } from './useArtworkRefresh'
|
||||
import { Artwork } from './Artwork'
|
||||
|
||||
const withArt = {
|
||||
|
|
@ -28,6 +30,14 @@ describe('Artwork', () => {
|
|||
HTMLCanvasElement.prototype.getContext = vi.fn(() => null)
|
||||
})
|
||||
|
||||
it('loads the image under the version of the last refresh naming its record', () => {
|
||||
useImageUrl.mockReturnValue({ imgUrl: null, loading: true })
|
||||
useArtworkRefresh.mockReturnValue(1234)
|
||||
render(<Artwork record={withArt} />)
|
||||
expect(useArtworkRefresh).toHaveBeenCalledWith('al-1')
|
||||
expect(useImageUrl).toHaveBeenCalledWith('/rest/getCoverArt?id=al-1', 1234)
|
||||
})
|
||||
|
||||
it('renders nothing without a record', () => {
|
||||
useImageUrl.mockReturnValue({ imgUrl: null, loading: false })
|
||||
const { container } = render(<Artwork record={null} />)
|
||||
|
|
|
|||
23
ui/src/common/useArtworkRefresh.js
Normal file
23
ui/src/common/useArtworkRefresh.js
Normal file
|
|
@ -0,0 +1,23 @@
|
|||
import { useState } from 'react'
|
||||
import { useSelector } from 'react-redux'
|
||||
|
||||
const namesId = (resources, id) =>
|
||||
Object.values(resources || {}).some(
|
||||
(ids) => Array.isArray(ids) && ids.includes(id),
|
||||
)
|
||||
|
||||
// Artwork can change behind an unchanged URL (a background conversion replacing its stand-in), so
|
||||
// this returns when the server last announced new artwork for id, for useImageUrl to version by.
|
||||
export const useArtworkRefresh = (id) => {
|
||||
// A primitive, so an event naming other items does not re-render this cover.
|
||||
const named = useSelector(({ activity }) =>
|
||||
id && namesId(activity?.refresh?.resources, id)
|
||||
? activity.refresh.lastReceived
|
||||
: undefined,
|
||||
)
|
||||
const [refreshedAt, setRefreshedAt] = useState()
|
||||
if (named !== undefined && named !== refreshedAt) {
|
||||
setRefreshedAt(named)
|
||||
}
|
||||
return refreshedAt
|
||||
}
|
||||
44
ui/src/common/useArtworkRefresh.test.js
Normal file
44
ui/src/common/useArtworkRefresh.test.js
Normal file
|
|
@ -0,0 +1,44 @@
|
|||
import { renderHook } from '@testing-library/react-hooks'
|
||||
import { vi, describe, it, expect, beforeEach } from 'vitest'
|
||||
import { useSelector } from 'react-redux'
|
||||
import { useArtworkRefresh } from './useArtworkRefresh'
|
||||
|
||||
vi.mock('react-redux', () => ({ useSelector: vi.fn() }))
|
||||
|
||||
const withRefresh = (refresh) =>
|
||||
useSelector.mockImplementation((select) => select({ activity: { refresh } }))
|
||||
|
||||
describe('useArtworkRefresh', () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks()
|
||||
})
|
||||
|
||||
it('has no version before any refresh event', () => {
|
||||
withRefresh(undefined)
|
||||
const { result } = renderHook(() => useArtworkRefresh('ar-1'))
|
||||
expect(result.current).toBeUndefined()
|
||||
})
|
||||
|
||||
it('versions the artwork with the time of an event naming its id', () => {
|
||||
withRefresh({ lastReceived: 100, resources: { artist: ['ar-1'] } })
|
||||
const { result } = renderHook(() => useArtworkRefresh('ar-1'))
|
||||
expect(result.current).toBe(100)
|
||||
})
|
||||
|
||||
it('keeps its version when a later event names something else', () => {
|
||||
withRefresh({ lastReceived: 100, resources: { artist: ['ar-1'] } })
|
||||
const { result, rerender } = renderHook(() => useArtworkRefresh('ar-1'))
|
||||
|
||||
withRefresh({ lastReceived: 200, resources: { album: ['al-9'] } })
|
||||
rerender()
|
||||
|
||||
expect(result.current).toBe(100)
|
||||
})
|
||||
|
||||
// Honoring a wildcard would refetch every cover on screen after each scan.
|
||||
it('ignores wildcard events', () => {
|
||||
withRefresh({ lastReceived: 100, resources: { '*': '*', artist: ['*'] } })
|
||||
const { result } = renderHook(() => useArtworkRefresh('ar-1'))
|
||||
expect(result.current).toBeUndefined()
|
||||
})
|
||||
})
|
||||
|
|
@ -34,9 +34,10 @@ const evictIfNeeded = () => {
|
|||
/**
|
||||
* Loads an image via fetch() with AbortController so that in-flight requests
|
||||
* are canceled on unmount (e.g., during pagination). Uses a module-level cache
|
||||
* so remounting returns the cached blob URL instantly.
|
||||
* so remounting returns the cached blob URL instantly. A no-store response is a stand-in: it is
|
||||
* re-fetched on remount and on a new version, and painted until its replacement arrives.
|
||||
*/
|
||||
export const useImageUrl = (url) => {
|
||||
export const useImageUrl = (url, version) => {
|
||||
const cached = url ? cache.get(url) : null
|
||||
const [imgUrl, setImgUrl] = useState(cached?.blobUrl || null)
|
||||
const [loading, setLoading] = useState(!!url && !cached)
|
||||
|
|
@ -68,7 +69,7 @@ export const useImageUrl = (url) => {
|
|||
// Re-check: another component's effect may have populated the cache
|
||||
// between this component's render and effect execution.
|
||||
const entry = cache.get(url)
|
||||
if (entry) {
|
||||
if (entry && !entry.transient) {
|
||||
entry.refCount++
|
||||
setImgUrl(entry.blobUrl)
|
||||
setLoading(false)
|
||||
|
|
@ -77,12 +78,15 @@ export const useImageUrl = (url) => {
|
|||
entry.refCount--
|
||||
}
|
||||
}
|
||||
const standIn = !!entry
|
||||
|
||||
const controller = new AbortController()
|
||||
let queued = true
|
||||
setImgUrl(null)
|
||||
setLoading(true)
|
||||
setError(false)
|
||||
if (!standIn) {
|
||||
setImgUrl(null)
|
||||
setLoading(true)
|
||||
setError(false)
|
||||
}
|
||||
|
||||
const doFetch = () => {
|
||||
queued = false
|
||||
|
|
@ -92,9 +96,12 @@ export const useImageUrl = (url) => {
|
|||
if (!res.ok) {
|
||||
throw new Error(`HTTP ${res.status}`)
|
||||
}
|
||||
return res.blob()
|
||||
const transient = !!res.headers
|
||||
?.get('Cache-Control')
|
||||
?.includes('no-store')
|
||||
return res.blob().then((blob) => ({ blob, transient }))
|
||||
})
|
||||
.then((blob) => {
|
||||
.then(({ blob, transient }) => {
|
||||
activeFetches--
|
||||
processQueue()
|
||||
// Guard against late resolution after abort
|
||||
|
|
@ -105,12 +112,15 @@ export const useImageUrl = (url) => {
|
|||
// Handle concurrent fetches: if another component already cached
|
||||
// this URL, use its entry and discard our blob.
|
||||
const existing = cache.get(url)
|
||||
if (existing && existing.blobUrl) {
|
||||
if (existing?.blobUrl && !existing.transient) {
|
||||
existing.refCount++
|
||||
URL.revokeObjectURL(objectUrl)
|
||||
setImgUrl(existing.blobUrl)
|
||||
} else {
|
||||
cache.set(url, { blobUrl: objectUrl, refCount: 1 })
|
||||
if (existing?.blobUrl && existing.refCount <= 0) {
|
||||
URL.revokeObjectURL(existing.blobUrl)
|
||||
}
|
||||
cache.set(url, { blobUrl: objectUrl, refCount: 1, transient })
|
||||
evictIfNeeded()
|
||||
setImgUrl(objectUrl)
|
||||
}
|
||||
|
|
@ -122,6 +132,9 @@ export const useImageUrl = (url) => {
|
|||
if (err.name === 'AbortError') {
|
||||
return // Expected on unmount or URL change
|
||||
}
|
||||
if (standIn) {
|
||||
return // Keep painting the stand-in
|
||||
}
|
||||
// Cache the error so repeated mounts don't re-fetch broken URLs
|
||||
cache.set(url, { blobUrl: null, error: true, refCount: 0 })
|
||||
setError(true)
|
||||
|
|
@ -150,7 +163,7 @@ export const useImageUrl = (url) => {
|
|||
entry.refCount--
|
||||
}
|
||||
}
|
||||
}, [url])
|
||||
}, [url, version])
|
||||
|
||||
return { imgUrl, loading, error, fromCache }
|
||||
}
|
||||
|
|
|
|||
|
|
@ -234,4 +234,98 @@ describe('useImageUrl', () => {
|
|||
// A remembered failure has no blob, so callers must not treat it as instantly painted.
|
||||
expect(result2.current.fromCache).toBe(false)
|
||||
})
|
||||
|
||||
describe('no-store stand-ins', () => {
|
||||
const image = (data, cacheControl = 'public, no-cache') =>
|
||||
Promise.resolve({
|
||||
ok: true,
|
||||
headers: {
|
||||
get: (name) => (name === 'Cache-Control' ? cacheControl : null),
|
||||
},
|
||||
blob: () => Promise.resolve(new Blob([data])),
|
||||
})
|
||||
const standIn = (data) => image(data, 'no-store')
|
||||
const renderVersioned = (version) =>
|
||||
renderHook(
|
||||
({ version }) => useImageUrl('http://example.com/img.jpg', version),
|
||||
{ initialProps: { version } },
|
||||
)
|
||||
const settle = () =>
|
||||
act(async () => {
|
||||
await flushPromises()
|
||||
})
|
||||
|
||||
it('re-fetches a stand-in when its version changes, painting it meanwhile', async () => {
|
||||
let resolveFinal
|
||||
global.URL.createObjectURL = vi
|
||||
.fn()
|
||||
.mockReturnValueOnce('blob:stand-in')
|
||||
.mockReturnValueOnce('blob:final')
|
||||
global.fetch = vi
|
||||
.fn()
|
||||
.mockImplementationOnce(() => standIn('still'))
|
||||
.mockImplementationOnce(
|
||||
() => new Promise((resolve) => (resolveFinal = resolve)),
|
||||
)
|
||||
|
||||
const { result, rerender } = renderVersioned(undefined)
|
||||
await settle()
|
||||
expect(result.current.imgUrl).toBe('blob:stand-in')
|
||||
|
||||
rerender({ version: 1 })
|
||||
await settle()
|
||||
expect(global.fetch).toHaveBeenCalledTimes(2)
|
||||
expect(result.current.imgUrl).toBe('blob:stand-in')
|
||||
expect(result.current.loading).toBe(false)
|
||||
|
||||
await act(async () => {
|
||||
resolveFinal(image('animated'))
|
||||
await flushPromises()
|
||||
})
|
||||
expect(result.current.imgUrl).toBe('blob:final')
|
||||
})
|
||||
|
||||
it('keeps the stand-in when the re-fetch fails', async () => {
|
||||
global.fetch = vi
|
||||
.fn()
|
||||
.mockImplementationOnce(() => standIn('still'))
|
||||
.mockImplementationOnce(() =>
|
||||
Promise.resolve({ ok: false, status: 500 }),
|
||||
)
|
||||
|
||||
const { result, rerender } = renderVersioned(undefined)
|
||||
await settle()
|
||||
rerender({ version: 1 })
|
||||
await settle()
|
||||
|
||||
expect(global.fetch).toHaveBeenCalledTimes(2)
|
||||
expect(result.current.imgUrl).toBe('blob:mock-url')
|
||||
expect(result.current.error).toBe(false)
|
||||
})
|
||||
|
||||
// Plays, stars and ratings also name items, and must not re-download final images.
|
||||
it('does not re-fetch a final image when its version changes', async () => {
|
||||
global.fetch = vi.fn(() => image('final'))
|
||||
|
||||
const { rerender } = renderVersioned(undefined)
|
||||
await settle()
|
||||
rerender({ version: 1 })
|
||||
await settle()
|
||||
|
||||
expect(global.fetch).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('re-fetches a stand-in on remount', async () => {
|
||||
global.fetch = vi.fn(() => standIn('still'))
|
||||
|
||||
const { unmount } = renderVersioned(undefined)
|
||||
await settle()
|
||||
unmount()
|
||||
const { result } = renderVersioned(undefined)
|
||||
await settle()
|
||||
|
||||
expect(global.fetch).toHaveBeenCalledTimes(2)
|
||||
expect(result.current.imgUrl).toBe('blob:mock-url')
|
||||
})
|
||||
})
|
||||
})
|
||||
|
|
|
|||
7
utils/cache/file_caches.go
vendored
7
utils/cache/file_caches.go
vendored
|
|
@ -178,7 +178,12 @@ 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
|
||||
}
|
||||
|
|
|
|||
24
utils/cache/file_caches_test.go
vendored
24
utils/cache/file_caches_test.go
vendored
|
|
@ -106,6 +106,30 @@ var _ = Describe("File Caches", func() {
|
|||
Expect(called).To(BeTrue())
|
||||
})
|
||||
|
||||
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) {
|
||||
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
|
||||
|
||||
joined, err := fc.Get(context.Background(), &testArg{"k"})
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
close(release)
|
||||
Eventually(firstErr).Should(Receive(MatchError("loader failed")))
|
||||
|
||||
_, err = io.ReadAll(joined)
|
||||
Expect(err).To(MatchError("loader failed"))
|
||||
_ = joined.Close()
|
||||
})
|
||||
|
||||
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
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue