fix(artwork): stop a stuck item from stalling its drain pool

Worker.drain waits for every item in a batch before it dequeues the next
one, so a single item that never returns stops the whole pool. This is what
happened in discussion #6232: filepath.Rel looped forever on a UNC share
root, and all artist artwork stopped. Go can't stop a goroutine that
spins, so the only way to keep the pool going is to stop waiting.

process now runs acquire in its own goroutine and waits at most
itemTimeout (5 minutes). On timeout it logs an error, adds a trace step,
and settles the row as failed, so the normal backoff and give-up budget
apply. A per-worker busy set keeps one item from running twice at the
same time: if the row comes due again while the old run is still stuck,
it is settled as failed without starting a second run. This keeps the
existing guarantee that the drain resolves each item serially.

Shutdown still waits for in-flight items, but now at most itemTimeout
for one that is stuck.
This commit is contained in:
Deluan 2026-09-27 15:39:16 -04:00
commit d513b1ea48
2 changed files with 117 additions and 2 deletions

View file

@ -26,6 +26,8 @@ const (
// giveUpAfter bounds the retry budget from enqueue; past it the item settles and only an
// explicit reprocess retries it.
giveUpAfter = 12 * time.Hour
// itemTimeout is how long a drain waits on one item before rescheduling it.
itemTimeout = 5 * time.Minute
)
// drainPool drains one class of work with its own slot budget, so a blocking kind cannot
@ -50,8 +52,14 @@ type Worker struct {
gatesMu sync.Mutex
gates map[string]*extGate
// busy holds the items still being acquired, including ones a drain stopped waiting for.
busyMu sync.Mutex
busy map[itemKey]struct{}
}
type itemKey struct{ kind, id, imageType string }
func NewWorker(ds model.DataStore, store *ImageStore, ag *agents.Agents, ffmpeg ffmpeg.FFmpeg, broker events.Broker, imgCache cache.FileCache) *Worker {
w := &Worker{
proc: &processor{ds: ds, store: store},
@ -62,6 +70,7 @@ func NewWorker(ds model.DataStore, store *ImageStore, ag *agents.Agents, ffmpeg
runCtx: context.Background(),
paused: func() bool { return false },
gates: map[string]*extGate{},
busy: map[itemKey]struct{}{},
}
w.proc.resolver = newResolver(ds, ag, ffmpeg, w.gate)
w.proc.pruneLock = w.pruneMu.RLocker()
@ -244,8 +253,61 @@ func (w *Worker) process(ctx context.Context, item model.ArtworkQueueItem) (outc
item.ImageType = cmp.Or(item.ImageType, model.ImageTypePrimary)
trace := &ChainTrace{}
ctx = withTrace(ctx, trace)
out, got, retryIn := w.proc.acquire(ctx, item)
out, got, retryIn := w.acquireWithTimeout(ctx, item)
w.settle(ctx, item, out, retryIn, trace)
return out, got
}
// acquireWithTimeout stops waiting on a stuck item, which Go can't stop, and never runs one item twice at once.
func (w *Worker) acquireWithTimeout(ctx context.Context, item model.ArtworkQueueItem) (outcome, *acquired, time.Duration) {
key := itemKey{item.ItemKind, item.ItemID, item.ImageType}
if !w.markBusy(key) {
log.Debug(ctx, "Artwork: Item still running from an earlier drain", "kind", item.ItemKind, "id", item.ItemID)
traceFrom(ctx).add(TraceStep{Candidate: "worker", Outcome: OutcomeError, Detail: "previous attempt still running"})
return outcomeFailed, nil, 0
}
type result struct {
out outcome
got *acquired
retryIn time.Duration
}
done := make(chan result, 1)
go func() {
out, got, retryIn := w.proc.acquire(ctx, item)
// Cleared before the send, so the next drain never sees a finished item as busy.
w.clearBusy(key)
done <- result{out, got, retryIn}
}()
timer := time.NewTimer(itemTimeout)
defer timer.Stop()
select {
case r := <-done:
return r.out, r.got, r.retryIn
case <-timer.C:
log.Error(ctx, "Artwork: Item timed out, moving on", "kind", item.ItemKind, "id", item.ItemID, "timeout", itemTimeout)
traceFrom(ctx).add(TraceStep{Candidate: "worker", Outcome: OutcomeError, Detail: "timed out after " + itemTimeout.String()})
return outcomeFailed, nil, 0
}
}
func (w *Worker) markBusy(key itemKey) bool {
w.busyMu.Lock()
defer w.busyMu.Unlock()
if _, ok := w.busy[key]; ok {
return false
}
w.busy[key] = struct{}{}
return true
}
func (w *Worker) clearBusy(key itemKey) {
w.busyMu.Lock()
defer w.busyMu.Unlock()
delete(w.busy, key)
}
// settle deletes or reschedules the queue row for an acquire outcome.
func (w *Worker) settle(ctx context.Context, item model.ArtworkQueueItem, out outcome, retryIn time.Duration, trace *ChainTrace) {
queue := w.proc.ds.ArtworkQueue()
switch out {
case outcomeFound, outcomeAbsent:
@ -283,7 +345,6 @@ func (w *Worker) process(ctx context.Context, item model.ArtworkQueueItem) (outc
log.Warn(ctx, "Artwork: Could not remove exhausted queue item", "kind", item.ItemKind, "id", item.ItemID, err)
}
}
return out, got
}
// recordGiveUp keeps the last failure on the state row after the queue row is deleted. An item

View file

@ -7,7 +7,10 @@ import (
"testing/synctest"
"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/tests"
. "github.com/onsi/gomega"
)
@ -113,3 +116,54 @@ func TestArtworkGatePerAgentBreakerIsolation(t *testing.T) {
g.Expect(aCalls).To(Equal(1))
})
}
// Go can't stop a stuck item, so the drain must stop waiting for it, or the whole pool stalls.
func TestArtworkDrainMovesOnFromAStuckItem(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
g := NewWithT(t)
defer configtest.SetupConfig()()
conf.Server.ArtistArtPriority = "external"
conf.Server.DevArtworkExternalMaxRPS = 1000
block := make(chan struct{})
agent := &fakeImageAgent{name: "stuckAgent", block: block}
ag := imageAgents(agent)
artists := tests.CreateMockArtistRepo()
artists.SetData(model.Artists{{ID: "ar1", Name: "Artist"}})
queue := tests.CreateMockArtworkQueueRepo()
ds := &tests.MockDataStore{MockedArtist: artists, MockedFolder: &fakeFolderRepo{}, MockedArtwork: tests.CreateMockArtworkRepo(), MockedArtworkQueue: queue}
w := NewWorker(ds, NewImageStore(t.TempDir()), ag, tests.NewMockFFmpeg(""), &fakeEventBroker{}, nil)
ctx := t.Context()
g.Expect(queue.Enqueue(ctx, model.ArtworkQueueItem{ItemKind: "ar", ItemID: "ar1"})).To(Succeed())
start := time.Now()
n, err := w.drain(ctx, 1)
g.Expect(err).ToNot(HaveOccurred())
g.Expect(n).To(Equal(1))
g.Expect(time.Since(start)).To(Equal(itemTimeout), "the drain stops waiting once the item times out")
row := findQueued(queue, "ar", "ar1")
g.Expect(row).ToNot(BeNil(), "a timed-out item is rescheduled, not dropped")
g.Expect(row.Attempts).To(Equal(1))
g.Expect(row.RetryAt).To(BeTemporally(">", time.Now()))
g.Expect(row.Trace).To(ContainSubstring("timed out"))
// Due again while the first run is still stuck: it must not start a second one.
time.Sleep(time.Until(row.RetryAt))
start = time.Now()
_, err = w.drain(ctx, 1)
g.Expect(err).ToNot(HaveOccurred())
g.Expect(time.Since(start)).To(BeZero(), "a still-running item is rescheduled without waiting")
row = findQueued(queue, "ar", "ar1")
g.Expect(row.Attempts).To(Equal(2))
g.Expect(row.Trace).To(ContainSubstring("still running"))
// Once the stuck run returns, the next attempt runs normally and settles the row.
close(block)
synctest.Wait()
time.Sleep(time.Until(row.RetryAt))
_, err = w.drain(ctx, 1)
g.Expect(err).ToNot(HaveOccurred())
g.Expect(findQueued(queue, "ar", "ar1")).To(BeNil())
g.Expect(agent.artistCalls).To(Equal(2), "one stuck call and one retry, never two at once")
})
}