navidrome/core/artwork/worker_test.go
Deluan Quintão dc40bcaf80
feat(cli): add an artwork command group for diagnosing and re-driving artwork (#5957)
* feat(artwork): add a resolution chain trace collector

* feat(artwork): trace the local priority chain

* fix(artwork): record priority candidates the chain never evaluated

* refactor(artwork): report never-evaluated candidates as skipped

* feat(artwork): trace external agents at the gate seam

* feat(artwork): add repository queries to enqueue by current source

* feat(artwork): expose a tracing resolver for the CLI

* feat(artwork): read a single queue row by item

The explain CLI must report whether an item is queued, at what priority and when it
retries; the queue repository could only be drained in eligibility batches, which
cannot see a row that is still backing off.

* feat(cli): add artwork explain

Prints why an item has the artwork it has: the stored state, its queue row, the
governing config, the resolver's priority-chain walk and the verdict. Offline by
default so a diagnostic run cannot add load to an external provider; --live asks
the agents for real. Playlists and radios do not walk a priority chain, so they
report that instead of an empty chain table.

* fix(artwork): trace an external tier that never reaches an agent

A configured 'external' token vanished from the chain when no enabled agent provided
images for that entity type, and for synthetic artists, leaving the trace unable to
say whether the tier was even considered.

* fix(cli): never state an artwork outcome the walk did not observe

A transient external failure traced as 'error' fell through to 'not resolved', which
is the most common state behind a missing-artwork report. It is now indeterminate, and
an offline win that a skipped higher-priority external candidate could have taken says
so instead of naming a winner the live chain might not pick.

* feat(cli): add artwork refresh

* feat(cli): add artwork reprocess

Bulk re-enqueues artwork by kind and/or by the source an item currently
resolves from, previewing the matched count and confirming before queueing.

The preview counts with CountBySource (rows matched) and reports separately
what EnqueueBySource inserted: its DO NOTHING conflict policy leaves an
already-queued row untouched, so the two numbers differ and the output must
not claim the skipped rows were re-queued.

An unknown --source is rejected against the sources present in item_artwork,
rather than silently matching nothing and printing a reassuring 0.

* fix(cli): cover the reprocess selection rule and validate sources table-wide

The reconciliation that makes --source alone target every kind was only
exercised through runReprocess, which no test calls: mutating it to
`all := reprocessAll` left the suite green. It is now reprocessSelectsAll,
covered for all three selectors.

Scoping source validation to the selected kinds made the same well-formed
filter valid or invalid depending on which other kinds were selected, and its
error read the same for a typo as for a source that simply does not apply to
the chosen kind. Validation is now table-wide: a typo still aborts, while a
valid-but-inapplicable source falls through to "Nothing matches".

Also: the prompt now counts only the kinds that reach an external agent as
external cost, and --dry-run on an empty selection reports a dry run.

* fix(cli): cover the reprocess --yes guard and preview the external cost

Mutating the --yes check to `if true` left the suite green, so the one bypass
of the confirmation was unverified. The choice is now reprocessConfirm(yes, in),
covered in both directions.

The external estimate only reached the operator through the prompt, which
--dry-run skips — hiding the number in the one mode that exists to show it
before committing. The preview now carries it, and the prompt drops the clause
when no lookup will be made.

An empty selection says so again under --dry-run.

* feat(artwork): add read-only queue and absent counters

Both are needed by the artwork status CLI: a queue breakdown by kind and priority, and
the absent totals split against the recheck cutoff.

* feat(cli): add artwork status

Reports the queue, where artwork currently resolves from, absent counts against the 24h
recheck window, and the stored config fingerprint versus the current one — the line that
turns 'why is my server re-resolving everything?' into one command.

fingerprint() and staleAbsentAge are exported so the CLI reports the values backfill
itself compares, instead of a second copy of the formula that can silently drift.

* fix(cli): lead the artwork status backfill line with the queued backlog

By the time anyone runs a diagnostic, backfill has usually already stored the new
fingerprint, so 'up to date' was printed while thousands of items churned through external
providers. The backlog is the finding; the fingerprint is context.

Also echoes the config inputs the fingerprint covers, so a change can be traced to the
setting that caused it, and pins the rendered rows: the Absent values, the queue TOTAL and
a queue-scoped kind/priority pair were all unasserted, so kindName and priorityName were
effectively untested. FingerprintInputs is now the single listing ConfigFingerprint hashes;
a pinned hash proves the value did not change.

* refactor(artwork): export the trace outcome vocabulary

The CLI hardcoded the outcome literals and the "external:" prefix, so renaming a
constant's value in core/artwork left cmd compiling and the suite green while
`artwork explain` silently degraded its verdict.

Renaming a value now fails the golden vocabulary test in core/artwork and the
explainResult tests in cmd.

* fix(cli): keep the re-enqueue warning when a backfill is already running

A stale stored fingerprint with items already queued is the worst state the
system can be in: a second full re-enqueue is pending on top of the one running.
The line carried the weakest wording of the three, and was untested.

* refactor(artwork): drop the unreachable breaker branch from the tracing gate

--live wires the tracing gate straight to passthroughGate, so errBreakerOpen can
never reach it; the test only passed by injecting a fake gate.

* refactor(artwork): delete the never-emitted not-reached outcome

Candidates after the winner are lower priority and say nothing about why a source
won; the ones that matter sit above it and are already recorded.

* refactor(artwork): make the trace nil-safe in one place only

add already handles a nil trace, so record's own guard was dead; Steps was the
odd one out and would panic where every other method tolerates nil.

* refactor(artwork): export the trace types directly

ChainTrace and TraceStep were unexported types re-exported through aliases,
which existed only so the CLI had a name to refer to them by. The types are
public API — Resolver.Steps returns []TraceStep and the CLI constructs a
ChainTrace — so name them that way and drop the indirection.

Encapsulation is unchanged: add, mu and steps stay unexported, so only this
package can write a step.

* refactor(cli): simplify parseArtworkKind with slices.Contains

Replaces a nested loop and a manual append with slices.Contains and the
repo's slice.Map helper. Same behaviour, same error message.

* fix(cli): print the absent artwork source under the name --source accepts

`artwork explain` rendered the stored empty source as "(absent)", while
`artwork reprocess --source` only accepts "absent", so pasting what explain
printed straight back into reprocess was rejected as an unknown source.

* refactor(artwork): own the kind list and the chain predicate in the package

Export RecheckKinds and add WalksPriorityChain so the CLI stops keeping its
own copies of both, and unexport externalCandidate, which nothing outside the
package consumes.

* refactor(cli): drop the artwork command's duplicated state and formatting

Reuse artwork.RecheckKinds and artwork.WalksPriorityChain, extract
newTabWriter and externalEstimate, fold reprocessSelectsAll into
selectedKinds, and derive the queue total and the walks-chain flag instead of
carrying them in the report structs.

* test(persistence): drop two artwork-queue specs that cannot fail

One seeded hash and source together and then asserted the two counts agree,
so its setup guaranteed the result; the other repeated the count-does-not-
enqueue property already covered by the CountBySource spec.

* refactor(artwork): rename Resolver to TracingResolver for clarity

* fix(cli): count playlists in the artwork reprocess external estimate

The estimate used WalksPriorityChain, which is true only for artist and album,
so a playlist-only reprocess reported "External lookups: none" and the
confirmation prompt dropped the external-cost warning. Playlists do reach the
network: through the m3u ExternalImageURL fetch when EnableM3UExternalAlbumArt
is on, and — verified by test — through the generated grid, whose tiles resolve
album art via the full album priority chain.

Adds artwork.MayFetchExternal, a config-aware predicate for "can this kind's
resolver reach the network", and uses it for the estimate. WalksPriorityChain
keeps its separate job of deciding whether explain prints a chain block.

* fix(cli): estimate artwork reprocess external lookups per agent, not per item

The reprocess prompt billed one external lookup per externally-capable item.
fetchArtistImage/fetchAlbumImage try every enabled image agent and stop early
only on a hit, and resolvePlaylist can fetch the m3u image and then resolve up
to four sampled albums for the grid, each walking the album agents again. The
number the operator confirmed could understate real provider traffic several
fold, in the prompt whose whole job is to stop a provider flood.

ExternalLookupsPerItem now multiplies by the visible image-agent count and adds
the playlist grid factor. It stays a floor: the CLI never calls Manager.Start(),
so the plugin registry is empty and plugin-provided agents are dropped by
getEnabledAgentNames. On an install with 5 agents of which 3 are plugins the
count is well under the truth, so the wording is now "at least N" rather than
"up to N" — a zero visible count still bills one lookup for the same reason.

Fixing the plugin visibility is out of scope: Manager.Start() needs a Subsonic
router and writes to the DB via syncPlugins, breaking this command group's
read-only guarantee.

* fix(cli): state the artwork reprocess estimate as an estimate, not a bound

Neither bound is true. A ceiling is false because plugin agents are invisible to
a CLI that never starts the plugin manager, and a floor is false because a local
hit ends the walk before any agent is asked and a hit on the first agent skips
the rest. "at least N" traded one wrong claim for another.

The line now names its blind spots instead:

  External lookups: ~340 estimated (plugin agents not counted; local hits may
  need fewer).

The same line is reused in the confirmation prompt, and the zero case still
reads "External lookups: none." with the prompt dropping the clause entirely.
The count itself is unchanged.

* fix(cli): account for every configured agent in artwork explain

The Agents: line printed the raw config while the Chain only showed the agents the CLI could
construct, with nothing explaining the gap: plugin agents are never registered in a CLI that does
not start the plugin manager, and a built-in without credentials returns nil. Three of five agents
could vanish, including ones ranked above the one shown.

Also treat a live external error before the winning hit like the already-handled would-try case:
the resolver serves such a hit provisionally and retries later, so the verdict is indeterminate.

The Result line is still not qualified when an unavailable agent might have won; that needs agent
ranking, and is left to the follow-up that makes the CLI load plugin agents for real.

* fix(cli): do not call an external artwork win indeterminate

explainResult qualified the verdict whenever an external OutcomeError
appeared before the winning hit. When a later external agent returns an
image, fetchArtistImage/fetchAlbumImage discard the earlier error, so
extError is false: the worker settles the item and schedules no retry.
Telling the operator it may resolve differently on a retry was wrong.

The warning is only correct when a lower-priority local source won while
an external error was recorded, which is the case that carries extError.

* fix(cli): accept --source absent when nothing is currently absent

validateSources checks the requested sources against the ones item_artwork
actually uses, to catch a typo. The reserved empty source (spelled 'absent' on
the CLI) is a valid filter even when it matches nothing, so a scheduled
'artwork reprocess --source absent --yes' stopped working the moment the
library finished resolving. Treat it as intrinsically valid and let the
existing zero-match path report it.

* feat(artwork): explain disc and media file artwork from the CLI

`artwork explain` rejected `dc` and `mf` because it validated against RecheckKinds,
the list of kinds the backfill revisits. Those are different questions: a kind with no
recheck path still has artwork someone can report as wrong.

Disc artwork now walks DiscArtPriority under a trace, so explain reports which entry won
and why the others lost, including entries that map to no source at all (external is
unsupported, a disc with no subtitle, an album folder with no images). Media file artwork
traces its single embedded candidate, separating "EnableMediaFileCoverArt is off" from
"the track has no embedded art" — stored state cannot tell those apart.

Each command now validates against the kinds it can actually serve: explain takes all six,
refresh takes artwork.RefreshableKinds (which nativeapi now shares instead of keeping its
own copy), reprocess still takes RecheckKinds. Disc artwork stays out of refresh: the
worker cannot resolve it, so the queue row would be rejected on every drain.

WalksPriorityChain becomes Explainable, and ResolveArtist/ResolveAlbum collapse into
Resolve(kind, id).

* refactor(artwork): one disc-artwork walk for serving and explain

resolveDisc duplicated the loop selectImageReader already ran: try each source in
priority order, take the first that yields an image. The serving path and the CLI
diverged on two details as a result — only selectImageReader checked ctx between
candidates and logged each attempt.

Both now call discArtworkReader.selectImage, which takes the chainState the CLI already
uses for the other kinds. The serving path passes an untraced one, whose nil trace makes
recording a no-op. selectImageReader had no other caller and is gone.

The disc tests move from fromDiscArtPriority to discCandidates, so they assert the skip
reason for an entry that maps to no source rather than that it silently vanished, and
cancellation mid-walk is now covered.

* fix(artwork): reject a nil reader in the resize cache instead of panicking

resizedItem.Reader closes what open() hands back, so an open() that reports "no image"
as (nil, nil) rather than an error takes the request down with a nil-pointer panic. Every
caller returns an error today, and no test covered it: the resolution e2e harness stubs
the resize reader out entirely, so no e2e path reaches this code at all.

Guard it and cover Reader directly.

* refactor(artwork): move the keeps-state fact into core, drop a redundant guard

keepsArtworkState lived in package cmd and re-derived by hand what RefreshableKinds
already encodes: the same five-of-six kinds. It is now artwork.KeepsState, beside the
list, with a test pinning the two together — nothing else stopped them drifting, and a
drift would have explain report stored state for a kind that keeps none.

serveDisc's closure also hand-rolled a nil-reader error that both consumers of open()
now produce themselves: serveSource for a full-size request, resizedItem.Reader for a
resized one.

* fix(artwork): route disc candidates through the shared resolvers

openCandidate ran its own source loop and threw the error away, so a disc track that
exists but cannot be parsed traced as "miss" — indistinguishable from a track with no
embedded art. fromTag and fromFFmpegTag already report that case as errSourceUnreadable;
only this loop was discarding it. Telling those two apart is what the trace is for.

Candidates now carry a resolve func instead of raw sources: embedded goes to
resolveEmbedded, and the folder-backed entries to resolveFolderSource, extracted from
resolveFolderFile so both callers classify an unopenable file the same way. openCandidate
and its absolute-path special case go away with it.

Disc's own fromExternalFile and fromDiscSubtitle still swallow open errors, so folder
candidates cannot report unreadable yet; that is a change to their error contracts.

* fix(artwork): report an unreadable local candidate as indeterminate

processor.acquire treats resolution.localError exactly as it treats extError: a fault is
not a definitive "no image", so it retries instead of settling absent. explainResult
qualified only the external case, so a chain that ended on an unreadable local candidate
printed "not resolved" — the one verdict that says the walk was conclusive.

The qualification belongs only to the unresolved branch. chainState.try stamps extErr onto
a hit and deliberately drops localErr, so an unreadable step followed by a hit is settled
as found and must not carry a warning; a test pins that.

Found by Codex on 5f65d7cfa.
2026-08-14 21:07:56 -04:00

855 lines
31 KiB
Go

package artwork
import (
"context"
"errors"
"fmt"
"io"
"net/http"
"os"
"slices"
"sync"
"time"
"github.com/navidrome/navidrome/conf"
"github.com/navidrome/navidrome/conf/configtest"
"github.com/navidrome/navidrome/core/agents"
"github.com/navidrome/navidrome/model"
"github.com/navidrome/navidrome/server/events"
"github.com/navidrome/navidrome/tests"
"github.com/navidrome/navidrome/utils/cache"
"github.com/navidrome/navidrome/utils/slice"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"go.uber.org/goleak"
)
type recordingCache struct {
cache.FileCache
mu sync.Mutex
keys []string
disabled bool
}
func (c *recordingCache) Disabled(ctx context.Context) bool {
return c.disabled || c.FileCache.Disabled(ctx)
}
func (c *recordingCache) Get(ctx context.Context, arg cache.Item) (*cache.CachedStream, error) {
c.mu.Lock()
c.keys = append(c.keys, arg.Key())
c.mu.Unlock()
return c.FileCache.Get(ctx, arg)
}
func (c *recordingCache) getKeys() []string {
c.mu.Lock()
defer c.mu.Unlock()
return slices.Clone(c.keys)
}
// Simulates a concurrent Enqueue between DequeueBatch and the worker's delete, so
// DeleteIfUnchanged on the dequeued value no-ops.
type reenqueueOnDequeue struct {
*tests.MockArtworkQueueRepo
done bool
}
func (r *reenqueueOnDequeue) DequeueBatch(n int, kinds ...string) ([]model.ArtworkQueueItem, error) {
items, err := r.MockArtworkQueueRepo.DequeueBatch(n, kinds...)
if !r.done && len(items) > 0 {
r.done = true
for k, it := range r.Data {
if it.ItemKind == items[0].ItemKind && it.ItemID == items[0].ItemID {
it.RetryAt = items[0].RetryAt.Add(time.Minute)
r.Data[k] = it
}
}
}
return items, err
}
type fakeEventBroker struct {
http.Handler
mu sync.Mutex
events []events.Event
}
func (f *fakeEventBroker) SendMessage(_ context.Context, event events.Event) {
f.mu.Lock()
defer f.mu.Unlock()
f.events = append(f.events, event)
}
func (f *fakeEventBroker) SendBroadcastMessage(_ context.Context, event events.Event) {
f.mu.Lock()
defer f.mu.Unlock()
f.events = append(f.events, event)
}
func (f *fakeEventBroker) getEvents() []events.Event {
f.mu.Lock()
defer f.mu.Unlock()
return slices.Clone(f.events)
}
var _ events.Broker = (*fakeEventBroker)(nil)
func findQueued(q *tests.MockArtworkQueueRepo, kind, id string) *model.ArtworkQueueItem {
for _, it := range q.Data {
if it.ItemKind == kind && it.ItemID == id {
return &it
}
}
return nil
}
var _ = Describe("Worker", func() {
var (
ctx context.Context
ds *tests.MockDataStore
folderRepo *fakeFolderRepo
libRepo *tests.MockLibraryRepo
ffm *tests.MockFFmpeg
ag *agents.Agents
store *ImageStore
artRepo *tests.MockArtworkRepo
queueRepo *tests.MockArtworkQueueRepo
broker *fakeEventBroker
imgCache *recordingCache
repoRoot string
w *Worker
)
BeforeEach(func() {
DeferCleanup(configtest.SetupConfig())
ctx = context.Background()
var err error
repoRoot, err = os.Getwd()
Expect(err).ToNot(HaveOccurred())
conf.Server.CacheFolder = conf.NewDir(GinkgoT().TempDir())
folderRepo = &fakeFolderRepo{}
libRepo = &tests.MockLibraryRepo{}
libRepo.SetData(model.Libraries{{ID: 0, Path: testFileLibPath(repoRoot)}})
ffm = tests.NewMockFFmpeg("")
ag = agents.GetAgents(&tests.MockDataStore{}, nil)
artRepo = tests.CreateMockArtworkRepo()
queueRepo = tests.CreateMockArtworkQueueRepo()
ds = &tests.MockDataStore{
MockedFolder: folderRepo,
MockedLibrary: libRepo,
MockedArtwork: artRepo,
MockedArtworkQueue: queueRepo,
}
ds.MockedAlbum = tests.CreateMockAlbumRepo()
store = NewImageStore(GinkgoT().TempDir())
conf.Server.CoverArtPriority = "cover.jpg, embedded"
conf.Server.DevArtworkExternalMaxRPS = 1000 // keep the limiter out of the way of behavior tests
broker = &fakeEventBroker{}
imgCache = &recordingCache{FileCache: cache.NewFileCache("WorkerTest", "100MB", "images", 0,
func(ctx context.Context, arg cache.Item) (io.Reader, error) {
return arg.(artworkReader).Reader(ctx)
})}
// Init walks the cache dir on a goroutine; loaded CI runners can take >1s.
Eventually(func() bool { return imgCache.Available(ctx) }, 10*time.Second).Should(BeTrue())
w = NewWorker(ds, store, ag, ffm, broker, imgCache)
})
Describe("drain", func() {
It("processes a seeded queue item and removes it from the queue", func() {
folderRepo.result = []model.Folder{{
Path: "tests/fixtures/artist/an-album",
ImageFiles: []string{"cover.jpg"},
}}
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{
{ID: "al1", Name: "Album", FolderIDs: []string{"f1"}},
})
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{
ItemKind: "al", ItemID: "al1", Priority: model.ArtworkPriorityScan,
})).To(Succeed())
n, err := w.drain(ctx, 2)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al1", model.ImageTypePrimary)
Expect(err).ToNot(HaveOccurred())
Expect(ia.Source).To(Equal("folder"))
count, err := queueRepo.Count()
Expect(err).ToNot(HaveOccurred())
Expect(count).To(BeZero(), "a found item must be deleted from the queue")
})
It("processes an mf queue item, writing state and storing embedded bytes", func() {
conf.Server.EnableMediaFileCoverArt = true
ds.MockedMediaFile = tests.CreateMockMediaFileRepo()
ds.MockedMediaFile.(*tests.MockMediaFileRepo).SetData(model.MediaFiles{
{ID: "mf1", LibraryID: 0, Path: "tests/fixtures/artist/an-album/test.mp3", HasCoverArt: true},
})
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{
ItemKind: "mf", ItemID: "mf1", Priority: model.ArtworkPriorityBump,
})).To(Succeed())
n, err := w.drain(ctx, 2)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
ia, err := artRepo.GetItemArtwork(model.KindMediaFileArtwork, "mf1", model.ImageTypePrimary)
Expect(err).ToNot(HaveOccurred())
Expect(ia.Source).To(Equal("embedded"))
Expect(ia.Hash).ToNot(BeEmpty())
art, err := artRepo.GetImage(ia.Hash)
Expect(err).ToNot(HaveOccurred())
r, err := store.Open(ia.Hash, art.Mime)
Expect(err).ToNot(HaveOccurred())
defer r.Close()
data, err := io.ReadAll(r)
Expect(err).ToNot(HaveOccurred())
Expect(data).ToNot(BeEmpty(), "embedded bytes must be written to the store")
count, err := queueRepo.Count()
Expect(err).ToNot(HaveOccurred())
Expect(count).To(BeZero())
})
It("reschedules a failed item via MarkFailed with a backed-off retry_at", func() {
conf.Server.CoverArtPriority = "external"
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{{ID: "al4", Name: "Album"}})
imageAgents(&fakeImageAgent{name: "failAgent", err: errors.New("agent timed out")})
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "al", ItemID: "al4"})).To(Succeed())
n, err := w.drain(ctx, 2)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
it := findQueued(queueRepo, "al", "al4")
Expect(it).ToNot(BeNil())
Expect(it.Attempts).To(Equal(1))
Expect(it.RetryAt).To(BeTemporally(">", time.Now()))
_, err = artRepo.GetItemArtwork(model.KindAlbumArtwork, "al4", model.ImageTypePrimary)
Expect(err).To(MatchError(model.ErrNotFound), "a timeout must never settle on absent")
})
It("reschedules a found-stale item via MarkFailed while keeping its served state", func() {
conf.Server.CoverArtPriority = "external, cover.jpg"
folderRepo.result = []model.Folder{{
Path: "tests/fixtures/artist/an-album",
ImageFiles: []string{"cover.jpg"},
}}
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{
{ID: "alstale", Name: "Album", FolderIDs: []string{"f1"}},
})
imageAgents(&fakeImageAgent{name: "failAgent", err: errors.New("agent timed out")})
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "al", ItemID: "alstale"})).To(Succeed())
n, err := w.drain(ctx, 2)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
it := findQueued(queueRepo, "al", "alstale")
Expect(it).ToNot(BeNil(), "a found-stale row must survive for a higher-priority retry")
Expect(it.Attempts).To(Equal(1))
Expect(it.RetryAt).To(BeTemporally(">", time.Now()))
ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "alstale", model.ImageTypePrimary)
Expect(err).ToNot(HaveOccurred())
Expect(ia.Source).To(Equal("folder"), "the fallback art is served meanwhile")
evts := broker.getEvents()
Expect(evts).To(HaveLen(1), "the served fallback art must live-refresh the UI")
Expect(evts[0].(*events.RefreshResource).Data(evts[0])).To(ContainSubstring("alstale"))
})
It("keeps a row re-enqueued between dequeue and delete", func() {
folderRepo.result = []model.Folder{{
Path: "tests/fixtures/artist/an-album",
ImageFiles: []string{"cover.jpg"},
}}
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{
{ID: "al7", Name: "Album", FolderIDs: []string{"f1"}},
})
racing := &reenqueueOnDequeue{MockArtworkQueueRepo: queueRepo}
ds.MockedArtworkQueue = racing
w = NewWorker(ds, store, ag, ffm, broker, imgCache)
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{
ItemKind: "al", ItemID: "al7", Priority: model.ArtworkPriorityScan,
})).To(Succeed())
n, err := w.drain(ctx, 1)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
// The concurrent re-enqueue changed retry_at, so the found-path delete was a no-op.
Expect(findQueued(queueRepo, "al", "al7")).ToNot(BeNil())
ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al7", model.ImageTypePrimary)
Expect(err).ToNot(HaveOccurred())
Expect(ia.Source).To(Equal("folder"))
})
It("keeps a fresh re-enqueue ahead of a stale failure backoff", func() {
conf.Server.CoverArtPriority = "external"
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{{ID: "al8", Name: "Album"}})
imageAgents(&fakeImageAgent{name: "failAgent", err: errors.New("agent timed out")})
racing := &reenqueueOnDequeue{MockArtworkQueueRepo: queueRepo}
ds.MockedArtworkQueue = racing
w = NewWorker(ds, store, ag, ffm, broker, imgCache)
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "al", ItemID: "al8"})).To(Succeed())
dequeued := findQueued(queueRepo, "al", "al8").RetryAt
n, err := w.drain(ctx, 1)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
// The re-enqueue reset retry_at; the failure path must not stomp it nor bump attempts.
it := findQueued(queueRepo, "al", "al8")
Expect(it).ToNot(BeNil())
Expect(it.Attempts).To(BeZero())
Expect(it.RetryAt).To(BeTemporally("==", dequeued.Add(time.Minute)))
})
It("gives up and settles absent once the retry budget is exhausted", func() {
conf.Server.CoverArtPriority = "external"
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{{ID: "al9", Name: "Album"}})
imageAgents(&fakeImageAgent{name: "failAgent", err: errors.New("agent timed out")})
w = NewWorker(ds, store, ag, ffm, broker, imgCache)
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "al", ItemID: "al9"})).To(Succeed())
// Age the row past the retry budget.
for k, v := range queueRepo.Data {
if v.ItemID == "al9" {
v.EnqueuedAt = time.Now().Add(-(giveUpAfter + time.Hour))
queueRepo.Data[k] = v
}
}
n, err := w.drain(ctx, 1)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
Expect(findQueued(queueRepo, "al", "al9")).To(BeNil())
ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al9", model.ImageTypePrimary)
Expect(err).ToNot(HaveOccurred())
Expect(ia.Hash).To(BeEmpty())
})
It("keeps already-served art when the retry budget is exhausted", func() {
conf.Server.CoverArtPriority = "external"
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{{ID: "al10", Name: "Album"}})
Expect(artRepo.PutItemArtwork(&model.ItemArtwork{
ItemKind: "al", ItemID: "al10", ImageType: model.ImageTypePrimary,
Hash: "cafebabe", Source: "external:lastfm",
})).To(Succeed())
imageAgents(&fakeImageAgent{name: "failAgent", err: errors.New("agent timed out")})
w = NewWorker(ds, store, ag, ffm, broker, imgCache)
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "al", ItemID: "al10"})).To(Succeed())
for k, v := range queueRepo.Data {
if v.ItemID == "al10" {
v.EnqueuedAt = time.Now().Add(-(giveUpAfter + time.Hour))
queueRepo.Data[k] = v
}
}
n, err := w.drain(ctx, 1)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
Expect(findQueued(queueRepo, "al", "al10")).To(BeNil())
ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al10", model.ImageTypePrimary)
Expect(err).ToNot(HaveOccurred())
Expect(ia.Hash).To(Equal("cafebabe"), "a persistent outage must not discard served art")
})
// Media files are excluded from RecheckKinds, so an absent row here would never be
// revisited: a transient read error would look permanent.
It("does not settle absent on exhaustion for a kind with no recheck path", func() {
conf.Server.EnableMediaFileCoverArt = true
ds.MockedMediaFile = tests.CreateMockMediaFileRepo()
ds.MockedMediaFile.(*tests.MockMediaFileRepo).SetData(model.MediaFiles{
{ID: "mfX", LibraryID: 0, Path: "tests/fixtures/artist/an-album/gone.mp3", HasCoverArt: true},
})
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "mf", ItemID: "mfX"})).To(Succeed())
for k, v := range queueRepo.Data {
if v.ItemID == "mfX" {
v.EnqueuedAt = time.Now().Add(-(giveUpAfter + time.Hour))
queueRepo.Data[k] = v
}
}
n, err := w.drain(ctx, 1)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
Expect(findQueued(queueRepo, "mf", "mfX")).To(BeNil(), "the row must stop retrying")
_, err = artRepo.GetItemArtwork(model.KindMediaFileArtwork, "mfX", model.ImageTypePrimary)
Expect(err).To(MatchError(model.ErrNotFound),
"no row leaves the track unresolved, so a later view can still recover it")
})
It("resolves a private playlist under an admin context instead of failing forever", func() {
ds.MockedUser = adminUserRepo()
vds := &visibilityPlaylistDS{
MockDataStore: ds,
private: model.Playlist{ID: "plPriv", OwnerID: "admin"},
tracks: &tests.MockPlaylistTrackRepo{},
}
w = NewWorker(vds, store, ag, ffm, broker, imgCache)
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "pl", ItemID: "plPriv"})).To(Succeed())
n, err := w.drain(ctx, 1)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
Expect(findQueued(queueRepo, "pl", "plPriv")).To(BeNil())
ia, err := artRepo.GetItemArtwork(model.KindPlaylistArtwork, "plPriv", model.ImageTypePrimary)
Expect(err).ToNot(HaveOccurred())
Expect(ia.Hash).To(BeEmpty())
})
It("returns zero when the queue is empty", func() {
n, err := w.drain(ctx, 2)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(BeZero())
})
It("broadcasts a single refresh event for the found items in a batch", func() {
folderRepo.result = []model.Folder{{
Path: "tests/fixtures/artist/an-album",
ImageFiles: []string{"cover.jpg"},
}}
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{
{ID: "al1", Name: "Album 1", FolderIDs: []string{"f1"}},
{ID: "al2", Name: "Album 2", FolderIDs: []string{"f1"}},
})
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "al", ItemID: "al1", Priority: model.ArtworkPriorityScan})).To(Succeed())
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "al", ItemID: "al2", Priority: model.ArtworkPriorityScan})).To(Succeed())
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "ar", ItemID: "ar1", Priority: model.ArtworkPriorityScan})).To(Succeed())
n, err := w.drain(ctx, 3)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(3))
evts := broker.getEvents()
Expect(evts).To(HaveLen(1), "exactly one coalesced event per drain batch")
rr, ok := evts[0].(*events.RefreshResource)
Expect(ok).To(BeTrue())
data := rr.Data(rr)
Expect(data).To(ContainSubstring(`"album"`))
Expect(data).To(ContainSubstring("al1"))
Expect(data).To(ContainSubstring("al2"))
Expect(data).ToNot(ContainSubstring("artist"), "a failed (unresolved) artist must not be refreshed")
Expect(data).ToNot(ContainSubstring("ar1"))
})
DescribeTable("only lists the kinds it actually resolved",
func(kinds []string, wantSong bool) {
items := slice.Map(kinds, func(k string) model.ArtworkQueueItem {
return model.ArtworkQueueItem{ItemKind: k, ItemID: k + "1"}
})
w.broadcastRefresh(ctx, items)
evts := broker.getEvents()
Expect(evts).To(HaveLen(1))
data := evts[0].(*events.RefreshResource).Data(evts[0])
if wantSong {
Expect(data).To(ContainSubstring(`"song"`))
} else {
Expect(data).ToNot(ContainSubstring(`"song"`))
}
},
// An album's tracks inherit its art, but the dependent ids are unbounded: the client
// fans an album refresh out to the tracks it has loaded.
Entry("album alone does not name songs", []string{"al"}, false),
Entry("artist alone does not", []string{"ar"}, false),
Entry("playlist alone does not", []string{"pl"}, false),
Entry("album mixed with others still does not", []string{"ar", "al"}, false),
Entry("songs resolving on their own are listed by id", []string{"mf"}, true),
)
It("broadcasts a refresh when an item resolves to absent (removed cover)", func() {
conf.Server.CoverArtPriority = "cover.*" // local-only; no folder image → absent
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{{ID: "al3", Name: "Artless"}})
folderRepo.result = nil
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "al", ItemID: "al3", Priority: model.ArtworkPriorityScan})).To(Succeed())
n, err := w.drain(ctx, 2)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
evts := broker.getEvents()
Expect(evts).To(HaveLen(1), "a removed cover must live-refresh clients so they drop it")
Expect(evts[0].(*events.RefreshResource).Data(evts[0])).To(ContainSubstring("al3"))
ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "al3", model.ImageTypePrimary)
Expect(err).ToNot(HaveOccurred())
Expect(ia.Hash).To(BeEmpty(), "the outcome was absent, not found")
})
It("does not broadcast when no item is found", func() {
conf.Server.CoverArtPriority = "external"
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{{ID: "alx", Name: "Album"}})
imageAgents(&fakeImageAgent{name: "failAgent", err: errors.New("agent timed out")})
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{ItemKind: "al", ItemID: "alx"})).To(Succeed())
n, err := w.drain(ctx, 2)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
Expect(broker.getEvents()).To(BeEmpty(), "a drain with no found items sends no event")
})
})
Describe("gate/breaker", func() {
It("opens after 5 consecutive external errors and short-circuits the step", func() {
var calls int
failing := func() (io.ReadCloser, string, error) {
calls++
return nil, "", errors.New("boom")
}
for range 5 {
_, _, err := w.gate("A", failing)
Expect(err).To(HaveOccurred())
}
Expect(calls).To(Equal(5))
_, _, err := w.gate("A", failing)
Expect(err).To(MatchError(errBreakerOpen))
Expect(calls).To(Equal(5), "an open breaker must not call the external step")
})
It("resets the failure count on a successful call", func() {
failing := func() (io.ReadCloser, string, error) { return nil, "", errors.New("boom") }
ok := func() (io.ReadCloser, string, error) { return io.NopCloser(nil), "p", nil }
for range 4 {
_, _, _ = w.gate("A", failing)
}
_, _, err := w.gate("A", ok)
Expect(err).ToNot(HaveOccurred())
var calls int
counting := func() (io.ReadCloser, string, error) {
calls++
return nil, "", errors.New("boom")
}
for range 5 {
_, _, _ = w.gate("A", counting)
}
Expect(calls).To(Equal(5), "the breaker should have re-closed after the success")
})
It("does not open the breaker when the run is cancelled", func() {
cancelled := func() (io.ReadCloser, string, error) { return nil, "", context.Canceled }
for range breakerThreshold + 3 {
_, _, err := w.gate("A", cancelled)
Expect(err).To(MatchError(context.Canceled), "a cancellation passes through, never errBreakerOpen")
}
var calls int
counting := func() (io.ReadCloser, string, error) {
calls++
return nil, "", errors.New("boom")
}
_, _, _ = w.gate("A", counting)
Expect(calls).To(Equal(1), "the breaker stayed closed, so the step still runs")
})
It("ignores a cancellation mid-run, neither counting nor clearing the failures", func() {
failing := func() (io.ReadCloser, string, error) { return nil, "", errors.New("boom") }
cancelled := func() (io.ReadCloser, string, error) { return nil, "", context.Canceled }
for range breakerThreshold - 1 {
_, _, _ = w.gate("A", failing)
}
_, _, _ = w.gate("A", cancelled)
var calls int
counting := func() (io.ReadCloser, string, error) {
calls++
return nil, "", errors.New("boom")
}
_, _, _ = w.gate("A", counting)
Expect(calls).To(Equal(1), "the cancellation must not have counted as the final failure")
_, _, err := w.gate("A", counting)
Expect(err).To(MatchError(errBreakerOpen), "the cancellation must not have cleared the earlier failures")
Expect(calls).To(Equal(1), "an open breaker must not call the external step")
})
It("does not open the breaker on a run of agent not-found misses", func() {
// agents.ErrNotFound is a definitive miss, not a fault: artless items must not
// trip the breaker, or they would loop in retry instead of settling absent.
notFound := func() (io.ReadCloser, string, error) { return nil, "", agents.ErrNotFound }
for range breakerThreshold + 3 {
_, _, err := w.gate("A", notFound)
Expect(err).To(MatchError(agents.ErrNotFound), "a miss passes through, never errBreakerOpen")
}
var calls int
counting := func() (io.ReadCloser, string, error) {
calls++
return nil, "", errors.New("boom")
}
_, _, _ = w.gate("A", counting)
Expect(calls).To(Equal(1), "the breaker stayed closed, so the step still runs")
})
It("isolates each agent's breaker: one open gate does not block another", func() {
failing := func() (io.ReadCloser, string, error) { return nil, "", errors.New("boom") }
for range breakerThreshold {
_, _, _ = w.gate("A", failing)
}
_, _, err := w.gate("A", failing)
Expect(err).To(MatchError(errBreakerOpen), "agent A's breaker is open")
var bCalls int
bStep := func() (io.ReadCloser, string, error) {
bCalls++
return io.NopCloser(nil), "p", nil
}
for range breakerThreshold + 2 {
_, _, err := w.gate("B", bStep)
Expect(err).ToNot(HaveOccurred(), "agent B keeps being called while A is open")
}
Expect(bCalls).To(Equal(breakerThreshold + 2))
})
})
Describe("precache", func() {
BeforeEach(func() {
folderRepo.result = []model.Folder{{
Path: "tests/fixtures/artist/an-album",
ImageFiles: []string{"cover.jpg"},
}}
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{
{ID: "alpc", Name: "Album", FolderIDs: []string{"f1"}},
})
conf.Server.UICoverArtSize = 300
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{
ItemKind: "al", ItemID: "alpc", Priority: model.ArtworkPriorityScan,
})).To(Succeed())
})
It("warms the resize cache at the UI cover size after a found acquisition", func() {
conf.Server.EnableArtworkPrecache = true
n, err := w.drain(ctx, 1)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
// The list surfaces request square covers, so warming any other key is wasted work.
Expect(imgCache.getKeys()).To(ContainElement(ContainSubstring(".300.true.")))
})
It("skips warming when precache is disabled", func() {
conf.Server.EnableArtworkPrecache = false
n, err := w.drain(ctx, 1)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(1))
Expect(imgCache.getKeys()).To(BeEmpty())
})
// It warms from the bytes acquisition already held, so no state row or store file is needed.
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, so a hit can only come from the bytes handed in;
// the probe refuses to open, proving nothing is re-read.
probe := &resizedItem{
hash: ia.Hash, size: 300, square: true, 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("drain pools", func() {
// A kind in neither pool is never dequeued, with nothing to catch it at compile time.
It("covers every kind the worker can process, exactly once", func() {
var pooled []string
for _, p := range newDrainPools() {
pooled = append(pooled, p.kinds...)
}
for kind := range artworkKindToResource {
Expect(pooled).To(ContainElement(kind.Prefix()), "kind %q belongs to no drain pool", kind.Prefix())
}
Expect(pooled).To(HaveLen(len(artworkKindToResource)), "a kind is claimed by more than one pool")
})
// Artists resolve through a rate-limited agent that holds its slot while waiting, so one
// shared pool would park every album behind them.
It("resolves albums while artists are stuck on a slow agent", func() {
conf.Server.CoverArtPriority = "cover.jpg"
conf.Server.ArtistArtPriority = "external"
folderRepo.result = []model.Folder{{
Path: "tests/fixtures/artist/an-album",
ImageFiles: []string{"cover.jpg"},
}}
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{{ID: "alx", Name: "Album", FolderIDs: []string{"f1"}}})
ds.MockedArtist = tests.CreateMockArtistRepo()
// Every artist lookup blocks until released, standing in for the rate limiter.
block := make(chan struct{})
artists := model.Artists{}
for i := range 8 {
artists = append(artists, model.Artist{ID: fmt.Sprintf("arx%d", i), Name: "A"})
}
ds.MockedArtist.(*tests.MockArtistRepo).SetData(artists)
imageAgents(&fakeImageAgent{name: "slowAgent", block: block})
// Artists first, exactly as Backfill orders them.
for _, a := range artists {
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{
ItemKind: "ar", ItemID: a.ID, Priority: model.ArtworkPriorityBackfill,
})).To(Succeed())
}
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{
ItemKind: "al", ItemID: "alx", Priority: model.ArtworkPriorityBackfill,
})).To(Succeed())
runCtx, cancel := context.WithCancel(ctx)
done := make(chan struct{})
go func() { defer close(done); _ = w.Run(runCtx) }()
// Join Run: a leaked pool goroutine would race the config snapshot Ginkgo restores.
DeferCleanup(func() {
cancel()
close(block) // unpark the blocked lookups so the pools can unwind
<-done
})
Eventually(func() bool {
ia, err := artRepo.GetItemArtwork(model.KindAlbumArtwork, "alx", model.ImageTypePrimary)
return err == nil && ia.Hash != ""
}, 5*time.Second, 50*time.Millisecond).Should(BeTrue(),
"a blocked external pool must not hold up local artwork")
_, err := artRepo.GetItemArtwork(model.KindArtistArtwork, "arx0", model.ImageTypePrimary)
Expect(err).To(MatchError(model.ErrNotFound), "artists are still blocked, as intended")
})
})
Describe("batching", func() {
It("leaves undispatched items queued when cancelled mid-batch", func() {
// Every album must resolve, so any dispatched item deletes its row regardless of dequeue order.
albums := model.Albums{}
for i := range 8 {
id := fmt.Sprintf("alc%d", i)
albums = append(albums, model.Album{ID: id, Name: "Album"})
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{
ItemKind: "al", ItemID: id, Priority: model.ArtworkPriorityScan,
})).To(Succeed())
}
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(albums)
cancelledCtx, cancel := context.WithCancel(ctx)
cancel()
_, err := w.drain(cancelledCtx, 1)
Expect(err).ToNot(HaveOccurred())
for i := range 8 {
id := fmt.Sprintf("alc%d", i)
Expect(findQueued(queueRepo, "al", id)).ToNot(BeNil(), "row "+id+" must survive a cancelled drain")
}
})
It("dequeues past the worker pool so one drain covers many items", func() {
for i := range 16 {
ds.MockedAlbum.(*tests.MockAlbumRepo).SetData(model.Albums{{ID: fmt.Sprintf("alb%d", i), Name: "Album"}})
Expect(queueRepo.Enqueue(model.ArtworkQueueItem{
ItemKind: "al", ItemID: fmt.Sprintf("alb%d", i), Priority: model.ArtworkPriorityScan,
})).To(Succeed())
}
n, err := w.drain(ctx, 2)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(16), "a batch sized to the pool would have stopped at 4")
})
})
Describe("RunPrune", func() {
It("runs a prune under the worker mutex", func() {
Expect(w.RunPrune(ctx)).To(Succeed())
})
})
Describe("Run", func() {
It("exits cleanly when the context is cancelled", func() {
runCtx, cancel := context.WithCancel(ctx)
done := make(chan error, 1)
go func() { done <- w.Run(runCtx) }()
cancel()
Eventually(done, time.Second).Should(Receive(BeNil()))
})
It("does not leak goroutines after Run exits", func() {
DeferCleanup(configtest.SetupConfig())
ignore := goleak.IgnoreCurrent()
DeferCleanup(func() { goleak.VerifyNone(GinkgoT(), ignore) })
localDS := &tests.MockDataStore{MockedArtworkQueue: tests.CreateMockArtworkQueueRepo()}
lw := NewWorker(localDS, NewImageStore(GinkgoT().TempDir()), agents.GetAgents(localDS, nil), tests.NewMockFFmpeg(""), &fakeEventBroker{}, imgCache)
runCtx, cancel := context.WithCancel(ctx)
done := make(chan error, 1)
go func() { done <- lw.Run(runCtx) }()
time.Sleep(20 * time.Millisecond) // let the loop settle on the idle select
cancel()
Eventually(done, 2*time.Second).Should(Receive(BeNil()))
})
})
})
var _ = Describe("backoff", func() {
It("returns the expected schedule with no jitter", func() {
for _, c := range []struct {
attempts int
want time.Duration
}{
{0, 5 * time.Second},
{1, 20 * time.Second},
{2, 80 * time.Second},
{3, 320 * time.Second},
{4, 1280 * time.Second},
{5, 5120 * time.Second},
{6, 20480 * time.Second},
{7, 12 * time.Hour},
{8, 12 * time.Hour},
} {
Expect(backoffFor(c.attempts, 0)).To(Equal(c.want), "attempt %d", c.attempts)
}
})
It("applies jitter proportionally", func() {
base := backoffFor(2, 0)
Expect(backoffFor(2, 0.2)).To(Equal(time.Duration(float64(base) * 1.2)))
Expect(backoffFor(2, -0.2)).To(Equal(time.Duration(float64(base) * 0.8)))
})
It("keeps random jitter within +/-40%", func() {
lo := time.Duration(float64(320*time.Second) * 0.6)
hi := time.Duration(float64(320*time.Second) * 1.4)
for range 200 {
d := backoff(3)
Expect(d).To(BeNumerically(">=", lo))
Expect(d).To(BeNumerically("<=", hi))
}
})
})