From cf0264412ba9ccad453ca230ab36784cc5711092 Mon Sep 17 00:00:00 2001 From: Deluan Date: Sat, 25 Jul 2026 11:47:40 -0400 Subject: [PATCH] perf(artwork): precache from the bytes just acquired Warming the resize cache re-read the two rows and the file the acquisition had just written, so every acquired image cost two extra queries and a second full read of a file whose bytes were still in memory. processItem now hands back what it persisted and precache warms from that, under the same cache key the serving path computes. Resolving the admin user also moves behind the empty-queue check: it is needed only to resolve private playlists, so an idle server no longer runs a user lookup on every poll. --- core/artwork/processor.go | 41 ++++++++++++++++---------- core/artwork/processor_test.go | 30 ++++++++++--------- core/artwork/worker.go | 50 +++++++++++++++----------------- core/artwork/worker_soak_test.go | 2 +- core/artwork/worker_test.go | 27 +++++++++++++++++ 5 files changed, 93 insertions(+), 57 deletions(-) diff --git a/core/artwork/processor.go b/core/artwork/processor.go index ab3bf6e5d..89ef24c3a 100644 --- a/core/artwork/processor.go +++ b/core/artwork/processor.go @@ -55,37 +55,45 @@ type workerDeps struct { gate gateFunc } +// acquired is what processItem persisted, handed back so the caller can warm the resize +// cache without re-reading the rows and the file it just wrote. +type acquired struct { + ia *model.ItemArtwork + mime string + data []byte +} + // processItem resolves one queue item end to end: find an image, hash/decode/ // blurhash it, place its bytes, and persist the resulting state. -func processItem(ctx context.Context, deps *workerDeps, item model.ArtworkQueueItem) outcome { +func processItem(ctx context.Context, deps *workerDeps, item model.ArtworkQueueItem) (outcome, *acquired) { repo := deps.ds.Artwork(ctx) res, err := resolveItem(ctx, deps.ds, deps.agents, deps.ffmpeg, item, deps.gate) if err != nil { log.Warn(ctx, "artwork: could not resolve item", "kind", item.ItemKind, "id", item.ItemID, err) - return outcomeFailed + return outcomeFailed, nil } if res.reader == nil { if res.extError || res.localError { // A source errored/timed out rather than answering "no image": never settle on // absent, keep serving old state. - return outcomeFailed + return outcomeFailed, nil } - return writeAbsent(ctx, repo, item) + return writeAbsent(ctx, repo, item), nil } defer res.reader.Close() data, err := readCapped(res.reader) if err != nil { log.Warn(ctx, "artwork: failed to read resolved image", "kind", item.ItemKind, "id", item.ItemID, "source", res.source, err) - return outcomeFailed + return outcomeFailed, nil } log.Debug(ctx, "artwork: read resolved image", "kind", item.ItemKind, "id", item.ItemID, "source", res.source, "bytes", len(data)) hash, err := HashImage(bytes.NewReader(data)) if err != nil { log.Warn(ctx, "artwork: failed to hash image", "kind", item.ItemKind, "id", item.ItemID, err) - return outcomeFailed + return outcomeFailed, nil } art, err := repo.GetImage(hash) @@ -96,24 +104,24 @@ func processItem(ctx context.Context, deps *workerDeps, item model.ArtworkQueueI art, err = decodeArtwork(ctx, hash, data) if err != nil { log.Warn(ctx, "artwork: failed to decode resolved image", "kind", item.ItemKind, "id", item.ItemID, err) - return outcomeFailed + return outcomeFailed, nil } default: log.Warn(ctx, "artwork: failed to look up image hash", "kind", item.ItemKind, "id", item.ItemID, err) - return outcomeFailed + return outcomeFailed, nil } art.SizeBytes = int64(len(data)) sourcePath, refMtime, err := placeBytes(deps.store, art, res, data) if err != nil { log.Warn(ctx, "artwork: failed to write image store", "kind", item.ItemKind, "id", item.ItemID, err) - return outcomeFailed + return outcomeFailed, nil } if err := repo.PutImage(art); err != nil { log.Warn(ctx, "artwork: failed to persist artwork image", "kind", item.ItemKind, "id", item.ItemID, err) - return outcomeFailed + return outcomeFailed, nil } - if err := repo.PutItemArtwork(&model.ItemArtwork{ + ia := &model.ItemArtwork{ ItemKind: item.ItemKind, ItemID: item.ItemID, ImageType: item.ImageType, @@ -122,14 +130,17 @@ func processItem(ctx context.Context, deps *workerDeps, item model.ArtworkQueueI SourcePath: sourcePath, RefMtime: refMtime, AttemptedAt: time.Now(), - }); err != nil { + } + // PutItemArtwork stamps UpdatedAt on ia, so what it holds now matches the persisted row. + if err := repo.PutItemArtwork(ia); err != nil { log.Warn(ctx, "artwork: failed to persist item artwork state", "kind", item.ItemKind, "id", item.ItemID, err) - return outcomeFailed + return outcomeFailed, nil } + got := &acquired{ia: ia, mime: art.Mime, data: data} if res.extError { - return outcomeFoundStale + return outcomeFoundStale, got } - return outcomeFound + return outcomeFound, got } // writeAbsent records a known-absent state: every local/external source answered definitively "no". diff --git a/core/artwork/processor_test.go b/core/artwork/processor_test.go index 33675fda8..b7d9ece19 100644 --- a/core/artwork/processor_test.go +++ b/core/artwork/processor_test.go @@ -83,7 +83,7 @@ var _ = Describe("processItem", func() { {ID: "al1", Name: "Album", FolderIDs: []string{"f1"}}, }) - out := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al1"}) + out, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al1"}) Expect(out).To(Equal(outcomeFound)) ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al1", model.ImageTypePrimary) @@ -105,7 +105,7 @@ var _ = Describe("processItem", func() { }) folderRepo.result = nil - out := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al2"}) + out, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al2"}) Expect(out).To(Equal(outcomeFound)) ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al2", model.ImageTypePrimary) @@ -128,7 +128,7 @@ var _ = Describe("processItem", func() { {ID: "al3", Name: "Album"}, }) - out := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al3"}) + out, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al3"}) Expect(out).To(Equal(outcomeAbsent)) ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al3", model.ImageTypePrimary) @@ -150,7 +150,7 @@ var _ = Describe("processItem", func() { {ID: "al-io", Name: "Album", FolderIDs: []string{"f1"}}, }) - out := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al-io"}) + out, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al-io"}) Expect(out).To(Equal(outcomeFailed)) _, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al-io", model.ImageTypePrimary) @@ -164,7 +164,7 @@ var _ = Describe("processItem", func() { }) imageAgents(&fakeImageAgent{name: "failAgent", err: errors.New("agent timed out")}) - out := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al4"}) + out, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al4"}) Expect(out).To(Equal(outcomeFailed)) _, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al4", model.ImageTypePrimary) @@ -182,7 +182,7 @@ var _ = Describe("processItem", func() { }) imageAgents(&fakeImageAgent{name: "failAgent", err: errors.New("agent timed out")}) - out := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "alstale"}) + out, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "alstale"}) Expect(out).To(Equal(outcomeFoundStale)) ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "alstale", model.ImageTypePrimary) @@ -203,7 +203,7 @@ var _ = Describe("processItem", func() { ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{{ID: "alext", Name: "Album"}}) imageAgents(&fakeImageAgent{name: "deezerFake", imgs: []agents.ExternalImage{{URL: srv.URL, Size: 500}}}) - out := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "alext"}) + out, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "alext"}) Expect(out).To(Equal(outcomeFound)) ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "alext", model.ImageTypePrimary) @@ -229,7 +229,7 @@ var _ = Describe("processItem", func() { {ID: "al6", Name: "Album B", FolderIDs: []string{"f1"}}, }) - out1 := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al5"}) + out1, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al5"}) Expect(out1).To(Equal(outcomeFound)) ia1, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al5", model.ImageTypePrimary) Expect(err).ToNot(HaveOccurred()) @@ -240,7 +240,7 @@ var _ = Describe("processItem", func() { poisoned.BlurHash = "SENTINEL" artRepo.Data[ia1.Hash] = poisoned - out2 := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al6"}) + out2, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al6"}) Expect(out2).To(Equal(outcomeFound)) ia2, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al6", model.ImageTypePrimary) Expect(err).ToNot(HaveOccurred()) @@ -271,7 +271,8 @@ var _ = Describe("processItem", func() { }) folderRepo.result = []model.Folder{{Path: "album-a", ImageFiles: []string{"cover.jpg"}}} - Expect(processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "alA"})).To(Equal(outcomeFound)) + outN, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "alA"}) + Expect(outN).To(Equal(outcomeFound)) iaA, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "alA", model.ImageTypePrimary) Expect(err).ToNot(HaveOccurred()) Expect(iaA.Source).To(Equal("folder")) @@ -284,7 +285,8 @@ var _ = Describe("processItem", func() { artRepo.Data[iaA.Hash] = poisoned folderRepo.result = []model.Folder{{Path: "album-b", ImageFiles: []string{"cover.jpg"}}} - Expect(processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "alB"})).To(Equal(outcomeFound)) + outN, _ = processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "alB"}) + Expect(outN).To(Equal(outcomeFound)) iaB, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "alB", model.ImageTypePrimary) Expect(err).ToNot(HaveOccurred()) Expect(iaB.Hash).To(Equal(iaA.Hash)) @@ -315,7 +317,7 @@ var _ = Describe("processItem", func() { radioRepo.Data = map[string]*model.Radio{"ra1": {ID: "ra1", Name: "Radio", UploadedImage: "ra1_test.jpg"}} ds.MockedRadio = radioRepo - out := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "ra", ItemID: "ra1"}) + out, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "ra", ItemID: "ra1"}) Expect(out).To(Equal(outcomeFailed)) _, err := artRepo.GetItemArtwork(model.KindRadioArtwork, "ra1", model.ImageTypePrimary) @@ -336,7 +338,7 @@ var _ = Describe("processItem", func() { radioRepo.Data = map[string]*model.Radio{"big": {ID: "big", Name: "Radio", UploadedImage: "big_test.jpg"}} ds.MockedRadio = radioRepo - out := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "ra", ItemID: "big"}) + out, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "ra", ItemID: "big"}) Expect(out).To(Equal(outcomeFailed)) _, err = artRepo.GetItemArtwork(model.KindRadioArtwork, "big", model.ImageTypePrimary) @@ -361,7 +363,7 @@ var _ = Describe("processItem", func() { Expect(os.WriteFile(blockedRoot, []byte("x"), 0600)).To(Succeed()) deps.store = NewImageStore(blockedRoot) - out := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al7"}) + out, _ := processItem(ctx, deps, model.ArtworkQueueItem{ItemKind: "al", ItemID: "al7"}) Expect(out).To(Equal(outcomeFailed)) _, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al7", model.ImageTypePrimary) diff --git a/core/artwork/worker.go b/core/artwork/worker.go index 989cc604b..4adf53b5e 100644 --- a/core/artwork/worker.go +++ b/core/artwork/worker.go @@ -1,6 +1,7 @@ package artwork import ( + "bytes" "context" "errors" "io" @@ -124,9 +125,6 @@ func (w *Worker) RunPrune(ctx context.Context) error { } func (w *Worker) drain(ctx context.Context, concurrency int) (int, error) { - // Resolved per drain, not once in Run: the worker starts at boot, possibly before any - // admin exists, so a late-created admin is picked up on the next poll (private playlists). - ctx = auth.WithAdminUser(ctx, w.deps.ds) batch, err := w.deps.ds.ArtworkQueue(ctx).DequeueBatch(2 * concurrency) if err != nil { return 0, err @@ -135,6 +133,9 @@ func (w *Worker) drain(ctx context.Context, concurrency int) (int, error) { if len(items) == 0 { return 0, nil } + // Resolved only once there is work, and per drain rather than per item: the worker needs an + // admin identity for private playlists, and can start before any admin exists. + ctx = auth.WithAdminUser(ctx, w.deps.ds) sem := make(chan struct{}, concurrency) var wg sync.WaitGroup var refreshMu sync.Mutex @@ -146,7 +147,7 @@ func (w *Worker) drain(ctx context.Context, concurrency int) (int, error) { defer wg.Done() defer func() { <-sem }() defer w.release(it) - out := w.process(ctx, it) + out, got := w.process(ctx, it) // Refresh clients on any visible state change: found/foundStale (new art) and absent // (removed art — clients must drop a previously-served immutable cover). foundStale // also wrote a served state row. @@ -157,8 +158,8 @@ func (w *Worker) drain(ctx context.Context, concurrency int) (int, error) { } // Precache only actual images. Post-outcome only: the queue row was already settled // by process, so warming the resize cache here can never block or alter queue ops. - if out == outcomeFound || out == outcomeFoundStale { - w.precache(ctx, it) + if got != nil { + w.precache(ctx, got) } }(item) } @@ -200,12 +201,12 @@ func (w *Worker) broadcastRefresh(ctx context.Context, found []model.ArtworkQueu w.broker.SendBroadcastMessage(ctx, event) } -func (w *Worker) process(ctx context.Context, item model.ArtworkQueueItem) outcome { +func (w *Worker) process(ctx context.Context, item model.ArtworkQueueItem) (outcome, *acquired) { if item.ImageType == "" { item.ImageType = model.ImageTypePrimary } w.pruneMu.RLock() - out := processItem(ctx, &w.deps, item) + out, got := processItem(ctx, &w.deps, item) w.pruneMu.RUnlock() queue := w.deps.ds.ArtworkQueue(ctx) @@ -237,7 +238,7 @@ func (w *Worker) process(ctx context.Context, item model.ArtworkQueueItem) outco log.Warn(ctx, "artwork: could not remove exhausted queue item", "kind", item.ItemKind, "id", item.ItemID, err) } } - return out + return out, got } // hasResolvedArtwork reports whether the item already has a hash-bearing state row. @@ -250,29 +251,24 @@ func (w *Worker) hasResolvedArtwork(ctx context.Context, item model.ArtworkQueue return err == nil && ia.Hash != "" } -// precache warms the resize cache for a newly-acquired image at the UI cover size, so the -// first UI request is a cache hit. Skipped when disabled; failures are debug-only. -func (w *Worker) precache(ctx context.Context, item model.ArtworkQueueItem) { +// precache warms the resize cache at the UI cover size from the bytes just acquired, so the +// first UI request is a cache hit without re-reading the rows or the file. Skipped when +// disabled; failures are debug-only. +func (w *Worker) precache(ctx context.Context, got *acquired) { if !conf.Server.EnableArtworkPrecache || w.deps.cache == nil || w.deps.cache.Disabled(ctx) { return } - imageType := item.ImageType - if imageType == "" { - imageType = model.ImageTypePrimary + // Same key as the serving path (hash/size/square); only the source of the bytes differs. + item := &resizedItem{ + hash: got.ia.Hash, + size: conf.Server.UICoverArtSize, + lastUpdate: got.ia.UpdatedAt, + ffmpeg: w.deps.ffmpeg, + open: func() (io.ReadCloser, error) { return io.NopCloser(bytes.NewReader(got.data)), nil }, } - repo := w.deps.ds.Artwork(ctx) - kind, _ := model.ParseKind(item.ItemKind) - ia, err := repo.GetItemArtwork(kind, item.ItemID, imageType) - if err != nil || ia.Hash == "" { - return - } - art, err := repo.GetImage(ia.Hash) + stream, err := w.deps.cache.Get(ctx, item) if err != nil { - return - } - stream, err := w.deps.cache.Get(ctx, newResizedItem(ia, art.Mime, conf.Server.UICoverArtSize, false, w.deps.store, w.deps.ffmpeg)) - if err != nil { - log.Debug(ctx, "artwork: precache failed", "kind", item.ItemKind, "id", item.ItemID, err) + log.Debug(ctx, "artwork: precache failed", "kind", got.ia.ItemKind, "id", got.ia.ItemID, err) return } _, _ = io.Copy(io.Discard, stream) diff --git a/core/artwork/worker_soak_test.go b/core/artwork/worker_soak_test.go index b6642d315..f728e6899 100644 --- a/core/artwork/worker_soak_test.go +++ b/core/artwork/worker_soak_test.go @@ -100,7 +100,7 @@ var _ = Describe("Worker soak", func() { start := time.Now() for i := range soakCycles { it := items[i%len(items)] - out := processItem(context.Background(), deps, it) + out, _ := processItem(context.Background(), deps, it) // "Serve-adjacent" read-back: exercise the Phase 2 surfaces a caller would // use after acquisition, not the old serving pipeline. diff --git a/core/artwork/worker_test.go b/core/artwork/worker_test.go index 5d280edcc..41032bdd4 100644 --- a/core/artwork/worker_test.go +++ b/core/artwork/worker_test.go @@ -573,6 +573,33 @@ var _ = Describe("Worker", func() { Expect(imgCache.getKeys()).To(BeEmpty()) }) + + // It warms from the bytes the acquisition already held, so no state row or store file + // needs to exist for it to work. + It("warms from the acquired bytes without reading them back", func() { + conf.Server.EnableArtworkPrecache = true + ia := &model.ItemArtwork{ + ItemKind: "al", ItemID: "unpersisted", ImageType: model.ImageTypePrimary, + Hash: "abcdef0123456789", UpdatedAt: time.Now(), + } + + data, err := os.ReadFile("tests/fixtures/artist/an-album/cover.jpg") + Expect(err).ToNot(HaveOccurred()) + + w.precache(ctx, &acquired{ia: ia, mime: "image/jpeg", data: data}) + + // Nothing backs that hash on disk or in the store, so the entry can only have come + // from the bytes handed in. Probing with a source that refuses to open proves it + // is really cached rather than re-read on demand. + probe := &resizedItem{ + hash: ia.Hash, size: 300, lastUpdate: ia.UpdatedAt, ffmpeg: ffm, + open: func() (io.ReadCloser, error) { return nil, errors.New("precache must not re-read the source") }, + } + stream, err := imgCache.Get(ctx, probe) + Expect(err).ToNot(HaveOccurred()) + defer stream.Close() + Expect(io.ReadAll(stream)).ToNot(BeEmpty()) + }) }) Describe("RunPrune", func() {