diff --git a/core/artwork/worker.go b/core/artwork/worker.go index 8c2f2bd40..62c6a7ea6 100644 --- a/core/artwork/worker.go +++ b/core/artwork/worker.go @@ -28,8 +28,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 + // itemTimeout must outlast a slow but healthy item: each image agent's call plus its download, in turn. + itemTimeout = 3 * time.Minute ) // drainPool drains one class of work with its own slot budget, so a blocking kind cannot @@ -253,6 +253,9 @@ func (w *Worker) process(ctx context.Context, item model.ArtworkQueueItem) (outc item.ImageType = cmp.Or(item.ImageType, model.ImageTypePrimary) ctx = withTrace(ctx, &ChainTrace{}) out, got, retryIn := w.acquireWithTimeout(ctx, item) + if ctx.Err() != nil { + return outcomeFailed, nil // shutting down: leave the row for the next run + } w.settle(ctx, item, out, retryIn) return out, got } @@ -283,6 +286,8 @@ func (w *Worker) acquireWithTimeout(ctx context.Context, item model.ArtworkQueue select { case r := <-done: return r.out, r.got, r.retryIn + case <-ctx.Done(): + return outcomeFailed, nil, 0 case <-time.After(itemTimeout): log.Error(ctx, "Artwork: Item timed out, moving on", "kind", item.ItemKind, "id", item.ItemID, "timeout", itemTimeout) traceStage(ctx, "worker", fmt.Errorf("timed out after %s", itemTimeout)) diff --git a/core/artwork/worker_test.go b/core/artwork/worker_test.go index a6c07b763..95b2df0a6 100644 --- a/core/artwork/worker_test.go +++ b/core/artwork/worker_test.go @@ -838,6 +838,11 @@ var _ = Describe("Worker", func() { cancel() close(block) // unpark the blocked lookups so the pools can unwind <-done + // Run returns without waiting on the lookups, so wait for them before config restores. + Eventually(func() (n int) { + w.busy.Range(func(_, _ any) bool { n++; return true }) + return n + }).Should(BeZero()) }) Eventually(func() bool { diff --git a/core/artwork/worker_timing_test.go b/core/artwork/worker_timing_test.go index d991fc38d..d4aa79d26 100644 --- a/core/artwork/worker_timing_test.go +++ b/core/artwork/worker_timing_test.go @@ -1,6 +1,7 @@ package artwork import ( + "context" "errors" "io" "testing" @@ -117,24 +118,32 @@ func TestArtworkGatePerAgentBreakerIsolation(t *testing.T) { }) } +// newStuckItemWorker queues one artist whose agent blocks until the returned channel is closed. +func newStuckItemWorker(t *testing.T) (*Worker, *tests.MockArtworkQueueRepo, *fakeImageAgent, chan struct{}) { + t.Cleanup(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) + if err := queue.Enqueue(t.Context(), model.ArtworkQueueItem{ItemKind: "ar", ItemID: "ar1"}); err != nil { + t.Fatal(err) + } + return w, queue, agent, block +} + // 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) + w, queue, agent, block := newStuckItemWorker(t) 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) @@ -167,3 +176,26 @@ func TestArtworkDrainMovesOnFromAStuckItem(t *testing.T) { g.Expect(agent.artistCalls).To(Equal(2), "one stuck call and one retry, never two at once") }) } + +// Shutdown must not wait on a stuck item, and must leave its row for the next run. +func TestArtworkDrainStopsAtOnceOnShutdown(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + g := NewWithT(t) + w, queue, _, block := newStuckItemWorker(t) + before := *findQueued(queue, "ar", "ar1") + ctx, cancel := context.WithCancel(t.Context()) + go func() { + time.Sleep(time.Second) + cancel() + }() + + start := time.Now() + _, err := w.drain(ctx, 1) + g.Expect(err).ToNot(HaveOccurred()) + g.Expect(time.Since(start)).To(Equal(time.Second), "the drain returns as soon as the context is cancelled") + g.Expect(*findQueued(queue, "ar", "ar1")).To(Equal(before), "the row is left as it was") + + close(block) + synctest.Wait() + }) +}