mirror of
https://github.com/navidrome/navidrome.git
synced 2026-10-08 18:37:09 +02:00
Compare commits
6 commits
master
...
t3code/fix
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d90ee1819e |
||
|
|
83d1aa4c3b | ||
|
|
3817637207 | ||
|
|
89e17f7ed8 | ||
|
|
d513b1ea48 | ||
|
|
e42418db77 |
11 changed files with 258 additions and 11 deletions
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
})
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
3
utils/cache/spread_fs.go
vendored
3
utils/cache/spread_fs.go
vendored
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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"),
|
||||
)
|
||||
})
|
||||
|
|
|
|||
40
utils/files_windows_test.go
Normal file
40
utils/files_windows_test.go
Normal file
|
|
@ -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"),
|
||||
)
|
||||
})
|
||||
Loading…
Add table
Add a link
Reference in a new issue