diff --git a/core/artwork/folders_artist.go b/core/artwork/folders_artist.go index 3403935db..40972f71d 100644 --- a/core/artwork/folders_artist.go +++ b/core/artwork/folders_artist.go @@ -30,7 +30,7 @@ func fromArtistFolder(ctx context.Context, libFS fs.FS, libPath, artistFolder, p if libFS == nil { return nil, "", fmt.Errorf("artist folder lookup unavailable") } - rel, err := filepath.Rel(libPath, artistFolder) + rel, err := utils.RelPath(libPath, artistFolder) if err != nil || rel == ".." || strings.HasPrefix(rel, ".."+string(filepath.Separator)) { return nil, "", fmt.Errorf(`artist folder '%s' is outside library '%s'`, artistFolder, libPath) } diff --git a/core/artwork/worker.go b/core/artwork/worker.go index 28e51958c..6122f88f2 100644 --- a/core/artwork/worker.go +++ b/core/artwork/worker.go @@ -4,6 +4,8 @@ import ( "bytes" "cmp" "context" + "errors" + "fmt" "io" "math" "math/rand/v2" @@ -26,6 +28,10 @@ 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 must outlast a slow but healthy item: each image agent's call plus its download, in turn. + itemTimeout = 3 * time.Minute + // shutdownGrace lets cancelled work unwind before Run returns and the DB closes, without waiting on a stuck item. + shutdownGrace = 5 * time.Second ) // drainPool drains one class of work with its own slot budget, so a blocking kind cannot @@ -50,8 +56,13 @@ type Worker struct { gatesMu sync.Mutex gates map[string]*extGate + + // busy holds the itemKeys still being acquired, including ones a drain stopped waiting for. + busy sync.Map } +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}, @@ -242,10 +253,56 @@ func (w *Worker) broadcastRefresh(ctx context.Context, found []model.ArtworkQueu func (w *Worker) process(ctx context.Context, item model.ArtworkQueueItem) (outcome, *acquired) { item.ImageType = cmp.Or(item.ImageType, model.ImageTypePrimary) - trace := &ChainTrace{} - ctx = withTrace(ctx, trace) - out, got, retryIn := w.proc.acquire(ctx, item) + 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 +} +// 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 _, running := w.busy.LoadOrStore(key, struct{}{}); running { + log.Debug(ctx, "Artwork: Item still running from an earlier drain", "kind", item.ItemKind, "id", item.ItemID) + traceStage(ctx, "worker", errors.New("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() { + // The deadline stops work that honors ctx; the select below covers work that doesn't. + actx, cancel := context.WithTimeout(ctx, itemTimeout) + defer cancel() + out, got, retryIn := w.proc.acquire(actx, item) + // Cleared before the send, so the next drain never sees a finished item as busy. + w.busy.Delete(key) + done <- result{out, got, retryIn} + }() + select { + case r := <-done: + return r.out, r.got, r.retryIn + case <-ctx.Done(): + select { + case <-done: + case <-time.After(shutdownGrace): + } + 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)) + return outcomeFailed, nil, 0 + } +} + +// 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) { queue := w.proc.ds.ArtworkQueue() switch out { case outcomeFound, outcomeAbsent: @@ -256,7 +313,7 @@ func (w *Worker) process(ctx context.Context, item model.ArtworkQueueItem) (outc } case outcomeFoundStale, outcomeFailed: retryAt := time.Now().Add(retryDelay(item.Attempts, retryIn)) - encoded := trace.encode("") + encoded := traceFrom(ctx).encode("") if retryAt.Before(item.EnqueuedAt.Add(giveUpAfter)) { // A mid-flight re-enqueue reset retry_at; stale backoff must not stomp its // fresh, immediate eligibility. @@ -283,7 +340,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 diff --git a/core/artwork/worker_timing_test.go b/core/artwork/worker_timing_test.go index 63f12d03c..7896f6f09 100644 --- a/core/artwork/worker_timing_test.go +++ b/core/artwork/worker_timing_test.go @@ -1,13 +1,17 @@ package artwork import ( + "context" "errors" "io" "testing" "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 +117,109 @@ func TestArtworkGatePerAgentBreakerIsolation(t *testing.T) { g.Expect(aCalls).To(Equal(1)) }) } + +// 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) + w, queue, agent, block := newStuckItemWorker(t) + ctx := t.Context() + + 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") + }) +} + +// Shutdown joins cancelled work that unwinds, so nothing touches the DB after Run returns. +func TestArtworkDrainShutdownWaitsForWorkToUnwind(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() + time.Sleep(100 * time.Millisecond) + close(block) + }() + + start := time.Now() + _, err := w.drain(ctx, 1) + g.Expect(err).ToNot(HaveOccurred()) + g.Expect(time.Since(start)).To(Equal(time.Second+100*time.Millisecond), "the drain waits for the cancelled item to return") + _, running := w.busy.Load(itemKey{"ar", "ar1", model.ImageTypePrimary}) + g.Expect(running).To(BeFalse()) + g.Expect(*findQueued(queue, "ar", "ar1")).To(Equal(before), "the row is left for the next run") + }) +} + +// Shutdown must not wait on a stuck item for longer than the grace period. +func TestArtworkDrainShutdownDoesNotWaitOnAStuckItem(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+shutdownGrace), "the drain gives up on the stuck item after the grace period") + g.Expect(*findQueued(queue, "ar", "ar1")).To(Equal(before), "the row is left for the next run") + + close(block) + synctest.Wait() + }) +} diff --git a/core/playlists/parse_m3u.go b/core/playlists/parse_m3u.go index 9610e9dbb..9ca1c9b6c 100644 --- a/core/playlists/parse_m3u.go +++ b/core/playlists/parse_m3u.go @@ -15,6 +15,7 @@ import ( "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/model/request" + "github.com/navidrome/navidrome/utils" "github.com/navidrome/navidrome/utils/slice" "golang.org/x/text/unicode/norm" ) @@ -148,7 +149,7 @@ func (r pathResolution) ToQualifiedString() (string, error) { if !r.valid { return "", fmt.Errorf("invalid path resolution") } - relativePath, err := filepath.Rel(r.libraryPath, r.absolutePath) + relativePath, err := utils.RelPath(r.libraryPath, r.absolutePath) if err != nil { return "", err } diff --git a/core/playlists/playlists.go b/core/playlists/playlists.go index c9bc03b97..8e7ccf93c 100644 --- a/core/playlists/playlists.go +++ b/core/playlists/playlists.go @@ -15,6 +15,7 @@ import ( "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/model/request" + "github.com/navidrome/navidrome/utils" ) type Playlists interface { @@ -77,7 +78,7 @@ func InPath(folder model.Folder) bool { if conf.Server.PlaylistsPath == "" { return true } - rel, _ := filepath.Rel(folder.LibraryPath, folder.AbsolutePath()) + rel, _ := utils.RelPath(folder.LibraryPath, folder.AbsolutePath()) for path := range strings.SplitSeq(conf.Server.PlaylistsPath, string(filepath.ListSeparator)) { if match, _ := doublestar.Match(path, rel); match { return true diff --git a/scanner/scanner.go b/scanner/scanner.go index cd2fe3c8d..f156d15fd 100644 --- a/scanner/scanner.go +++ b/scanner/scanner.go @@ -15,6 +15,7 @@ import ( "github.com/navidrome/navidrome/core/playlists" "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" + "github.com/navidrome/navidrome/utils" "github.com/navidrome/navidrome/utils/run" "github.com/navidrome/navidrome/utils/slice" ) @@ -65,7 +66,7 @@ func libraryRelativePath(libPath, folderPath string) string { if err != nil { return folderPath } - rel, err := filepath.Rel(absLib, folderPath) + rel, err := utils.RelPath(absLib, folderPath) if err != nil || !filepath.IsLocal(rel) { return folderPath } diff --git a/scanner/watcher.go b/scanner/watcher.go index 1ac5468f0..3c0fbdbe7 100644 --- a/scanner/watcher.go +++ b/scanner/watcher.go @@ -13,6 +13,7 @@ import ( "github.com/navidrome/navidrome/core/storage" "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" + "github.com/navidrome/navidrome/utils" "github.com/navidrome/navidrome/utils/singleton" ) @@ -245,7 +246,7 @@ func (w *watcher) processLibraryEvents(ctx context.Context, lib *model.Library, log.Debug(ctx, "Watcher stopped due to context cancellation", "libraryID", lib.ID, "name", lib.Name) return nil case path := <-events: - path, err := filepath.Rel(absLibPath, path) + path, err := utils.RelPath(absLibPath, path) if err != nil { log.Error(ctx, "Error getting relative path", "libraryID", lib.ID, "absolutePath", absLibPath, "path", path, err) continue diff --git a/utils/cache/spread_fs.go b/utils/cache/spread_fs.go index 11801b324..3ebd00e1a 100644 --- a/utils/cache/spread_fs.go +++ b/utils/cache/spread_fs.go @@ -12,6 +12,7 @@ import ( "github.com/djherbis/fscache" "github.com/djherbis/stream" "github.com/navidrome/navidrome/log" + "github.com/navidrome/navidrome/utils" ) const completeMarkerSuffix = ".complete" @@ -97,7 +98,7 @@ func (sfs *spreadFS) walkDataFiles(visit func(absoluteFilePath string)) error { log.Error("Error loading cache", "dir", sfs.root, err) return nil } - path, err := filepath.Rel(sfs.root, absoluteFilePath) + path, err := utils.RelPath(sfs.root, absoluteFilePath) if err != nil { return nil //nolint:nilerr } diff --git a/utils/files.go b/utils/files.go index 2fce307ea..8ef522e5e 100644 --- a/utils/files.go +++ b/utils/files.go @@ -40,3 +40,16 @@ func FileExists(path string) bool { _, err := os.Stat(path) return err == nil || !os.IsNotExist(err) } + +// RelPath is filepath.Rel without the UNC share root hang (golang/go#79784). Remove it after moving to Go 1.28. +func RelPath(basePath, targPath string) (string, error) { + return filepath.Rel(normalizeUNCRoot(basePath), normalizeUNCRoot(targPath)) +} + +func normalizeUNCRoot(p string) string { + vol := filepath.VolumeName(p) + if len(vol) > 2 && strings.Trim(p[len(vol):], `\/`) == "" { + return vol + string(filepath.Separator) + } + return p +} diff --git a/utils/files_test.go b/utils/files_test.go index c6e578f05..04b4d4e0b 100644 --- a/utils/files_test.go +++ b/utils/files_test.go @@ -223,3 +223,26 @@ var _ = Describe("FileExists", func() { Expect(result).To(Or(BeTrue(), BeFalse())) // Should not panic }) }) + +var _ = Describe("RelPath", func() { + DescribeTable("returns the same result as filepath.Rel", + func(base, target string) { + base, target = filepath.FromSlash(base), filepath.FromSlash(target) + expected, expectedErr := filepath.Rel(base, target) + rel, err := utils.RelPath(base, target) + Expect(rel).To(Equal(expected)) + if expectedErr == nil { + Expect(err).ToNot(HaveOccurred()) + } else { + Expect(err).To(MatchError(expectedErr.Error())) + } + }, + Entry("same path", "/music", "/music"), + Entry("same path with trailing separator", "/music", "/music/"), + Entry("child path", "/music", "/music/artist/album"), + Entry("sibling path", "/music/a", "/music/b"), + Entry("parent path", "/music/artist", "/music"), + Entry("relative paths", "music", "music/artist"), + Entry("absolute and relative paths", "/music", "music"), + ) +}) diff --git a/utils/files_windows_test.go b/utils/files_windows_test.go new file mode 100644 index 000000000..419291b98 --- /dev/null +++ b/utils/files_windows_test.go @@ -0,0 +1,40 @@ +package utils_test + +import ( + "time" + + "github.com/navidrome/navidrome/utils" + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +var _ = Describe("RelPath on Windows", func() { + type result struct { + rel string + err error + } + + DescribeTable("handles UNC share roots without hanging", + func(base, target, expected string) { + done := make(chan result, 1) + go func() { + rel, err := utils.RelPath(base, target) + done <- result{rel, err} + }() + var res result + Eventually(done).WithTimeout(5 * time.Second).Should(Receive(&res)) + Expect(res.err).ToNot(HaveOccurred()) + Expect(res.rel).To(Equal(expected)) + }, + Entry("root and root with trailing separator", `\\server\Music`, `\\server\Music\`, "."), + Entry("root with trailing separator and root", `\\server\Music\`, `\\server\Music`, "."), + Entry("root and root with many separators", `\\server\Music`, `\\server\Music\\`, "."), + Entry("root with forward slashes", `//server/Music`, `//server/Music/`, "."), + Entry("root and child folder", `\\server\Music`, `\\server\Music\Artist`, "Artist"), + Entry("root with trailing separator and child folder", `\\server\Music\`, `\\server\Music\Artist\Album`, `Artist\Album`), + Entry("child folder and root", `\\server\Music\Artist`, `\\server\Music\`, ".."), + Entry("child folder and root without trailing separator", `\\server\Music\Artist`, `\\server\Music`, ".."), + Entry("root and root with different case", `\\server\Music`, `\\SERVER\music\`, "."), + Entry("drive root and child folder", `C:\`, `C:\Music`, "Music"), + ) +})