From ab5aa03b9f901af1df30b108efb07a64da8d8ad6 Mon Sep 17 00:00:00 2001 From: Deluan Date: Sun, 26 Jul 2026 21:48:24 -0400 Subject: [PATCH] refactor(artwork): collect the external gate contract in one file gateFunc, passthroughGate and isTransientExternal sat in agent_images.go while every implementation lived in worker.go: extGate, breaker, Worker.gate, gateFor. isTransientExternal even carries a comment saying it must stay consistent with breaker.record, which was in the other file -- a rule spanning two files with only a comment holding it together. Pure move into gate.go: no symbol added or removed. --- core/artwork/agent_images.go | 15 ----- core/artwork/gate.go | 112 +++++++++++++++++++++++++++++++++++ core/artwork/worker.go | 86 +-------------------------- 3 files changed, 113 insertions(+), 100 deletions(-) create mode 100644 core/artwork/gate.go diff --git a/core/artwork/agent_images.go b/core/artwork/agent_images.go index dd8c50727..681490936 100644 --- a/core/artwork/agent_images.go +++ b/core/artwork/agent_images.go @@ -2,7 +2,6 @@ package artwork import ( "context" - "errors" "io" "net/url" @@ -24,14 +23,6 @@ func externalName(name string) string { return str.Clear(name) } -// gateFunc gates one named external fetch (rate limit + circuit breaker per name). -// resolveItem defaults to passthroughGate; the worker injects the per-agent gate. -type gateFunc = func(name string, f func() (io.ReadCloser, string, error)) (io.ReadCloser, string, error) - -func passthroughGate(_ string, f func() (io.ReadCloser, string, error)) (io.ReadCloser, string, error) { - return f() -} - // bestImageURL returns the largest-Size image URL, skipping empty or unparseable // URLs; nil when none qualifies. Parsing happens per candidate so a malformed largest // URL never shadows a valid smaller one. @@ -112,9 +103,3 @@ func fetchAlbumImage(ctx context.Context, ag *agents.Agents, gate gateFunc, al m } return nil, "", extErr } - -// isTransientExternal reports whether an external step failed in a way worth retrying; -// a not-found (from either package) is a definitive answer, not a fault. -func isTransientExternal(err error) bool { - return err != nil && !errors.Is(err, agents.ErrNotFound) && !errors.Is(err, model.ErrNotFound) -} diff --git a/core/artwork/gate.go b/core/artwork/gate.go new file mode 100644 index 000000000..0d7cf15df --- /dev/null +++ b/core/artwork/gate.go @@ -0,0 +1,112 @@ +package artwork + +import ( + "errors" + "io" + "sync" + "time" + + "github.com/navidrome/navidrome/conf" + "github.com/navidrome/navidrome/core/agents" + "github.com/navidrome/navidrome/model" + "golang.org/x/time/rate" +) + +const ( + breakerThreshold = 5 + breakerProbeAfter = time.Minute +) + +var errBreakerOpen = errors.New("artwork: external circuit breaker open") + +// gateFunc gates one named external fetch (rate limit + circuit breaker per name). +// resolveItem defaults to passthroughGate; the worker injects the per-agent gate. +type gateFunc = func(name string, f func() (io.ReadCloser, string, error)) (io.ReadCloser, string, error) + +func passthroughGate(_ string, f func() (io.ReadCloser, string, error)) (io.ReadCloser, string, error) { + return f() +} + +// isTransientExternal reports whether an external step failed in a way worth retrying; +// a not-found (from either package) is a definitive answer, not a fault. +func isTransientExternal(err error) bool { + return err != nil && !errors.Is(err, agents.ErrNotFound) && !errors.Is(err, model.ErrNotFound) +} + +// extGate is one agent's rate limiter + circuit breaker; each external agent gets its +// own so a provider whose API or CDN is down backs off in isolation from the others. +type extGate struct { + limiter *rate.Limiter + breaker *breaker +} + +// gate wraps a named external step with that agent's own rate limiter and circuit +// breaker, matching gateFunc so it can be handed to the processor's resolver. +func (w *Worker) gate(name string, f func() (io.ReadCloser, string, error)) (io.ReadCloser, string, error) { + g := w.gateFor(name) + if !g.breaker.allow() { + return nil, "", errBreakerOpen + } + if err := g.limiter.Wait(w.runCtx); err != nil { + return nil, "", err + } + r, path, err := f() + g.breaker.record(err) + return r, path, err +} + +// gateFor lazily creates the per-name gate on first use, each with its own limiter at +// ArtworkExternalMaxRPS and its own breaker. +func (w *Worker) gateFor(name string) *extGate { + w.gatesMu.Lock() + defer w.gatesMu.Unlock() + if g, ok := w.gates[name]; ok { + return g + } + rps := conf.Server.ArtworkExternalMaxRPS + limit := rate.Inf + if rps > 0 { + limit = rate.Limit(rps) + } + g := &extGate{limiter: rate.NewLimiter(limit, max(1, rps)), breaker: newBreaker()} + w.gates[name] = g + return g +} + +// breaker opens after breakerThreshold consecutive external errors and admits a +// single probe once breakerProbeAfter has elapsed; a success re-closes it. +type breaker struct { + mu sync.Mutex + failures int + openedAt time.Time +} + +func newBreaker() *breaker { return &breaker{} } + +func (b *breaker) allow() bool { + b.mu.Lock() + defer b.mu.Unlock() + if b.failures < breakerThreshold { + return true + } + if time.Since(b.openedAt) >= breakerProbeAfter { + b.openedAt = time.Now() // start a fresh probe window so only one caller passes + return true + } + return false +} + +func (b *breaker) record(err error) { + b.mu.Lock() + defer b.mu.Unlock() + // A not-found (from either package) is a definitive answer, not a fault; only real + // errors trip the breaker. Must stay consistent with isTransientExternal. + if err == nil || errors.Is(err, model.ErrNotFound) || errors.Is(err, agents.ErrNotFound) { + b.failures = 0 + return + } + b.failures++ + if b.failures == breakerThreshold { + b.openedAt = time.Now() + } +} diff --git a/core/artwork/worker.go b/core/artwork/worker.go index d113a6fe0..aa9b9bc35 100644 --- a/core/artwork/worker.go +++ b/core/artwork/worker.go @@ -3,7 +3,6 @@ package artwork import ( "bytes" "context" - "errors" "io" "math" "math/rand/v2" @@ -18,7 +17,6 @@ import ( "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/server/events" "github.com/navidrome/navidrome/utils/cache" - "golang.org/x/time/rate" ) const ( @@ -26,20 +24,9 @@ const ( backoffBase = 5 * time.Second // giveUpAfter bounds the retry budget from enqueue: past it the worker stops retrying and // hands the item to the periodic stale-absent recheck (settling absent on a bare failure). - giveUpAfter = 12 * time.Hour - breakerThreshold = 5 - breakerProbeAfter = time.Minute + giveUpAfter = 12 * time.Hour ) -var errBreakerOpen = errors.New("artwork: external circuit breaker open") - -// extGate is one agent's rate limiter + circuit breaker; each external agent gets its -// own so a provider whose API or CDN is down backs off in isolation from the others. -type extGate struct { - limiter *rate.Limiter - breaker *breaker -} - // drainPool drains one class of work with its own slot budget, so a kind whose resolution // blocks cannot occupy slots another kind needs. type drainPool struct { @@ -326,39 +313,6 @@ func (w *Worker) precache(ctx context.Context, got *acquired) { _ = stream.Close() } -// gate wraps a named external step with that agent's own rate limiter and circuit -// breaker, matching gateFunc so it can be handed to the processor's resolver. -func (w *Worker) gate(name string, f func() (io.ReadCloser, string, error)) (io.ReadCloser, string, error) { - g := w.gateFor(name) - if !g.breaker.allow() { - return nil, "", errBreakerOpen - } - if err := g.limiter.Wait(w.runCtx); err != nil { - return nil, "", err - } - r, path, err := f() - g.breaker.record(err) - return r, path, err -} - -// gateFor lazily creates the per-name gate on first use, each with its own limiter at -// ArtworkExternalMaxRPS and its own breaker. -func (w *Worker) gateFor(name string) *extGate { - w.gatesMu.Lock() - defer w.gatesMu.Unlock() - if g, ok := w.gates[name]; ok { - return g - } - rps := conf.Server.ArtworkExternalMaxRPS - limit := rate.Inf - if rps > 0 { - limit = rate.Limit(rps) - } - g := &extGate{limiter: rate.NewLimiter(limit, max(1, rps)), breaker: newBreaker()} - w.gates[name] = g - return g -} - // backoffFor returns min(5s×4^n, giveUpAfter) scaled by (1+jitter), with jitter in [-0.4, 0.4]. func backoffFor(attempts int, jitter float64) time.Duration { d := math.Min(float64(backoffBase)*math.Pow(4, float64(attempts)), float64(giveUpAfter)) @@ -368,41 +322,3 @@ func backoffFor(attempts int, jitter float64) time.Duration { func backoff(attempts int) time.Duration { return backoffFor(attempts, rand.Float64()*0.8-0.4) //nolint:gosec // retry jitter, not security-sensitive } - -// breaker opens after breakerThreshold consecutive external errors and admits a -// single probe once breakerProbeAfter has elapsed; a success re-closes it. -type breaker struct { - mu sync.Mutex - failures int - openedAt time.Time -} - -func newBreaker() *breaker { return &breaker{} } - -func (b *breaker) allow() bool { - b.mu.Lock() - defer b.mu.Unlock() - if b.failures < breakerThreshold { - return true - } - if time.Since(b.openedAt) >= breakerProbeAfter { - b.openedAt = time.Now() // start a fresh probe window so only one caller passes - return true - } - return false -} - -func (b *breaker) record(err error) { - b.mu.Lock() - defer b.mu.Unlock() - // A not-found (from either package) is a definitive answer, not a fault; only real - // errors trip the breaker. Must stay consistent with isTransientExternal. - if err == nil || errors.Is(err, model.ErrNotFound) || errors.Is(err, agents.ErrNotFound) { - b.failures = 0 - return - } - b.failures++ - if b.failures == breakerThreshold { - b.openedAt = time.Now() - } -}