diff --git a/backend/download/candidate_shape_test.go b/backend/download/candidate_shape_test.go new file mode 100644 index 0000000..6b979a0 --- /dev/null +++ b/backend/download/candidate_shape_test.go @@ -0,0 +1,248 @@ +package download + +import ( + "context" + "os" + "path/filepath" + "testing" +) + +// Multi-disc rips, single-track results and coverage counted in tracks +// rather than files (#270). + +func TestParsePathReadsTheDiscFromItsFolder(t *testing.T) { + t.Parallel() + + cases := []struct { + path string + disc int + track int + folder string + }{ + {`\share\Pink Floyd - The Wall (1979)\CD2\03 Hey You.flac`, 2, 3, "The Wall"}, + {`\share\The Wall\Disc 1\01 In The Flesh.flac`, 1, 1, "The Wall"}, + {`\share\The Wall\[Disk-2]\01 Hey You.flac`, 2, 1, "The Wall"}, + {`\share\The Wall\CD1 - Live\04 Mother.flac`, 1, 4, "The Wall"}, + // The filename's own disc number is more specific than the folder. + {`\share\The Wall\CD1\2-05 Comfortably Numb.flac`, 2, 5, "The Wall"}, + // Not a disc folder: a number is required. + {`\share\CDs\The Wall\01 In The Flesh.flac`, 0, 1, "The Wall"}, + } + + for _, tc := range cases { + t.Run(tc.path, func(t *testing.T) { + t.Parallel() + + got := ParsePath(tc.path) + if got.Disc != tc.disc || got.Track != tc.track || got.Folder != tc.folder { + t.Errorf( + "ParsePath = disc %d track %d folder %q, want %d %d %q", + got.Disc, got.Track, got.Folder, tc.disc, tc.track, tc.folder, + ) + } + }) + } +} + +func TestAlbumDir(t *testing.T) { + t.Parallel() + + cases := map[string]string{ + `\share\Album\CD1\01 A.flac`: "/share/Album", + `\share\Album\01 A.flac`: "/share/Album", + `CD1\01 A.flac`: "CD1", + `\share\CD Collection\01.mp3`: "/share/CD Collection", + } + + for in, want := range cases { + if got := AlbumDir(in); got != want { + t.Errorf("AlbumDir(%q) = %q, want %q", in, got, want) + } + } +} + +// One album shared as CD1/CD2 is one candidate, named after the album. +func TestSlskdGroupsDiscFoldersIntoOneCandidate(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + stub.responses = []slskdResponse{{ + Username: "peer", + Files: []slskdFile{ + {Filename: `\m\The Wall\CD1\01 In The Flesh.flac`, Size: 1}, + {Filename: `\m\The Wall\CD1\02 The Thin Ice.flac`, Size: 1}, + {Filename: `\m\The Wall\CD2\01 Hey You.flac`, Size: 1}, + {Filename: `\m\The Wall\CD2\02 Is There Anybody Out There.flac`, Size: 1}, + }, + }} + + s, _ := newStubSlskd(t, stub) + + got, err := s.Search(context.Background(), Download{Query: "the wall"}) + if err != nil { + t.Fatalf("Search: %v", err) + } + + if len(got) != 1 { + t.Fatalf("got %d candidates, want the two discs as one", len(got)) + } + + if got[0].Title != "The Wall" || len(got[0].Files) != 4 { + t.Errorf( + "candidate = %q with %d files, want \"The Wall\" with 4", + got[0].Title, len(got[0].Files), + ) + } +} + +// A track search matches one file per folder, so a single-track request +// must accept a one-file folder that an album request rightly drops. +func TestSlskdKeepsASingleFileForATrackRequest(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + stub.responses = []slskdResponse{{ + Username: "peer", + Files: []slskdFile{ + {Filename: `\m\OK Computer\02 Paranoid Android.flac`, Size: 1}, + }, + }} + + s, _ := newStubSlskd(t, stub) + + track, err := s.Search(context.Background(), Download{ + RecordingMBID: "rec-1", Artist: "Radiohead", Album: "Paranoid Android", + }) + if err != nil { + t.Fatalf("Search: %v", err) + } + + if len(track) != 1 { + t.Errorf("track request: got %d candidates, want 1", len(track)) + } + + album, err := s.Search(context.Background(), Download{ + ReleaseMBID: "rel-1", Artist: "Radiohead", Album: "OK Computer", + }) + if err != nil { + t.Fatalf("Search: %v", err) + } + + if len(album) != 0 { + t.Errorf("album request: got %d candidates, want the one-file folder dropped", len(album)) + } +} + +// Two discs with a file of the same name both reach staging, each under +// its disc folder, where the importer reads the disc number from. +func TestSlskdCollectKeepsDiscFolders(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + s, downloads := newStubSlskd(t, stub) + + for _, disc := range []string{"CD1", "CD2"} { + dir := filepath.Join(downloads, disc) + if err := os.MkdirAll(dir, 0o750); err != nil { + t.Fatalf("mkdir: %v", err) + } + + if err := os.WriteFile( + filepath.Join(dir, "01 Intro.flac"), []byte(disc), 0o600, + ); err != nil { + t.Fatalf("write: %v", err) + } + } + + dst := t.TempDir() + + got, err := s.collect(Candidate{Files: []CandidateFile{ + {Path: `\m\Album\CD1\01 Intro.flac`, IsAudio: true}, + {Path: `\m\Album\CD2\01 Intro.flac`, IsAudio: true}, + }}, dst) + if err != nil { + t.Fatalf("collect: %v", err) + } + + if len(got.Files) != 2 { + t.Fatalf("collected %d files, want 2", len(got.Files)) + } + + for _, disc := range []string{"CD1", "CD2"} { + data, err := os.ReadFile(filepath.Join(dst, disc, "01 Intro.flac")) + if err != nil || string(data) != disc { + t.Errorf("%s's file missing or overwritten: %q, %v", disc, data, err) + } + + if hint := ParsePath(filepath.Join(dst, disc, "01 Intro.flac")); hint.Disc == 0 { + t.Errorf("staged %s file lost its disc number", disc) + } + } +} + +// A two-disc release whose discs both number from 01 aligns completely +// once the disc comes from the folder; before, disc 2's 01 collided with +// disc 1's. +func TestMultiDiscCandidateAlignsEveryTrack(t *testing.T) { + t.Parallel() + + dl := Download{ + ReleaseMBID: "the-wall", + Artist: "Pink Floyd", + Album: "The Wall", + Expected: []ExpectedTrack{ + {DiscNumber: 1, Position: 1, Title: "In the Flesh?"}, + {DiscNumber: 1, Position: 2, Title: "The Thin Ice"}, + {DiscNumber: 2, Position: 1, Title: "Hey You"}, + {DiscNumber: 2, Position: 2, Title: "Is There Anybody Out There?"}, + }, + } + + c := Candidate{ + Title: "The Wall", + Files: []CandidateFile{ + {Path: `\m\Pink Floyd - The Wall\CD1\01 In the Flesh.flac`, Size: 1}, + {Path: `\m\Pink Floyd - The Wall\CD1\02 The Thin Ice.flac`, Size: 1}, + {Path: `\m\Pink Floyd - The Wall\CD2\01 Hey You.flac`, Size: 1}, + {Path: `\m\Pink Floyd - The Wall\CD2\02 Is There Anybody Out There.flac`, Size: 1}, + }, + } + + got := Score(dl, c, 50, AutoDownloadPrefs{}) + + if got.Match.Completeness != 1 { + t.Errorf("completeness = %f, want 1", got.Match.Completeness) + } + + if got.Match.AlbumFit < 0.99 { + t.Errorf("album fit = %f, want the album's own name to match", got.Match.AlbumFit) + } + + if got.Match.Overall < minMatch { + t.Errorf("match = %f, want it to clear the auto-pick bar %f", got.Match.Overall, minMatch) + } +} + +// Ten files against a ten-track album is not a complete album when only +// three of them are its tracks. +func TestCompletenessCountsTracksNotFiles(t *testing.T) { + t.Parallel() + + dl := okComputer() + + c := Candidate{Title: "OK Computer", Files: []CandidateFile{ + {Path: `\m\Radiohead - OK Computer\Airbag.flac`, Size: 1}, + {Path: `\m\Radiohead - OK Computer\Paranoid Android.flac`, Size: 1}, + {Path: `\m\Radiohead - OK Computer\Exit Music (For a Film).flac`, Size: 1}, + {Path: `\m\Radiohead - OK Computer\Creep.flac`, Size: 1}, + }} + + got := Score(dl, c, 50, AutoDownloadPrefs{}) + + if got.Match.Completeness > 0.76 { + t.Errorf( + "completeness = %f with 3 of 4 tracks present, want at most 0.75", + got.Match.Completeness, + ) + } +} diff --git a/backend/download/concurrency_test.go b/backend/download/concurrency_test.go index 308e733..2181582 100644 --- a/backend/download/concurrency_test.go +++ b/backend/download/concurrency_test.go @@ -47,7 +47,7 @@ func grabAll( go func() { defer wg.Done() - f.manager.grab(ctx, dl, candidate, nil) + f.manager.grab(ctx, dl, candidate, nil, false) }() } @@ -79,9 +79,9 @@ func TestConcurrencyForPrefersOverrideThenKind(t *testing.T) { want int }{ { - name: "slskd defaults to one", + name: "slskd defaults to a few peers", cfg: Config{Kind: KindSlskd}, - want: 1, + want: 3, }, { name: "usenet defaults higher", @@ -92,9 +92,9 @@ func TestConcurrencyForPrefersOverrideThenKind(t *testing.T) { name: "explicit override wins", cfg: Config{ Kind: KindSlskd, - Settings: map[string]string{concurrencyKey: "3"}, + Settings: map[string]string{concurrencyKey: "1"}, }, - want: 3, + want: 1, }, { name: "nonsense override falls back", @@ -102,7 +102,7 @@ func TestConcurrencyForPrefersOverrideThenKind(t *testing.T) { Kind: KindSlskd, Settings: map[string]string{concurrencyKey: "not a number"}, }, - want: 1, + want: 3, }, { name: "zero override falls back", @@ -110,7 +110,7 @@ func TestConcurrencyForPrefersOverrideThenKind(t *testing.T) { Kind: KindSlskd, Settings: map[string]string{concurrencyKey: "0"}, }, - want: 1, + want: 3, }, { name: "unknown kind falls back to the global default", @@ -126,9 +126,9 @@ func TestConcurrencyForPrefersOverrideThenKind(t *testing.T) { } } -// The reason the per-provider cap exists: a Soulseek daemon capped at -// one transfer must serialize, even when the global cap would allow -// more and the user has queued several albums at once. +// The reason the per-provider cap exists: a daemon capped at one +// transfer must serialize, even when the global cap would allow more and +// the user has queued several albums at once. func TestPerProviderCapSerializesTransfers(t *testing.T) { t.Parallel() @@ -142,6 +142,7 @@ func TestPerProviderCapSerializesTransfers(t *testing.T) { ID: 1, Kind: KindSlskd, Priority: 50, + Settings: map[string]string{concurrencyKey: "1"}, }, slow) // Three requests against the same one-at-a-time provider. @@ -210,8 +211,8 @@ func TestSyncSemaphoresReplacesChangedLimits(t *testing.T) { f.manager.installProvider(Config{ID: 1, Kind: KindSlskd}, nil) first := f.manager.semaphoreFor(1) - if cap(first) != 1 { - t.Fatalf("slskd semaphore cap = %d, want 1", cap(first)) + if want := kindConcurrency[KindSlskd]; cap(first) != want { + t.Fatalf("slskd semaphore cap = %d, want %d", cap(first), want) } // Same limit: the semaphore is kept, so in-flight accounting is not diff --git a/backend/download/fallback_test.go b/backend/download/fallback_test.go new file mode 100644 index 0000000..48565d5 --- /dev/null +++ b/backend/download/fallback_test.go @@ -0,0 +1,261 @@ +package download + +import ( + "context" + "errors" + "os" + "testing" +) + +// A transfer that fails on one copy of an album is not a failed +// download while another acceptable copy exists. On Soulseek the usual +// failure is one peer being offline, with several others offering the +// same folder. + +var errPeerOffline = errors.New("peer went offline") + +func TestManagerFallsBackToTheNextCandidate(t *testing.T) { + t.Parallel() + + f := newManagerFixture(t) + + // The failing source ranks first on priority, so the fallback is + // what reaches the one that works. + bad := fakeWithAlbum(1, "offline-peer", ".flac") + bad.GrabErr = errPeerOffline + good := fakeWithAlbum(2, "online-peer", ".flac") + + f.manager.installProvider(Config{ID: 1, Priority: 90}, bad) + f.manager.installProvider(Config{ID: 2, Priority: 10}, good) + + dl := fourTrackDownload() + + if _, err := f.manager.Start(context.Background(), dl); err != nil { + t.Fatalf("Start: %v", err) + } + + waitForDownloadState(t, f.store, dl.ID, StateComplete) + + if bad.GrabCalls != 1 || good.GrabCalls != 1 { + t.Errorf( + "grabs: failing=%d working=%d, want 1 and 1", + bad.GrabCalls, good.GrabCalls, + ) + } + + // The abandoned attempt's staging goes with it; only a request that + // fails outright keeps its staging for inspection. + waitFor(t, func() bool { + entries, err := os.ReadDir(f.staging.Root()) + + return err == nil && len(entries) == 0 + }, "the failed attempt's staging was never released") +} + +// Falling back must not lower the bar. A second choice outside the +// user's guardrails is not a choice auto-pick may make, first or second. +func TestManagerFallbackRespectsTheGuardrails(t *testing.T) { + t.Parallel() + + f := newManagerFixture(t) + f.manager.SetPreferences(AutoDownloadPrefs{MaxSizeMB: 50}) + + bad := fakeWithAlbum(1, "offline-peer", ".flac") + bad.GrabErr = errPeerOffline + bad.Candidates[0].TotalSize = 40 << 20 + + huge := fakeWithAlbum(2, "oversized", ".flac") + huge.Candidates[0].TotalSize = 900 << 20 + + f.manager.installProvider(Config{ID: 1, Priority: 90}, bad) + f.manager.installProvider(Config{ID: 2, Priority: 10}, huge) + + dl := fourTrackDownload() + + if _, err := f.manager.Start(context.Background(), dl); err != nil { + t.Fatalf("Start: %v", err) + } + + waitForDownloadState(t, f.store, dl.ID, StateFailed) + + if huge.GrabCalls != 0 { + t.Errorf("fell back to a candidate over the size ceiling") + } +} + +// A copy the user picked by hand is the copy they asked for. Quietly +// substituting another is a decision they did not make. +func TestManagerPickDoesNotFallBack(t *testing.T) { + t.Parallel() + + f := newManagerFixture(t) + + bad := fakeWithAlbum(1, "offline-peer", ".flac") + bad.GrabErr = errPeerOffline + good := fakeWithAlbum(2, "online-peer", ".flac") + + f.manager.installProvider(Config{ID: 1, Priority: 90}, bad) + f.manager.installProvider(Config{ID: 2, Priority: 10}, good) + + // A ceiling below both copies parks the result set for the user. + f.manager.SetPreferences(AutoDownloadPrefs{MaxSizeMB: 1}) + + bad.Candidates[0].TotalSize = 30 << 20 + good.Candidates[0].TotalSize = 30 << 20 + + dl := fourTrackDownload() + + if _, err := f.manager.Start(context.Background(), dl); err != nil { + t.Fatalf("Start: %v", err) + } + + if err := f.manager.Pick( + context.Background(), dl.ID, "offline-peer-cand", + ); err != nil { + t.Fatalf("Pick: %v", err) + } + + waitForDownloadState(t, f.store, dl.ID, StateFailed) + + if good.GrabCalls != 0 { + t.Errorf("a hand-picked grab fell back to another candidate") + } +} + +// Fallback is for surviving an offline peer or two, not for walking a +// forty-peer list for six hours. +func TestManagerFallbackIsBounded(t *testing.T) { + t.Parallel() + + f := newManagerFixture(t) + + var providers []*FakeProvider + + for i := int64(1); i <= maxGrabAttempts+2; i++ { + p := fakeWithAlbum(i, "peer-"+itoa(int(i)), ".flac") + p.GrabErr = errPeerOffline + + f.manager.installProvider(Config{ID: i, Priority: 50}, p) + providers = append(providers, p) + } + + dl := fourTrackDownload() + + if _, err := f.manager.Start(context.Background(), dl); err != nil { + t.Fatalf("Start: %v", err) + } + + waitForDownloadState(t, f.store, dl.ID, StateFailed) + + grabs := 0 + for _, p := range providers { + grabs += p.GrabCalls + } + + if grabs != maxGrabAttempts { + t.Errorf("grabs = %d, want %d", grabs, maxGrabAttempts) + } +} + +// The veto judges the best candidate *inside* the guardrails, so the +// grab has to take that one — not the overall best, which may be the +// very copy the user said not to take unattended. +func TestManagerAutoPickTakesTheBestEligibleCandidate(t *testing.T) { + t.Parallel() + + f := newManagerFixture(t) + f.manager.SetPreferences(AutoDownloadPrefs{MaxSizeMB: 50}) + + huge := fakeWithAlbum(1, "oversized", ".flac") + huge.Candidates[0].TotalSize = 900 << 20 + + fits := fakeWithAlbum(2, "fits", ".flac") + fits.Candidates[0].TotalSize = 40 << 20 + + f.manager.installProvider(Config{ID: 1, Priority: 90}, huge) + f.manager.installProvider(Config{ID: 2, Priority: 10}, fits) + + dl := fourTrackDownload() + + ranked, err := f.manager.Start(context.Background(), dl) + if err != nil { + t.Fatalf("Start: %v", err) + } + + if ranked[0].ID != "oversized-cand" { + t.Fatalf("fixture: best overall is %s, want the oversized copy", ranked[0].ID) + } + + waitForDownloadState(t, f.store, dl.ID, StateComplete) + + if huge.GrabCalls != 0 || fits.GrabCalls != 1 { + t.Errorf( + "grabs: oversized=%d fits=%d, want 0 and 1", + huge.GrabCalls, fits.GrabCalls, + ) + } +} + +// On Soulseek a failure is the peer's, so every folder that peer offered +// goes with it. Elsewhere a failure is the release's, and one indexer's +// other releases are still worth trying. +func TestRuledOutBy(t *testing.T) { + t.Parallel() + + failed := []Candidate{ + {ID: "slskd:alice:Album", Kind: KindSlskd, ProviderID: 1, Origin: "alice"}, + {ID: "tracker-1", Kind: KindProwlarr, ProviderID: 2, Origin: "indexer"}, + } + + cases := []struct { + name string + c Candidate + want bool + }{ + { + name: "the same candidate", + c: Candidate{ID: "tracker-1", Kind: KindProwlarr, ProviderID: 2, Origin: "indexer"}, + want: true, + }, + { + name: "another folder from a failed peer", + c: Candidate{ + ID: "slskd:alice:Album (2)", + Kind: KindSlskd, + ProviderID: 1, + Origin: "alice", + }, + want: true, + }, + { + name: "another peer", + c: Candidate{ID: "slskd:bob:Album", Kind: KindSlskd, ProviderID: 1, Origin: "bob"}, + want: false, + }, + { + name: "another release from the same indexer", + c: Candidate{ID: "tracker-2", Kind: KindProwlarr, ProviderID: 2, Origin: "indexer"}, + want: false, + }, + { + name: "a peer of the same name on a different daemon", + c: Candidate{ + ID: "slskd:alice:Album", + Kind: KindSlskd, + ProviderID: 3, + Origin: "alice", + }, + want: false, + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + if got := ruledOutBy(tc.c, failed); got != tc.want { + t.Errorf("ruledOutBy = %v, want %v", got, tc.want) + } + }) + } +} diff --git a/backend/download/keyedlock.go b/backend/download/keyedlock.go new file mode 100644 index 0000000..95f45f8 --- /dev/null +++ b/backend/download/keyedlock.go @@ -0,0 +1,77 @@ +package download + +import ( + "context" + "sync" +) + +// keyedLock is a set of mutexes created on demand, one per key, that +// honour a context while waiting. An entry lives only while someone +// holds or waits on it, so a key per Soulseek peer or per folder name +// does not accumulate for the life of the process. +type keyedLock[K comparable] struct { + mu sync.Mutex + held map[K]*keyedEntry +} + +type keyedEntry struct { + ch chan struct{} + + // refs counts holders and waiters; the entry is dropped at zero. + refs int +} + +// acquire blocks until k is free or ctx ends, and returns the function +// that frees it. +func (l *keyedLock[K]) acquire(ctx context.Context, k K) (func(), error) { + l.mu.Lock() + + if l.held == nil { + l.held = map[K]*keyedEntry{} + } + + e, ok := l.held[k] + if !ok { + e = &keyedEntry{ch: make(chan struct{}, 1)} + l.held[k] = e + } + + e.refs++ + + l.mu.Unlock() + + select { + case e.ch <- struct{}{}: + case <-ctx.Done(): + l.drop(k, e) + + return nil, ctx.Err() + } + + var once sync.Once + + return func() { + once.Do(func() { + <-e.ch + l.drop(k, e) + }) + }, nil +} + +func (l *keyedLock[K]) drop(k K, e *keyedEntry) { + l.mu.Lock() + defer l.mu.Unlock() + + e.refs-- + if e.refs == 0 { + delete(l.held, k) + } +} + +// size reports how many keys are held or awaited, for tests. +func (l *keyedLock[K]) size() int { + l.mu.Lock() + defer l.mu.Unlock() + + return len(l.held) +} diff --git a/backend/download/manager.go b/backend/download/manager.go index 2b284b5..3b429de 100644 --- a/backend/download/manager.go +++ b/backend/download/manager.go @@ -59,13 +59,15 @@ const concurrencyKey = "maxConcurrent" // A single global cap is the wrong shape here: usenet and torrent // clients are built to run many transfers at once and are throttled by // bandwidth, while Soulseek transfers come from one person's home -// upload slot. Hitting the same peer with parallel requests gets you -// queued behind everyone else at best and banned at worst, so slskd is -// capped at one — the polite number, and the one that actually -// completes fastest, because a Soulseek peer serves one file at a time -// regardless of how many you ask for. +// upload slot. Politeness there is per *peer* — asking one user for two +// folders at once gets you queued behind everyone else at best and +// banned at worst — and the manager holds that line separately, one +// grab per peer (peerLocks). Two different users do not compete for +// anyone's slot, so the daemon-wide number only bounds how many peers +// are asked at once, and one slow peer no longer serialises every other +// Soulseek download behind it. var kindConcurrency = map[Kind]int{ - KindSlskd: 1, + KindSlskd: 3, KindYtDlp: 2, KindQBittorrent: 4, KindSABnzbd: 4, @@ -155,6 +157,11 @@ type Manager struct { semMu sync.Mutex provSem map[int64]chan struct{} + // peerLocks holds one grab per Soulseek peer, taken before any + // slot: a grab waiting for a busy peer must not sit on a provider + // slot another peer could be using. + peerLocks keyedLock[peerKey] + // delegatePoll is how often delegating managers are asked for // status. A field rather than the constant so tests can drive the // full delegate flow without sleeping through it. @@ -577,12 +584,12 @@ func (m *Manager) Start( )) } - if m.AutoPickable(dl, ranked) { + if pick, ok := autoPick(dl, ranked, m.preferences()); ok { if job != nil { job.Logf(jobs.LevelInfo, "Auto-selected best candidate") } - go m.grab(context.WithoutCancel(ctx), dl, ranked[0], job) + go m.grab(context.WithoutCancel(ctx), dl, pick, job, true) return ranked, nil } @@ -622,6 +629,13 @@ func (m *Manager) Attempt( return false, veto, nil } + pick, ok := autoPick(dl, ranked, m.preferences()) + if !ok { + // Unreachable while autoPick and AutoPickVeto agree; kept so a + // future divergence refuses rather than grabbing blind. + return false, "no candidate clears the auto-download bar", nil + } + if err := m.store.CreateDownload(ctx, dl); err != nil { return false, "", err } @@ -643,7 +657,7 @@ func (m *Manager) Attempt( )) } - go m.grab(context.WithoutCancel(ctx), dl, ranked[0], job) + go m.grab(context.WithoutCancel(ctx), dl, pick, job, true) return true, "", nil } @@ -678,7 +692,7 @@ func (m *Manager) Pick( job := m.startJob(dl) - go m.grab(context.WithoutCancel(ctx), dl, *chosen, job) + go m.grab(context.WithoutCancel(ctx), dl, *chosen, job, false) return nil } @@ -702,13 +716,20 @@ func (m *Manager) Cancel(ctx context.Context, downloadID string) error { return nil } -// grab drives one candidate all the way to the library. It runs on its +// grab drives one request all the way to the library. It runs on its // own goroutine and owns the job from here on. +// +// When fallback is set and a candidate's transfer fails, the next +// candidate that auto-pick would itself have accepted is tried in its +// place (see nextCandidate). It is set for the two unattended routes +// and not for a candidate the user picked by hand: they chose that copy, +// and quietly substituting another is a decision they did not make. func (m *Manager) grab( ctx context.Context, dl Download, c Candidate, job *jobs.Handle, + fallback bool, ) { ctx, cancel := context.WithTimeout(ctx, grabTimeout) defer cancel() @@ -723,6 +744,99 @@ func (m *Manager) grab( m.actMu.Unlock() }() + var failed []Candidate + + for { + out := m.attemptGrab(ctx, dl, c, job) + if out.err == nil { + m.finishGrab(ctx, dl, out.item, out.imported, job) + + return + } + + failed = append(failed, c) + + next, ok := m.nextCandidate(ctx, dl, failed, out, fallback) + if !ok { + m.failDownload(ctx, job, dl.ID, out.err) + + return + } + + m.logger.Info( + "download candidate failed; trying the next", + "download", dl.ID, + "failed", c.ID, + "next", next.ID, + "error", out.err, + ) + + if job != nil { + job.Logf(jobs.LevelWarn, fmt.Sprintf( + "%s failed (%v); trying %s instead", + describeCandidate(c), out.err, describeCandidate(next), + )) + } + + // The failed attempt's staging holds at most a partial folder + // nobody is going to import, and the next attempt reserves its + // own. Only the final failure keeps its staging for inspection. + if out.item.StagingDir != "" { + if err := m.staging.Release(out.item.StagingDir); err != nil { + m.logger.Warn("could not release staging dir", "error", err) + } + } + + c = next + } +} + +// maxGrabAttempts bounds how many candidates one request will try. A +// popular album can have dozens of peers; the point of falling back is +// to survive the ordinary one or two that are offline, not to walk the +// whole list for six hours. +const maxGrabAttempts = 3 + +// peerKey names one Soulseek user on one daemon. The same username on +// two daemons is two logins and two queues. +type peerKey struct { + provider int64 + peer string +} + +// peerKeyFor returns the peer a candidate is fetched from, when the +// source is one where asking a peer for two things at once is rude. +func peerKeyFor(c Candidate) (peerKey, bool) { + if c.Kind != KindSlskd || c.Origin == "" { + return peerKey{}, false + } + + return peerKey{provider: c.ProviderID, peer: c.Origin}, true +} + +// grabOutcome is how one candidate's attempt ended. +type grabOutcome struct { + item DownloadItem + imported ImportResult + err error + + // retryable reports whether another candidate might succeed where + // this one failed: the transfer failed, or delivered too little of + // the album. Anything else — no staging space, no library root, a + // tag write failing — would fail the next candidate identically. + retryable bool +} + +// attemptGrab takes one candidate through transfer and import. It +// records the item's own failure, but not the download's: whether the +// download has failed is the caller's decision, since another candidate +// may yet succeed. +func (m *Manager) attemptGrab( + ctx context.Context, + dl Download, + c Candidate, + job *jobs.Handle, +) grabOutcome { // Who will move the bytes is decided before any slot is taken, so // the transfer waits in its own provider's queue rather than in a // global one. A delegate takes no slot at all: the transfer is @@ -731,21 +845,26 @@ func (m *Manager) grab( // work against our budget. plan, err := m.planTransfer(dl, c) if err != nil { - m.failDownload(ctx, job, dl.ID, err) - - return + return grabOutcome{err: err} } if !plan.delegated() { + if key, ok := peerKeyFor(c); ok { + release, err := m.peerLocks.acquire(ctx, key) + if err != nil { + return grabOutcome{err: err} + } + + defer release() + } + provSem := m.semaphoreFor(plan.transportID) select { case provSem <- struct{}{}: defer func() { <-provSem }() case <-ctx.Done(): - m.failDownload(ctx, job, dl.ID, ctx.Err()) - - return + return grabOutcome{err: ctx.Err()} } globalSem := m.globalSem() @@ -754,9 +873,7 @@ func (m *Manager) grab( case globalSem <- struct{}{}: defer func() { <-globalSem }() case <-ctx.Done(): - m.failDownload(ctx, job, dl.ID, ctx.Err()) - - return + return grabOutcome{err: ctx.Err()} } } @@ -771,24 +888,30 @@ func (m *Manager) grab( dir, err := m.staging.Reserve(item.ID) if err != nil { - m.failDownload(ctx, job, dl.ID, err) - - return + return grabOutcome{err: err} } item.StagingDir = dir if err := m.store.CreateItem(ctx, item); err != nil { - m.failDownload(ctx, job, dl.ID, err) + return grabOutcome{item: item, err: err} + } - return + fail := func(err error, retryable bool) grabOutcome { + if serr := m.store.SetItemState( + ctx, item.ID, StateFailed, err.Error(), + ); serr != nil { + m.logger.Warn("could not record item failure", "error", serr) + } + + return grabOutcome{item: item, err: err, retryable: retryable} } result, err := m.transfer(ctx, dl, item, plan, job) if err != nil { - m.failItem(ctx, job, item, dl.ID, err) - - return + // A delegate's failure is the external manager's verdict on the + // whole request, not on one copy of it. + return fail(err, !plan.delegated()) } m.setStates(ctx, dl.ID, item.ID, StateImporting) @@ -798,42 +921,116 @@ func (m *Manager) grab( job.SetStages(importStages(2)) } - var imported ImportResult - if result.Delegated { // The external manager already placed and tagged these files in // its own library. Moving them out from under a system that is // still managing them would be worse than useless, so the files // are recorded where they are and the library scan picks them // up in place. - imported = ImportResult{Paths: result.Files} - if job != nil { job.Logf(jobs.LevelInfo, fmt.Sprintf( "External manager imported %d files; recording them in place", len(result.Files), )) } - } else { - opts := m.importOptions() - opts.WriteTags = true - opts.LibraryRoot, err = m.library.LibraryPath(dl.LibraryID) - if err != nil { - m.failItem(ctx, job, item, dl.ID, - fmt.Errorf("resolve library root: %w", err)) - - return - } - - imported, err = m.importer.Import(ctx, dl, result, opts) - if err != nil { - m.failItem(ctx, job, item, dl.ID, err) - - return + return grabOutcome{ + item: item, + imported: ImportResult{Paths: result.Files}, } } + opts := m.importOptions() + opts.WriteTags = true + + opts.LibraryRoot, err = m.library.LibraryPath(dl.LibraryID) + if err != nil { + return fail(fmt.Errorf("resolve library root: %w", err), false) + } + + imported, err := m.importer.Import(ctx, dl, result, opts) + if err != nil { + return fail(err, errors.Is(err, ErrTooIncomplete)) + } + + return grabOutcome{item: item, imported: imported} +} + +// nextCandidate picks the candidate to try after the ones in failed. +// +// It only ever offers a candidate auto-pick would have taken on its own +// (autoAcceptable), so falling back cannot lower the bar an unattended +// download is held to: the second choice has to clear the same gates +// the first did. +// +// On Soulseek a failure belongs to the *peer* — offline, refusing, or +// holding us in a queue — so every folder that peer offered is skipped +// with it. Elsewhere a failure belongs to the release, and only that +// candidate is. +func (m *Manager) nextCandidate( + ctx context.Context, + dl Download, + failed []Candidate, + out grabOutcome, + fallback bool, +) (Candidate, bool) { + if !fallback || !out.retryable || ctx.Err() != nil || + len(failed) >= maxGrabAttempts { + return Candidate{}, false + } + + m.resMu.RLock() + ranked := m.results[dl.ID] + m.resMu.RUnlock() + + prefs := m.preferences() + + for _, c := range ranked { + if ruledOutBy(c, failed) || !autoAcceptable(dl, c, prefs) { + continue + } + + return c, true + } + + return Candidate{}, false +} + +// ruledOutBy reports whether a failure among failed also rules out c. +func ruledOutBy(c Candidate, failed []Candidate) bool { + for _, f := range failed { + if c.ID == f.ID && c.ProviderID == f.ProviderID { + return true + } + + if c.Kind == KindSlskd && f.Kind == KindSlskd && + c.ProviderID == f.ProviderID && c.Origin != "" && + c.Origin == f.Origin { + return true + } + } + + return false +} + +// describeCandidate names a candidate for the job log. +func describeCandidate(c Candidate) string { + if c.Origin != "" { + return fmt.Sprintf("%q from %s", c.Title, c.Origin) + } + + return fmt.Sprintf("%q", c.Title) +} + +// finishGrab records a successful import and retires what the request +// was holding. +func (m *Manager) finishGrab( + ctx context.Context, + dl Download, + item DownloadItem, + imported ImportResult, + job *jobs.Handle, +) { if err := m.store.SetItemImported( ctx, item.ID, imported.Paths, ); err != nil { @@ -1177,23 +1374,6 @@ func (m *Manager) failDownload( } } -// failItem records an item-level failure and fails its download. -func (m *Manager) failItem( - ctx context.Context, - job *jobs.Handle, - item DownloadItem, - downloadID string, - err error, -) { - if serr := m.store.SetItemState( - ctx, item.ID, StateFailed, err.Error(), - ); serr != nil { - m.logger.Warn("could not record item failure", "error", serr) - } - - m.failDownload(ctx, job, downloadID, err) -} - // startJob registers the request in the background jobs panel. func (m *Manager) startJob(dl Download) *jobs.Handle { if m.jobsReg == nil { diff --git a/backend/download/pathmatch.go b/backend/download/pathmatch.go index 78a2249..5348c4d 100644 --- a/backend/download/pathmatch.go +++ b/backend/download/pathmatch.go @@ -66,6 +66,14 @@ var ( // separatorPattern splits "Artist - Album" style folder names. separatorPattern = regexp.MustCompile(`\s+[-–—]\s+`) + + // discFolderPattern matches a directory that holds one disc of an + // album rather than the album: "CD1", "CD 2", "Disc 3", "Disk-1", + // "[Disc 2]", "CD1 - The Early Years". A number is required, so a + // folder merely called "CDs" is not one. + discFolderPattern = regexp.MustCompile( + `(?i)^\s*[\[(]?\s*(?:cd|disc|disk)\s*[-_.#]?\s*(\d{1,2})\b`, + ) ) // FormatForPath returns the audio format implied by a path's extension, @@ -94,16 +102,63 @@ type TrackHint struct { Folder string } +// discFolder reports whether a directory name is one disc of an album, +// and which. +func discFolder(name string) (int, bool) { + m := discFolderPattern.FindStringSubmatch(name) + if m == nil { + return 0, false + } + + n, err := strconv.Atoi(m[1]) + if err != nil || n == 0 { + return 0, false + } + + return n, true +} + +// AlbumDir is the directory that holds a file's *album*: its parent, +// or its grandparent when the parent is a disc folder. +// +// Multi-disc rips are shared as `Album/CD1/…` and `Album/CD2/…`, and +// grouping candidates by the immediate parent split one album into two +// half-albums, each titled "CD1". Neither could clear the completeness +// or album-title bars, so a multi-disc release could not be auto-picked +// at all. A disc folder at the root has no album above it and is +// returned as it is. +func AlbumDir(p string) string { + dir := path.Dir(strings.ReplaceAll(p, `\`, "/")) + + if _, ok := discFolder(path.Base(dir)); !ok { + return dir + } + + parent := path.Dir(dir) + if parent == "." || parent == "/" || parent == "" { + return dir + } + + return parent +} + // ParsePath extracts what it can from one candidate file path. func ParsePath(p string) TrackHint { // Soulseek paths are Windows-style; normalize before splitting. norm := strings.ReplaceAll(p, `\`, "/") base := path.Base(norm) - folder := path.Base(path.Dir(norm)) name := strings.TrimSuffix(base, path.Ext(base)) - hint := TrackHint{Folder: cleanAlbumName(folder)} + // The album's name is the album directory's, not a disc folder's, + // and the disc folder is where a multi-disc rip says which disc a + // file is on. A disc number in the filename ("2-01 …") is more + // specific and overrides it below. + hint := TrackHint{Folder: cleanAlbumName(path.Base(AlbumDir(norm)))} + + if disc, ok := discFolder(path.Base(path.Dir(norm))); ok { + hint.Disc = disc + } if m := trackNumPattern.FindStringSubmatch(name); m != nil { if m[1] != "" { @@ -203,21 +258,72 @@ func AnnotateFiles(files []CandidateFile) []CandidateFile { // matchFiles aligns a candidate's audio files to the expected tracklist // and returns the per-file assignment plus the mean title similarity of -// the aligned pairs. +// the aligned pairs. alignFiles is the same alignment with the +// duration evidence as well. +func matchFiles( + files []CandidateFile, + expected []ExpectedTrack, +) ([]CandidateFile, float64) { + a := alignFiles(files, expected) + + return a.files, a.titleFit +} + +// alignment is what aligning a candidate to a tracklist found. +type alignment struct { + files []CandidateFile + + // titleFit is the mean title similarity over aligned pairs. + titleFit float64 + + // durationFit is the mean duration agreement over aligned pairs + // where both sides state a length, and timedPairs is how many such + // pairs there were. + durationFit float64 + timedPairs int + aligned int +} + +// durationAgreement scores how well a file's length matches the +// expected track's, in 0..1. Rips of the same master differ by a +// second or two of silence; a different edit, a live take or a +// truncated file differs by tens of seconds. +func durationAgreement(got, want int64) float64 { + const ( + exactMillis = 3_000 + wrongMillis = 30_000 + ) + + d := got - want + if d < 0 { + d = -d + } + + switch { + case d <= exactMillis: + return 1 + case d >= wrongMillis: + return 0 + default: + return 1 - float64(d-exactMillis)/float64(wrongMillis-exactMillis) + } +} + +// alignFiles aligns a candidate's audio files to the expected tracklist. // // Alignment is greedy by score rather than optimal: candidate folders // are small (a few dozen files at most) and the common cases — correct // track numbers, or clean "NN Title" names — are unambiguous, so the // extra machinery of Hungarian assignment buys nothing here. -func matchFiles( +func alignFiles( files []CandidateFile, expected []ExpectedTrack, -) ([]CandidateFile, float64) { +) alignment { annotated := make([]CandidateFile, len(files)) copy(annotated, files) if len(expected) == 0 { - return annotated, 0 + return alignment{files: annotated} } hints := make([]TrackHint, len(annotated)) @@ -230,8 +336,19 @@ func matchFiles( var ( total float64 matched int + + durTotal float64 + timed int ) + // timing adds a pair's duration evidence when both sides state one. + timing := func(f CandidateFile, e ExpectedTrack) { + if f.LengthMillis > 0 && e.LengthMillis > 0 { + durTotal += durationAgreement(f.LengthMillis, e.LengthMillis) + timed++ + } + } + // Pass 1: trust explicit track numbers when they are unique and in // range. A folder that numbers its files correctly is the strong // case, and title comparison only adds noise there. @@ -250,6 +367,8 @@ func matchFiles( total += autotag.TitleSimilarity(hints[i].Title, expected[idx].Title) matched++ + + timing(annotated[i], expected[idx]) } // Pass 2: title similarity for whatever is left. @@ -284,13 +403,26 @@ func matchFiles( total += bestSim matched++ + + timing(annotated[i], expected[bestIdx]) } if matched == 0 { - return annotated, 0 + return alignment{files: annotated} } - return annotated, total / float64(matched) + a := alignment{ + files: annotated, + titleFit: total / float64(matched), + timedPairs: timed, + aligned: matched, + } + + if timed > 0 { + a.durationFit = durTotal / float64(timed) + } + + return a } // indexForPosition finds the expected track at a disc/track position. diff --git a/backend/download/peer_concurrency_test.go b/backend/download/peer_concurrency_test.go new file mode 100644 index 0000000..113604f --- /dev/null +++ b/backend/download/peer_concurrency_test.go @@ -0,0 +1,224 @@ +package download + +import ( + "context" + "errors" + "path/filepath" + "sync" + "testing" + "time" +) + +// Soulseek politeness is per peer, not per daemon (#272). + +func TestKeyedLockSerialisesOneKeyOnly(t *testing.T) { + t.Parallel() + + var l keyedLock[string] + + ctx := context.Background() + + releaseA, err := l.acquire(ctx, "a") + if err != nil { + t.Fatalf("acquire a: %v", err) + } + + // Another key is free while "a" is held. + releaseB, err := l.acquire(ctx, "b") + if err != nil { + t.Fatalf("acquire b: %v", err) + } + + releaseB() + + // The same key waits, and gives up with its context. + short, cancel := context.WithTimeout(ctx, 20*time.Millisecond) + defer cancel() + + if _, err := l.acquire(short, "a"); !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("second acquire of a held key = %v, want the deadline", err) + } + + releaseA() + releaseA() // Idempotent: a second call must not free someone else's hold. + + if n := l.size(); n != 0 { + t.Errorf("%d keys left behind, want none once nobody holds or waits", n) + } +} + +// grabEach runs one grab per candidate and returns a function that waits +// for all of them; grabAll's reasons for waiting apply. +func grabEach(t *testing.T, f managerFixture, cands []Candidate) func() { + t.Helper() + + ctx := context.Background() + + var wg sync.WaitGroup + + for i, c := range cands { + dl := fourTrackDownload() + dl.ID = "dl-" + string(rune('a'+i)) + + if err := f.store.CreateDownload(ctx, dl); err != nil { + t.Fatalf("CreateDownload: %v", err) + } + + wg.Add(1) + + go func() { + defer wg.Done() + + f.manager.grab(ctx, dl, c, nil, false) + }() + } + + return func() { + done := make(chan struct{}) + + go func() { + wg.Wait() + close(done) + }() + + select { + case <-done: + case <-time.After(5 * time.Second): + t.Error("transfers did not finish") + } + } +} + +func slskdCandidates(p *FakeProvider, peers ...string) []Candidate { + out := make([]Candidate, 0, len(peers)) + + for i, peer := range peers { + c := p.Candidates[0] + c.ID = c.ID + "-" + itoa(i) + c.Kind = KindSlskd + c.ProviderID = 1 + c.Origin = peer + out = append(out, c) + } + + return out +} + +// Three albums from one user are asked for one at a time, even though +// the daemon would allow three transfers. +func TestOnePeerIsAskedForOneThingAtATime(t *testing.T) { + t.Parallel() + + f := newManagerFixture(t) + f.manager.SetMaxConcurrent(4) + + p := fakeWithAlbum(1, "slskd", ".flac") + p.GrabGate = make(chan struct{}) + + f.manager.installProvider(Config{ID: 1, Kind: KindSlskd, Priority: 50}, p) + + wait := grabEach(t, f, slskdCandidates(p, "alice", "alice", "alice")) + + waitFor(t, func() bool { return p.GrabCallCount() >= 1 }, "no grab started") + time.Sleep(150 * time.Millisecond) + + if got := p.MaxParallelGrabs(); got != 1 { + t.Errorf("%d simultaneous grabs from one peer, want 1", got) + } + + close(p.GrabGate) + + waitFor(t, func() bool { return p.GrabCallCount() == 3 }, "queued grabs never ran") + wait() + + if n := f.manager.peerLocks.size(); n != 0 { + t.Errorf("%d peer locks left behind", n) + } +} + +// Different users run at once, up to the daemon's cap — the point of +// the change: one slow peer no longer holds up every other. +func TestDifferentPeersRunTogether(t *testing.T) { + t.Parallel() + + f := newManagerFixture(t) + f.manager.SetMaxConcurrent(8) + + p := fakeWithAlbum(1, "slskd", ".flac") + p.GrabGate = make(chan struct{}) + + f.manager.installProvider(Config{ID: 1, Kind: KindSlskd, Priority: 50}, p) + + wait := grabEach(t, f, slskdCandidates(p, "alice", "bob", "carol", "dave")) + + waitFor( + t, + func() bool { return p.MaxParallelGrabs() >= kindConcurrency[KindSlskd] }, + "different peers were serialised", + ) + time.Sleep(100 * time.Millisecond) + + if got := p.MaxParallelGrabs(); got != kindConcurrency[KindSlskd] { + t.Errorf("%d simultaneous grabs, want the daemon cap %d", got, kindConcurrency[KindSlskd]) + } + + close(p.GrabGate) + wait() +} + +func TestSlskdLocalFolders(t *testing.T) { + t.Parallel() + + s := &slskd{downloadsPath: "/dl"} + + got := s.localFolders(Candidate{Files: []CandidateFile{ + {Path: `\m\The Wall\CD2\01 Hey You.flac`}, + {Path: `\m\The Wall\CD1\01 In The Flesh.flac`}, + {Path: `\m\The Wall\CD1\02 The Thin Ice.flac`}, + }}) + + want := []string{filepath.Join("/dl", "CD1"), filepath.Join("/dl", "CD2")} + if len(got) != len(want) || got[0] != want[0] || got[1] != want[1] { + t.Errorf("localFolders = %q, want %q", got, want) + } +} + +// Two peers' "Greatest Hits" land in one slskd directory, so the second +// grab does not enqueue until the first has collected its files. +func TestSlskdSameFolderNameWaits(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + s, downloads := newStubSlskd(t, stub) + + c := Candidate{ + Payload: map[string]string{"username": "bob"}, + Files: []CandidateFile{ + {Path: `\music\Greatest Hits\01 Intro.flac`, Size: 1, IsAudio: true}, + }, + } + + release, err := lockSlskdFolders( + context.Background(), []string{filepath.Join(downloads, "Greatest Hits")}, + ) + if err != nil { + t.Fatalf("lock: %v", err) + } + + ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond) + defer cancel() + + if _, err := s.Grab(ctx, c, t.TempDir(), nil); !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("Grab = %v, want it to wait on the held folder", err) + } + + release() + + stub.mu.Lock() + posted := stub.posted + stub.mu.Unlock() + + if posted { + t.Error("enqueued transfers into a folder another grab held") + } +} diff --git a/backend/download/provider.go b/backend/download/provider.go index f748978..b2d9965 100644 --- a/backend/download/provider.go +++ b/backend/download/provider.go @@ -236,17 +236,18 @@ func Register(d Descriptor, c Constructor) { } // concurrencyField describes the per-provider transfer limit, with help -// text explaining why the default is what it is — a user who raises -// slskd from 1 to 8 and gets themselves queued behind every other -// Soulseek user deserves to have been warned. +// text explaining what the number means where it means something +// unusual: on slskd it counts peers, since each peer is only ever asked +// for one folder at a time whatever it is set to. func concurrencyField(k Kind) Field { help := "Maximum simultaneous transfers from this client." if k == KindSlskd { - help = "Maximum simultaneous transfers. Soulseek peers serve " + - "one file at a time and queue or ban clients that ask for " + - "more, so 1 is both the polite setting and usually the " + - "fastest." + help = "How many Soulseek users to download from at once. " + + "Each user is only ever asked for one album at a time, " + + "since peers queue or ban clients that ask for more; " + + "this bounds how many different users are asked in " + + "parallel." } return Field{ diff --git a/backend/download/provider_slskd.go b/backend/download/provider_slskd.go index 02d1adb..51a6064 100644 --- a/backend/download/provider_slskd.go +++ b/backend/download/provider_slskd.go @@ -5,9 +5,12 @@ import ( "errors" "fmt" "log/slog" + "net/url" "os" "path" "path/filepath" + "regexp" + "slices" "strings" "time" @@ -69,12 +72,34 @@ const ( slskdTransferPoll = 3 * time.Second // slskdMinFiles is the fewest audio files a folder needs before it - // is offered as a candidate. Soulseek returns a lot of one-file - // noise for common queries. + // is offered as a candidate for an album. Soulseek returns a lot of + // one-file noise for common queries. A single-track request takes + // one (see minFilesFor). slskdMinFiles = 2 // slskdHTTPTimeout bounds one API call. slskdHTTPTimeout = 20 * time.Second + + // millisPerSecond converts slskd's whole-second file lengths. + millisPerSecond = 1000 + + // slskdStallAfter is how long a grab may go without a byte arriving + // before the peer is given up on. It is measured from enqueue, so + // it covers a peer that queues us and never starts as well as one + // that starts and stops. Ten minutes is long enough for a short + // queue ahead of us to clear and short enough that one unresponsive + // peer does not hold slskd's single transfer slot for an evening. + slskdStallAfter = 10 * time.Minute + + // slskdAbsentGrace is how long a requested file may be missing from + // slskd's transfer list before it is counted as failed. slskd lists + // a transfer as soon as it accepts it, so a file still absent after + // a few polls was refused. + slskdAbsentGrace = 30 * time.Second + + // slskdCancelTimeout bounds the cleanup that cancels abandoned + // transfers. + slskdCancelTimeout = 15 * time.Second ) func init() { @@ -134,6 +159,8 @@ type slskd struct { searchPoll time.Duration searchWait time.Duration transferPoll time.Duration + stallAfter time.Duration + absentGrace time.Duration } // newSlskd builds the provider from config. @@ -188,6 +215,8 @@ func newSlskd( searchPoll: slskdSearchPoll, searchWait: slskdSearchWait, transferPoll: slskdTransferPoll, + stallAfter: slskdStallAfter, + absentGrace: slskdAbsentGrace, }, nil } @@ -231,14 +260,13 @@ type slskdSearch struct { // slskdResponse is one peer's answer to a search. type slskdResponse struct { - Username string `json:"username"` - HasFreeUploadSlot bool `json:"hasFreeUploadSlot"` - QueueLength int `json:"queueLength"` - UploadSpeed int64 `json:"uploadSpeed"` - Files []slskdFile `json:"files"` - LockedFileCount int `json:"lockedFileCount"` - FileCount int `json:"fileCount"` - FreeUploadSlotFlag bool `json:"freeUploadSlots"` + Username string `json:"username"` + HasFreeUploadSlot bool `json:"hasFreeUploadSlot"` + QueueLength int `json:"queueLength"` + UploadSpeed int64 `json:"uploadSpeed"` + Files []slskdFile `json:"files"` + LockedFileCount int `json:"lockedFileCount"` + FileCount int `json:"fileCount"` } // slskdFile is one file a peer is offering. @@ -246,7 +274,9 @@ type slskdFile struct { Filename string `json:"filename"` Size int64 `json:"size"` BitRate int `json:"bitRate"` - Length int `json:"length"` + + // Length is the duration in whole seconds. + Length int `json:"length"` } // slskdTransfer is one download's state. @@ -278,24 +308,89 @@ func (t slskdTransfer) done() (finished, ok bool) { // per-folder candidates. A folder from one peer is the unit a user // actually wants: Soulseek has no album concept, but people organise // their shares by album directory. +// +// Up to two queries run at once — the request as written and a +// normalised form of it (see slskdQueries) — and their candidates are +// merged. They run concurrently rather than as a fallback because the +// manager gives a provider one search budget, and a Soulseek search +// spends most of it waiting for peers to answer; a second query after +// the first would not fit. func (s *slskd) Search(ctx context.Context, dl Download) ([]Candidate, error) { + queries := slskdQueries(dl) + if len(queries) == 0 { + return nil, nil + } + + type found struct { + candidates []Candidate + err error + } + + results := make(chan found, len(queries)) + + for _, q := range queries { + go func(q string) { + c, err := s.searchOnce(ctx, q, minFilesFor(dl)) + results <- found{candidates: c, err: err} + }(q) + } + + var ( + out []Candidate + seen = map[string]bool{} + firstErr error + answered int + ) + + for range queries { + r := <-results + if r.err != nil { + s.logger.Debug("slskd search failed", "error", r.err) + + if firstErr == nil { + firstErr = r.err + } + + continue + } + + answered++ + + // The same peer's folder turns up under both queries; the ID is + // peer and folder, so it is the same candidate. + for _, c := range r.candidates { + if seen[c.ID] { + continue + } + + seen[c.ID] = true + + out = append(out, c) + } + } + + if answered == 0 { + return nil, firstErr + } + + return out, nil +} + +// searchOnce runs one query to completion and returns its candidates. +func (s *slskd) searchOnce( + ctx context.Context, + text string, + minFiles int, +) ([]Candidate, error) { // slskd's search endpoint deserializes id as a .NET Guid server-side, // so it must be a dashed UUID — the app's own newID() (a plain hex // string, used for request/item IDs elsewhere) is rejected with an // HTTP 400 before any search happens. searchID := uuid.NewString() - body := map[string]any{ - "id": searchID, - "searchText": dl.SearchText(), - } - - if err := s.client.post(ctx, "/api/v0/searches", body, nil); err != nil { - return nil, err - } - - search, err := s.awaitSearch(ctx, searchID) - if err != nil { + if err := s.client.post( + ctx, "/api/v0/searches", s.searchRequest(searchID, text, minFiles), nil, + ); err != nil { return nil, err } @@ -307,52 +402,222 @@ func (s *slskd) Search(ctx context.Context, dl Download) ([]Candidate, error) { ) }() - return s.candidatesFrom(search), nil + if err := s.awaitSearch(ctx, searchID); err != nil { + return nil, err + } + + responses, err := s.searchResponses(ctx, searchID) + if err != nil { + return nil, err + } + + return s.candidatesFrom(responses, minFiles), nil +} + +// searchRequest is the body that starts a search. +// +// Every option is stated rather than left to the daemon, because +// slskd's defaults are its own and not ours. Its search timeout in +// particular has to finish inside our wait: a search that slskd is still +// running when we stop polling is results we asked for and discarded. +// The response and file limits are raised well above what a popular +// album produces, and the peer filters let slskd drop answers this +// provider would only score down to nothing — a folder too small to be +// a candidate, a peer with a queue it will not reach today. +func (s *slskd) searchRequest(id, text string, minFiles int) map[string]any { + const ( + responseLimit = 500 + fileLimit = 20_000 + maximumPeerQueueLength = 100 + ) + + // A tenth of the wait is left for the last poll and the responses + // fetch. + timeout := s.searchWait - s.searchWait/10 + + return map[string]any{ + "id": id, + "searchText": text, + "searchTimeout": timeout.Milliseconds(), + "responseLimit": responseLimit, + "fileLimit": fileLimit, + "filterResponses": true, + "minimumResponseFileCount": minFiles, + "maximumPeerQueueLength": maximumPeerQueueLength, + } +} + +// slskdQueries is what is searched for a request: the request's own +// search text, and a normalised form of it when that differs. +// +// Soulseek matches every term against the file's full path, so each +// extra word is a filter, and some words filter wrongly: +// +// - edition qualifiers — "(Deluxe Edition)", "[2011 Remaster]" — are +// in the catalog's title and rarely in anyone's folder name; +// - punctuation splits a term oddly, and a term that starts with "-" +// is an *exclusion*, so an album called "-ism" searches for +// everything without it; +// - "Various Artists" is in no one's path for a compilation. +// +// A query the user typed is theirs and is searched exactly as written. +func slskdQueries(dl Download) []string { + primary := strings.TrimSpace(dl.SearchText()) + if primary == "" { + return nil + } + + out := []string{primary} + + if dl.Query != "" { + return out + } + + artist := dl.Artist + if isVariousArtists(artist) { + artist = "" + } + + normal := Download{ + Artist: normalizeSearchTerms(artist), + Album: normalizeSearchTerms(editionPattern.ReplaceAllString(dl.Album, " ")), + } + + if alt := strings.TrimSpace(normal.SearchText()); alt != "" && + !strings.EqualFold(alt, primary) { + out = append(out, alt) + } + + return out +} + +var ( + // editionPattern finds an edition qualifier: a bracketed group that + // names an edition, or a trailing " - 2011 Remaster". + editionPattern = regexp.MustCompile( + `(?i)\s*[(\[][^)\]]*\b(?:deluxe|edition|remaster(?:ed)?|expanded|` + + `anniversary|bonus|explicit|reissue|special|collector'?s?|` + + `version|mono|stereo)\b[^)\]]*[)\]]` + + `|\s+-\s+(?:\d{4}\s+)?remaster(?:ed)?\b.*$`, + ) + + // nonWordPattern is everything that is not a letter or a digit. + nonWordPattern = regexp.MustCompile(`[^\p{L}\p{N}]+`) +) + +// normalizeSearchTerms reduces text to plain words. +func normalizeSearchTerms(s string) string { + return strings.Join(strings.Fields(nonWordPattern.ReplaceAllString(s, " ")), " ") +} + +// isVariousArtists reports whether an artist credit is a compilation's +// placeholder rather than an artist. +func isVariousArtists(artist string) bool { + switch strings.ToLower(strings.TrimSpace(artist)) { + case "various artists", "various", "va": + return true + default: + return false + } +} + +// minFilesFor is the fewest audio files a folder must offer to be a +// candidate for this request. +// +// Soulseek answers a search with the files that match it, not with the +// folders they sit in. An album query matches every file in the album's +// folder, because the folder name carries the terms; a *track* query +// usually matches one file per folder. The two-file floor that filters +// out one-file noise for an album therefore filtered out every result +// for a track, and a single-track request could never be served here. +func minFilesFor(dl Download) int { + if dl.RecordingMBID != "" { + return 1 + } + + return slskdMinFiles } // awaitSearch polls until the search completes or the budget runs out. // A timeout is not an error: partial Soulseek results are normal and // often good enough. -func (s *slskd) awaitSearch( - ctx context.Context, - searchID string, -) (slskdSearch, error) { +// +// The poll asks for the search's state only. It used to ask for every +// response on every one-second tick, which for a popular album is the +// same few thousand file entries serialised twenty times to be read +// once; searchResponses fetches them once at the end. +func (s *slskd) awaitSearch(ctx context.Context, searchID string) error { deadline := time.Now().Add(s.searchWait) - var last slskdSearch - for time.Now().Before(deadline) { select { case <-ctx.Done(): - return last, fmt.Errorf("%w: search cancelled", ErrSlskdTimeout) + return fmt.Errorf("%w: search cancelled", ErrSlskdTimeout) case <-time.After(s.searchPoll): } var search slskdSearch if err := s.client.get( - ctx, - "/api/v0/searches/"+searchID+"?includeResponses=true", - &search, + ctx, "/api/v0/searches/"+searchID, &search, ); err != nil { - return last, err + return err } - last = search - if search.IsComplete { - return search, nil + return nil } } - return last, nil + return nil } -// candidatesFrom groups a search's responses into candidates. -func (s *slskd) candidatesFrom(search slskdSearch) []Candidate { - out := make([]Candidate, 0, len(search.Responses)) +// searchResponses fetches a search's responses once. +// +// `/searches/{id}/responses` is the endpoint for that; a daemon that +// does not answer it is asked the older way, with the search itself +// carrying its responses, so an older slskd degrades to the previous +// behaviour rather than to no results at all. +func (s *slskd) searchResponses( + ctx context.Context, + searchID string, +) ([]slskdResponse, error) { + var responses []slskdResponse - for _, resp := range search.Responses { + err := s.client.get( + ctx, "/api/v0/searches/"+searchID+"/responses", &responses, + ) + if err == nil { + return responses, nil + } + + s.logger.Debug( + "slskd responses endpoint failed; asking with the search", + "error", err, + ) + + var search slskdSearch + + if err := s.client.get( + ctx, + "/api/v0/searches/"+searchID+"?includeResponses=true", + &search, + ); err != nil { + return nil, err + } + + return search.Responses, nil +} + +// candidatesFrom groups a search's responses into candidates, dropping +// folders with fewer than minFiles audio files. +func (s *slskd) candidatesFrom( + responses []slskdResponse, + minFiles int, +) []Candidate { + out := make([]Candidate, 0, len(responses)) + + for _, resp := range responses { for folder, files := range groupByFolder(resp.Files) { audio := 0 @@ -372,12 +637,14 @@ func (s *slskd) candidatesFrom(search slskdSearch) []Candidate { Format: format, Bitrate: f.BitRate, IsAudio: isAudio, + + LengthMillis: int64(f.Length) * millisPerSecond, }) total += f.Size } - if audio < slskdMinFiles { + if audio < minFiles { continue } @@ -398,13 +665,15 @@ func (s *slskd) candidatesFrom(search slskdSearch) []Candidate { return out } -// groupByFolder buckets a peer's files by their containing directory. +// groupByFolder buckets a peer's files by the album directory they sit +// in — the containing directory, or the one above it for a disc folder +// (see AlbumDir), so a multi-disc rip is one candidate and not two. func groupByFolder(files []slskdFile) map[string][]slskdFile { out := map[string][]slskdFile{} for _, f := range files { - norm := strings.ReplaceAll(f.Filename, `\`, "/") - out[path.Dir(norm)] = append(out[path.Dir(norm)], f) + dir := AlbumDir(f.Filename) + out[dir] = append(out[dir], f) } return out @@ -418,7 +687,7 @@ func groupByFolder(files []slskdFile) map[string][]slskdFile { func peerHealth(r slskdResponse) float64 { score := 0.35 - if r.HasFreeUploadSlot || r.FreeUploadSlotFlag { + if r.HasFreeUploadSlot { score += 0.4 } @@ -467,6 +736,21 @@ func (s *slskd) Grab( ) } + // slskd keeps finished transfers listed until someone removes them, + // and a transfer is matched to the request by filename. A record + // left by an earlier attempt at the same file from the same peer + // would otherwise be read as this attempt's answer the moment the + // first poll came back — an old failure failing a transfer that has + // not started. So what is already terminal is noted before enqueueing + // and ignored after. + release, err := lockSlskdFolders(ctx, s.localFolders(c)) + if err != nil { + return Result{}, err + } + defer release() + + stale := s.terminalTransferIDs(ctx, username) + wanted := make([]map[string]any, 0, len(c.Files)) for _, f := range c.Files { wanted = append(wanted, map[string]any{ @@ -476,24 +760,122 @@ func (s *slskd) Grab( } if err := s.client.post( - ctx, "/api/v0/transfers/downloads/"+username, wanted, nil, + ctx, slskdDownloadsPath(username), wanted, nil, ); err != nil { return Result{}, err } - if err := s.awaitTransfers(ctx, username, c, onProgress); err != nil { + if err := s.awaitTransfers( + ctx, username, stale, c, onProgress, + ); err != nil { return Result{}, err } return s.collect(c, dst) } +// slskdFolders serialises grabs that land in the same local folder. +// +// slskd names a download's directory after the remote *leaf* folder, so +// two different albums both shared as "Greatest Hits" — or any two +// multi-disc rips, whose leaves are "CD1" and "CD2" — are written into +// one directory, and collect finds files by name there. Run at once, +// a file one peer never sent is filled by the other peer's file of the +// same name. One grab per peer made that impossible; several peers at +// once makes it likely. It is package-level and keyed on the full +// path because two configured clients can share one daemon. +var slskdFolders keyedLock[string] + +// localFolders returns the directories under downloadsPath a candidate's +// files will be written to, sorted so every grab takes them in the same +// order and two cannot each hold what the other waits for. +func (s *slskd) localFolders(c Candidate) []string { + var out []string + + for _, f := range c.Files { + norm := strings.ReplaceAll(f.Path, `\`, "/") + out = append(out, filepath.Join(s.downloadsPath, path.Base(path.Dir(norm)))) + } + + slices.Sort(out) + + return slices.Compact(out) +} + +// lockSlskdFolders takes every folder in order, releasing what it holds +// if the context ends part way. +func lockSlskdFolders(ctx context.Context, folders []string) (func(), error) { + releases := make([]func(), 0, len(folders)) + + releaseAll := func() { + for _, r := range slices.Backward(releases) { + r() + } + } + + for _, f := range folders { + r, err := slskdFolders.acquire(ctx, f) + if err != nil { + releaseAll() + + return nil, err + } + + releases = append(releases, r) + } + + return releaseAll, nil +} + +// slskdDownloadsPath is the transfers endpoint for one peer. Soulseek +// usernames may contain spaces and punctuation, so the name is escaped +// rather than spliced into the path. +func slskdDownloadsPath(username string) string { + return "/api/v0/transfers/downloads/" + url.PathEscape(username) +} + +// terminalTransferIDs returns the ids of this peer's transfers that are +// already finished. Best effort: slskd answers 404 for a peer it has no +// transfers with, and any failure here means only that there is nothing +// to ignore. +func (s *slskd) terminalTransferIDs( + ctx context.Context, + username string, +) map[string]bool { + transfers, err := s.transfersFor(ctx, username) + if err != nil { + return nil + } + + out := make(map[string]bool, len(transfers)) + + for _, t := range transfers { + if finished, _ := t.done(); finished && t.ID != "" { + out[t.ID] = true + } + } + + return out +} + // awaitTransfers polls until every requested file reaches a terminal -// state. Soulseek queues are measured in hours, so the only deadline -// is the caller's context. +// state, the transfer stalls, or the caller gives up. +// +// Soulseek queues are measured in hours, so there is no deadline on the +// transfer as a whole — but there is one on *progress*. slskd's +// transfer limit is one, so a peer that holds us in its queue without +// sending a byte is not only failing this download, it is holding every +// other Soulseek download behind it. After stallAfter with nothing +// moving the peer is given up on, and the manager tries another. +// +// Whatever way this ends short of every file finishing, the transfers +// still live in slskd are cancelled there. Returning without doing so +// leaves the daemon downloading into its own folder for a request +// nobody is waiting on any more. func (s *slskd) awaitTransfers( ctx context.Context, username string, + stale map[string]bool, c Candidate, onProgress ProgressFunc, ) error { @@ -502,9 +884,18 @@ func (s *slskd) awaitTransfers( wanted[f.Path] = true } + var ( + started = time.Now() + lastProgress = started + lastBytes int64 + live []slskdTransfer + ) + for { select { case <-ctx.Done(): + s.cancelTransfers(username, live) + return fmt.Errorf("%w: transfer cancelled", ErrSlskdTimeout) case <-time.After(s.transferPoll): } @@ -512,60 +903,177 @@ func (s *slskd) awaitTransfers( transfers, err := s.transfersFor(ctx, username) if err != nil { // A blip talking to the daemon should not abandon a - // transfer that may be hours in. + // transfer that may be hours in — but a daemon that stays + // away is a stall like any other. s.logger.Debug("slskd transfer poll failed", "error", err) + if time.Since(lastProgress) >= s.stallAfter { + s.cancelTransfers(username, live) + + return fmt.Errorf( + "%w: slskd has not answered for %s: %w", + ErrSlskdTimeout, s.stallAfter, err, + ) + } + continue } - var ( - done, failed int - current int64 + tally := tallyTransfers( + transfers, wanted, stale, + time.Since(started) >= s.absentGrace, ) + live = tally.live - for _, t := range transfers { - if !wanted[t.Filename] { - continue - } - - current += t.BytesTransferred - - finished, ok := t.done() - if !finished { - continue - } - - if ok { - done++ - } else { - failed++ - } + if tally.bytes > lastBytes { + lastBytes = tally.bytes + lastProgress = time.Now() } if onProgress != nil { onProgress(Progress{ - Current: current, + Current: tally.bytes, Total: c.TotalSize, Phase: fmt.Sprintf( - "Transferring from %s (%d/%d)", username, done, len(wanted), + "Transferring from %s (%d/%d)", + username, tally.done, len(wanted), ), }) } - if done+failed < len(wanted) { + if tally.done+tally.failed >= len(wanted) { + // Some files failing is normal — a peer goes offline + // mid-folder. Let the importer's completeness check decide + // whether what arrived is enough, rather than discarding it + // here. + if tally.done == 0 { + return fmt.Errorf( + "%w: all %d files failed", + ErrSlskdTransferFailed, tally.failed, + ) + } + + return nil + } + + if time.Since(lastProgress) < s.stallAfter { continue } - // Some files failing is normal — a peer goes offline mid-folder. - // Let the importer's completeness check decide whether what - // arrived is enough, rather than discarding it here. - if done == 0 { - return fmt.Errorf( - "%w: all %d files failed", ErrSlskdTransferFailed, failed, + s.cancelTransfers(username, live) + + // A folder that stalls on its last track is the same shape as + // one whose last track failed, and goes forward the same way. + if tally.done > 0 { + s.logger.Info( + "slskd transfer stalled; keeping what arrived", + "peer", username, + "done", tally.done, + "wanted", len(wanted), ) + + return nil } - return nil + return fmt.Errorf( + "%w: %s sent nothing in %s", + ErrSlskdTimeout, username, s.stallAfter, + ) + } +} + +// transferTally is one poll's reading of the files a grab asked for. +type transferTally struct { + done, failed int + bytes int64 + + // live are the requested transfers slskd is still working on, + // which are what has to be cancelled if the grab is abandoned. + live []slskdTransfer +} + +// tallyTransfers reads a peer's transfer list against the files a grab +// asked for. +// +// A requested file slskd does not list at all is one it never accepted +// — refused at enqueue, or dropped — and it will never reach a terminal +// state to be counted by. Once absentExpired, such a file counts as +// failed, or the grab would wait on it until the six-hour ceiling. +func tallyTransfers( + transfers []slskdTransfer, + wanted map[string]bool, + stale map[string]bool, + absentExpired bool, +) transferTally { + seen := make(map[string]slskdTransfer, len(wanted)) + + for _, t := range transfers { + if !wanted[t.Filename] || stale[t.ID] { + continue + } + + seen[t.Filename] = t + } + + var out transferTally + + for name := range wanted { + t, ok := seen[name] + if !ok { + if absentExpired { + out.failed++ + } + + continue + } + + out.bytes += t.BytesTransferred + + finished, succeeded := t.done() + + switch { + case !finished: + out.live = append(out.live, t) + case succeeded: + out.done++ + default: + out.failed++ + } + } + + return out +} + +// cancelTransfers asks slskd to cancel and forget transfers this grab +// is abandoning. It runs on a context of its own: the usual reason to +// be here is that the caller's context has just been cancelled, and a +// cleanup that inherited it would never be sent. +func (s *slskd) cancelTransfers(username string, live []slskdTransfer) { + if len(live) == 0 { + return + } + + ctx, cancel := context.WithTimeout( + context.Background(), slskdCancelTimeout, + ) + defer cancel() + + for _, t := range live { + if t.ID == "" { + continue + } + + endpoint := slskdDownloadsPath(username) + "/" + + url.PathEscape(t.ID) + "?remove=true" + + if err := s.client.delete(ctx, endpoint); err != nil { + s.logger.Warn( + "could not cancel slskd transfer", + "peer", username, + "file", t.Filename, + "error", err, + ) + } } } @@ -581,9 +1089,7 @@ func (s *slskd) transfersFor( } `json:"directories"` } - if err := s.client.get( - ctx, "/api/v0/transfers/downloads/"+username, &raw, - ); err != nil { + if err := s.client.get(ctx, slskdDownloadsPath(username), &raw); err != nil { return nil, err } @@ -616,7 +1122,13 @@ func (s *slskd) collect(c Candidate, dst string) (Result, error) { continue } + // A multi-disc candidate keeps its disc folders in staging. + // Flattened, disc 2's "01 Intro.flac" overwrites disc 1's, and + // the importer loses the folder it reads the disc number from. target := filepath.Join(dst, base) + if _, ok := discFolder(folder); ok { + target = filepath.Join(dst, folder, base) + } if err := movePath(src, target); err != nil { return Result{}, fmt.Errorf("collect %s: %w", base, err) diff --git a/backend/download/provider_slskd_test.go b/backend/download/provider_slskd_test.go index 992e6d5..d06cec5 100644 --- a/backend/download/provider_slskd_test.go +++ b/backend/download/provider_slskd_test.go @@ -31,11 +31,30 @@ type slskdStub struct { transfers [][]slskdTransfer pollCount int + // before is what the downloads endpoint reports until something is + // enqueued: records slskd already held from earlier attempts. + before []slskdTransfer + // enqueued records what was requested for download. enqueued []map[string]any + posted bool + + // paths records the escaped path of every transfers call, and + // cancelled the escaped request URI of every DELETE. + paths []string + cancelled []string // unauthorized makes every call return 401. unauthorized bool + + // searches records every search request body, and searchGets the + // request URI of every search GET. + searches []map[string]any + searchGets []string + + // noResponsesEndpoint makes /searches/{id}/responses 404, as an + // older daemon would. + noResponsesEndpoint bool } func newSlskdStub(t *testing.T) *slskdStub { @@ -57,6 +76,16 @@ func newSlskdStub(t *testing.T) *slskdStub { return } + var body map[string]any + + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + t.Errorf("decode search body: %v", err) + } + + s.mu.Lock() + s.searches = append(s.searches, body) + s.mu.Unlock() + w.WriteHeader(http.StatusCreated) }) @@ -73,8 +102,26 @@ func newSlskdStub(t *testing.T) *slskdStub { s.mu.Lock() responses := s.responses + noEndpoint := s.noResponsesEndpoint + s.searchGets = append(s.searchGets, r.URL.RequestURI()) s.mu.Unlock() + if strings.HasSuffix(r.URL.Path, "/responses") { + if noEndpoint { + w.WriteHeader(http.StatusNotFound) + + return + } + + writeJSON(t, w, responses) + + return + } + + if r.URL.Query().Get("includeResponses") != "true" { + responses = nil + } + writeJSON(t, w, slskdSearch{ ID: "search-1", IsComplete: true, @@ -87,7 +134,12 @@ func newSlskdStub(t *testing.T) *slskdStub { return } - if r.Method == http.MethodPost { + s.mu.Lock() + s.paths = append(s.paths, r.URL.EscapedPath()) + s.mu.Unlock() + + switch r.Method { + case http.MethodPost: var body []map[string]any if err := json.NewDecoder(r.Body).Decode(&body); err != nil { @@ -96,15 +148,35 @@ func newSlskdStub(t *testing.T) *slskdStub { s.mu.Lock() s.enqueued = body + s.posted = true s.mu.Unlock() w.WriteHeader(http.StatusCreated) + return + case http.MethodDelete: + s.mu.Lock() + s.cancelled = append(s.cancelled, r.URL.RequestURI()) + s.mu.Unlock() + + w.WriteHeader(http.StatusNoContent) + return } s.mu.Lock() + if !s.posted { + before := s.before + s.mu.Unlock() + + writeJSON(t, w, map[string]any{ + "directories": []map[string]any{{"files": before}}, + }) + + return + } + idx := s.pollCount if idx >= len(s.transfers) { idx = len(s.transfers) - 1 @@ -191,6 +263,11 @@ func newStubSlskd(t *testing.T, stub *slskdStub) (*slskd, string) { s.searchWait = 200 * time.Millisecond s.transferPoll = time.Millisecond + // Long enough that no existing test trips them by accident; the + // tests about stalls and absences set their own. + s.stallAfter = time.Minute + s.absentGrace = time.Minute + return s, downloads } @@ -565,3 +642,293 @@ func TestSlskdRequiresConfiguration(t *testing.T) { }) } } + +// slskdAlbum is a two-file candidate from peer, with the files slskd +// would have written already in place under downloads. +func slskdAlbum(t *testing.T, downloads, peer string, arrived ...string) Candidate { + t.Helper() + + folder := filepath.Join(downloads, "Album") + if err := os.MkdirAll(folder, 0o750); err != nil { + t.Fatalf("mkdir: %v", err) + } + + for _, name := range arrived { + if err := os.WriteFile( + filepath.Join(folder, name), []byte("audio"), 0o600, + ); err != nil { + t.Fatalf("write: %v", err) + } + } + + return Candidate{ + Files: []CandidateFile{ + {Path: `\s\Album\01 A.flac`, Size: 500, IsAudio: true}, + {Path: `\s\Album\02 B.flac`, Size: 500, IsAudio: true}, + }, + TotalSize: 1000, + Payload: map[string]string{"username": peer}, + } +} + +// cancelledURIs returns what the stub was asked to cancel. +func (s *slskdStub) cancelledURIs() []string { + s.mu.Lock() + defer s.mu.Unlock() + + return append([]string(nil), s.cancelled...) +} + +// A peer that queues us and never sends a byte is given up on, and the +// queued transfers are cancelled in slskd rather than left to start +// hours later for a request nobody is waiting on. +func TestSlskdGrabGivesUpOnAStalledPeer(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + stub.transfers = [][]slskdTransfer{{ + {ID: "t1", Filename: `\s\Album\01 A.flac`, State: "Queued, Remotely"}, + {ID: "t2", Filename: `\s\Album\02 B.flac`, State: "Queued, Remotely"}, + }} + + s, downloads := newStubSlskd(t, stub) + s.stallAfter = 30 * time.Millisecond + + _, err := s.Grab( + context.Background(), slskdAlbum(t, downloads, "peer"), t.TempDir(), nil, + ) + if !errors.Is(err, ErrSlskdTimeout) { + t.Fatalf("error = %v, want ErrSlskdTimeout", err) + } + + got := stub.cancelledURIs() + if len(got) != 2 { + t.Fatalf("cancelled %v, want both queued transfers", got) + } + + for _, uri := range got { + if !strings.HasSuffix(uri, "?remove=true") { + t.Errorf("cancel %s does not remove the record", uri) + } + } +} + +// A folder that stalls on its last track goes forward with what +// arrived, the same as one whose last track failed; the importer's +// completeness check decides whether that is enough. +func TestSlskdGrabKeepsWhatArrivedBeforeAStall(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + stub.transfers = [][]slskdTransfer{{ + { + ID: "t1", Filename: `\s\Album\01 A.flac`, + State: "Completed, Succeeded", BytesTransferred: 500, + }, + {ID: "t2", Filename: `\s\Album\02 B.flac`, State: "Queued, Remotely"}, + }} + + s, downloads := newStubSlskd(t, stub) + s.stallAfter = 30 * time.Millisecond + + got, err := s.Grab( + context.Background(), + slskdAlbum(t, downloads, "peer", "01 A.flac"), + t.TempDir(), nil, + ) + if err != nil { + t.Fatalf("Grab: %v", err) + } + + if len(got.Files) != 1 { + t.Errorf("collected %d files, want the 1 that arrived", len(got.Files)) + } + + if cancelled := stub.cancelledURIs(); len(cancelled) != 1 || + !strings.Contains(cancelled[0], "/t2") { + t.Errorf("cancelled %v, want only the stalled t2", cancelled) + } +} + +// Progress is what holds the stall timer off. A transfer that keeps +// moving bytes is never abandoned, however long it takes. +func TestSlskdGrabWaitsOnATransferThatIsMoving(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + + for b := int64(1); b <= 100; b++ { + stub.transfers = append(stub.transfers, []slskdTransfer{ + {ID: "t1", Filename: `\s\Album\01 A.flac`, State: "InProgress", BytesTransferred: b}, + {ID: "t2", Filename: `\s\Album\02 B.flac`, State: "Queued, Remotely"}, + }) + } + + stub.transfers = append(stub.transfers, []slskdTransfer{ + { + ID: "t1", + Filename: `\s\Album\01 A.flac`, + State: "Completed, Succeeded", + BytesTransferred: 500, + }, + { + ID: "t2", + Filename: `\s\Album\02 B.flac`, + State: "Completed, Succeeded", + BytesTransferred: 500, + }, + }) + + s, downloads := newStubSlskd(t, stub) + // A hundred polls take several times the stall window; each one + // moves a byte. The window is kept well above one poll so a + // descheduled test runner does not read as a stall. + s.transferPoll = 5 * time.Millisecond + s.stallAfter = 150 * time.Millisecond + + got, err := s.Grab( + context.Background(), + slskdAlbum(t, downloads, "peer", "01 A.flac", "02 B.flac"), + t.TempDir(), nil, + ) + if err != nil { + t.Fatalf("Grab: %v", err) + } + + if len(got.Files) != 2 { + t.Errorf("collected %d files, want 2", len(got.Files)) + } +} + +// A file slskd never lists was refused at enqueue and will never reach +// a terminal state. It counts as failed once the grace period is up, +// rather than being waited on until the six-hour ceiling. +func TestSlskdGrabCountsAnUnlistedFileAsFailed(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + stub.transfers = [][]slskdTransfer{{ + { + ID: "t1", Filename: `\s\Album\01 A.flac`, + State: "Completed, Succeeded", BytesTransferred: 500, + }, + }} + + s, downloads := newStubSlskd(t, stub) + s.absentGrace = 20 * time.Millisecond + + got, err := s.Grab( + context.Background(), + slskdAlbum(t, downloads, "peer", "01 A.flac"), + t.TempDir(), nil, + ) + if err != nil { + t.Fatalf("Grab: %v", err) + } + + if len(got.Files) != 1 { + t.Errorf("collected %d files, want 1", len(got.Files)) + } +} + +// A finished record left by an earlier attempt at the same file is not +// this attempt's answer. Without the snapshot it would fail the grab on +// the first poll, before the new transfer had started. +func TestSlskdGrabIgnoresAnEarlierAttemptsRecord(t *testing.T) { + t.Parallel() + + stale := slskdTransfer{ + ID: "old", Filename: `\s\Album\01 A.flac`, State: "Completed, Errored", + } + + stub := newSlskdStub(t) + stub.before = []slskdTransfer{stale} + stub.transfers = [][]slskdTransfer{ + {stale}, + { + stale, + { + ID: "new", Filename: `\s\Album\01 A.flac`, + State: "Completed, Succeeded", BytesTransferred: 500, + }, + }, + } + + s, downloads := newStubSlskd(t, stub) + + c := slskdAlbum(t, downloads, "peer", "01 A.flac") + c.Files = c.Files[:1] + + got, err := s.Grab(context.Background(), c, t.TempDir(), nil) + if err != nil { + t.Fatalf("Grab: %v", err) + } + + if len(got.Files) != 1 { + t.Errorf("collected %d files, want 1", len(got.Files)) + } +} + +// Cancelling the download cancels the transfer in slskd too. The +// cleanup must not inherit the cancelled context, or it is never sent. +func TestSlskdGrabCancelsTransfersWhenTheCallerGivesUp(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + stub.transfers = [][]slskdTransfer{{ + {ID: "t1", Filename: `\s\Album\01 A.flac`, State: "InProgress", BytesTransferred: 10}, + {ID: "t2", Filename: `\s\Album\02 B.flac`, State: "Queued, Remotely"}, + }} + + s, downloads := newStubSlskd(t, stub) + + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + + _, err := s.Grab(ctx, slskdAlbum(t, downloads, "peer"), t.TempDir(), nil) + if !errors.Is(err, ErrSlskdTimeout) { + t.Fatalf("error = %v, want ErrSlskdTimeout", err) + } + + if got := stub.cancelledURIs(); len(got) != 2 { + t.Errorf("cancelled %v, want both live transfers", got) + } +} + +// Soulseek usernames carry spaces and punctuation; spliced raw into the +// path, a name with a slash addresses a different endpoint entirely. +func TestSlskdEscapesTheUsername(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + stub.transfers = [][]slskdTransfer{{ + { + ID: "t1", Filename: `\s\Album\01 A.flac`, + State: "Completed, Succeeded", BytesTransferred: 500, + }, + { + ID: "t2", Filename: `\s\Album\02 B.flac`, + State: "Completed, Succeeded", BytesTransferred: 500, + }, + }} + + s, downloads := newStubSlskd(t, stub) + + if _, err := s.Grab( + context.Background(), + slskdAlbum(t, downloads, "dj a/b", "01 A.flac", "02 B.flac"), + t.TempDir(), nil, + ); err != nil { + t.Fatalf("Grab: %v", err) + } + + stub.mu.Lock() + paths := append([]string(nil), stub.paths...) + stub.mu.Unlock() + + for _, p := range paths { + if p != "/api/v0/transfers/downloads/dj%20a%2Fb" { + t.Errorf("transfers call went to %s", p) + } + } +} diff --git a/backend/download/rank.go b/backend/download/rank.go index ecaa53c..b50abb1 100644 --- a/backend/download/rank.go +++ b/backend/download/rank.go @@ -35,6 +35,15 @@ const ( weightArtistFit = 0.12 ) +// Match sub-weights when the candidate's durations are known. Duration +// takes its weight from title fit, the signal it corroborates: a title +// says which song a file claims to be, a length says whether it is that +// recording — the right edit, the whole file, not the live take. +const ( + timedWeightTitleFit = 0.25 + timedWeightDurationFit = 0.15 +) + // Quality sub-weights. Each set sums to 1.0. // // There are two of them because a stated preference changes what the @@ -319,13 +328,13 @@ func Score(dl Download, c Candidate, priority int, prefs AutoDownloadPrefs) Cand audio := c.AudioFiles() - matched, titleFit := matchFiles(audio, dl.Expected) + a := alignFiles(audio, dl.Expected) // Write the alignment back so the picker can show which file maps // to which track. - c.Files = mergeMatched(c.Files, matched) + c.Files = mergeMatched(c.Files, a.files) - c.Match = scoreMatch(dl, c, audio, titleFit) + c.Match = scoreMatch(dl, c, audio, a) c.Quality = scoreQuality( c, audio, priority, prefs, dl.runtimeMillis(), ) @@ -340,14 +349,21 @@ func scoreMatch( dl Download, c Candidate, audio []CandidateFile, - titleFit float64, + a alignment, ) MatchScore { m := MatchScore{ - Anchored: dl.Anchored(), - TitleFit: titleFit, + Anchored: dl.Anchored(), + TitleFit: a.titleFit, + DurationFit: a.durationFit, + + // Durations count once at least half the aligned pairs state + // one; a single timed pair would be a coin toss carrying 15%. + DurationKnown: a.timedPairs > 0 && a.timedPairs*2 >= a.aligned, } - m.Completeness = completeness(len(audio), len(dl.Expected)) + m.Completeness = completeness( + alignedCount(c.Files), len(audio), len(dl.Expected), + ) // The candidate's own title, and the folder its files sit in, are // two independent guesses at the album name. Take the better one: @@ -367,9 +383,16 @@ func scoreMatch( // With no expected tracklist there is no title signal at all, so // redistribute its weight onto the album/artist evidence rather // than scoring every free-text result as half-wrong. - if len(dl.Expected) == 0 { + switch { + case len(dl.Expected) == 0: m.Overall = 0.55*m.AlbumFit + 0.45*m.ArtistFit - } else { + case m.DurationKnown: + m.Overall = timedWeightTitleFit*m.TitleFit + + timedWeightDurationFit*m.DurationFit + + weightCompleteness*m.Completeness + + weightAlbumFit*m.AlbumFit + + weightArtistFit*m.ArtistFit + default: m.Overall = weightTitleFit*m.TitleFit + weightCompleteness*m.Completeness + weightAlbumFit*m.AlbumFit + @@ -415,30 +438,54 @@ func artistFit(want string, c Candidate) float64 { return best } -// completeness scores audio file count against the expected track -// count. Extra files are penalized far more gently than missing ones: +// completeness scores how much of the expected tracklist a candidate +// covers. Extra files are penalized far more gently than missing ones: // a folder with bonus tracks or a stray intro is still the album, while // a folder missing half the tracks is not. -func completeness(got, want int) float64 { +// +// **Coverage is counted in aligned tracks, not in files.** It used to +// be the audio file count, so any ten files scored full marks against +// a ten-track album whether or not they were its tracks — and since +// title fit is the mean over the files that *did* align, a folder where +// three titles matched read as a near-perfect candidate on both counts. +// `aligned` is how many files matchFiles assigned to an expected track; +// `audio` still sets the penalty for extras, because a folder of thirty +// files holding the ten wanted is a worse copy than one holding ten. +func completeness(aligned, audio, want int) float64 { if want == 0 { - if got > 0 { + if audio > 0 { return 0.5 } return 0 } - if got == 0 { + if aligned == 0 { return 0 } - if got >= want { - extra := float64(got-want) / float64(want) + cover := float64(min(aligned, want)) / float64(want) - return math.Max(0.75, 1.0-0.25*extra) + if audio > want { + extra := float64(audio-want) / float64(want) + cover *= math.Max(0.75, 1.0-0.25*extra) } - return float64(got) / float64(want) + return cover +} + +// alignedCount is how many audio files were assigned to an expected +// track. +func alignedCount(files []CandidateFile) int { + n := 0 + + for _, f := range files { + if f.IsAudio && f.MatchedTo != 0 { + n++ + } + } + + return n } // scoreQuality answers whether this is a good copy. @@ -735,6 +782,45 @@ func AutoPickVeto( return "" } +// autoAcceptable reports whether auto-pick may take this one candidate +// without asking: the request is anchored to a tracklist, and the +// candidate is inside the user's guardrails and clears the match and +// quality bars. It is AutoPickVeto's test applied to a single +// candidate, which is what falling back to a second choice needs. +func autoAcceptable(dl Download, c Candidate, prefs AutoDownloadPrefs) bool { + return dl.Anchored() && + len(dl.Expected) > 0 && + prefs.eligible(c, dl.runtimeMillis()) && + c.Match.Overall >= minMatch && + c.Quality.Overall >= minQuality +} + +// autoPick returns the candidate auto-pick takes: the best-ranked one +// it may take at all. +// +// That is not `ranked[0]`. AutoPickVeto judges the best candidate +// *inside* the guardrails, so when the overall best is outside them — +// over the size ceiling, say — the veto passes on the strength of the +// second, and grabbing the first would download exactly the copy the +// user said not to take unattended. +func autoPick( + dl Download, + ranked []Candidate, + prefs AutoDownloadPrefs, +) (Candidate, bool) { + if AutoPickVeto(dl, ranked, prefs) != "" { + return Candidate{}, false + } + + for _, c := range ranked { + if autoAcceptable(dl, c, prefs) { + return c, true + } + } + + return Candidate{}, false +} + // mergeMatched copies MatchedTo assignments from the audio-only slice // back onto the full file list. func mergeMatched(all, matched []CandidateFile) []CandidateFile { diff --git a/backend/download/rank_test.go b/backend/download/rank_test.go index 9c0a2cf..4aeb93f 100644 --- a/backend/download/rank_test.go +++ b/backend/download/rank_test.go @@ -282,28 +282,32 @@ func TestCompleteness(t *testing.T) { tests := []struct { name string - got int + aligned int + audio int want int minScore float64 maxScore float64 }{ - {"exact", 10, 10, 1.0, 1.0}, - {"half missing", 5, 10, 0.49, 0.51}, - {"one bonus track", 11, 10, 0.95, 1.0}, - {"double", 20, 10, 0.74, 0.76}, - {"nothing", 0, 10, 0, 0}, - {"no expectation", 5, 0, 0.5, 0.5}, + {"exact", 10, 10, 10, 1.0, 1.0}, + {"half missing", 5, 5, 10, 0.49, 0.51}, + {"one bonus track", 10, 11, 10, 0.95, 1.0}, + {"double", 10, 20, 10, 0.74, 0.76}, + {"nothing", 0, 0, 10, 0, 0}, + {"no expectation", 0, 5, 0, 0.5, 0.5}, + // Ten files are not ten tracks: three that align are three. + {"right count, wrong tracks", 3, 10, 10, 0.29, 0.31}, + {"files that align to nothing", 0, 10, 10, 0, 0}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { t.Parallel() - got := completeness(tt.got, tt.want) + got := completeness(tt.aligned, tt.audio, tt.want) if got < tt.minScore || got > tt.maxScore { t.Errorf( - "completeness(%d, %d) = %f, want in [%f, %f]", - tt.got, tt.want, got, tt.minScore, tt.maxScore, + "completeness(%d, %d, %d) = %f, want in [%f, %f]", + tt.aligned, tt.audio, tt.want, got, tt.minScore, tt.maxScore, ) } }) diff --git a/backend/download/search_recall_test.go b/backend/download/search_recall_test.go new file mode 100644 index 0000000..4973733 --- /dev/null +++ b/backend/download/search_recall_test.go @@ -0,0 +1,238 @@ +package download + +import ( + "context" + "strings" + "testing" +) + +// What Soulseek is asked, how, and what is kept from the answer (#271). + +func TestSlskdQueries(t *testing.T) { + t.Parallel() + + cases := []struct { + name string + dl Download + want []string + }{ + { + name: "a plain request is searched once", + dl: Download{Artist: "Radiohead", Album: "OK Computer"}, + want: []string{"Radiohead OK Computer"}, + }, + { + name: "an edition qualifier gets a second query without it", + dl: Download{Artist: "Radiohead", Album: "OK Computer (Collector's Edition)"}, + want: []string{ + "Radiohead OK Computer (Collector's Edition)", + "Radiohead OK Computer", + }, + }, + { + name: "a trailing remaster note", + dl: Download{Artist: "Pink Floyd", Album: "Animals - 2018 Remaster"}, + want: []string{ + "Pink Floyd Animals - 2018 Remaster", + "Pink Floyd Animals", + }, + }, + { + name: "a leading dash would be an exclusion", + dl: Download{Artist: "Mocky", Album: "-ism"}, + want: []string{"Mocky -ism", "Mocky ism"}, + }, + { + name: "a compilation is not searched by its placeholder artist", + dl: Download{Artist: "Various Artists", Album: "Pulp Fiction"}, + want: []string{"Various Artists Pulp Fiction", "Pulp Fiction"}, + }, + { + name: "what the user typed is searched as written", + dl: Download{Query: "ok computer (deluxe)", Album: "OK Computer (Deluxe)"}, + want: []string{"ok computer (deluxe)"}, + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + got := slskdQueries(tc.dl) + if strings.Join(got, "|") != strings.Join(tc.want, "|") { + t.Errorf("slskdQueries = %q, want %q", got, tc.want) + } + }) + } +} + +// Both queries run, the options are stated rather than left to the +// daemon's defaults, and a folder both queries found is one candidate. +func TestSlskdSearchRunsBothQueriesAndMerges(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + stub.responses = []slskdResponse{{ + Username: "peer", + Files: []slskdFile{ + {Filename: `\m\Radiohead - OK Computer\01 Airbag.flac`, Size: 1, Length: 284}, + {Filename: `\m\Radiohead - OK Computer\02 Paranoid Android.flac`, Size: 1, Length: 383}, + }, + }} + + s, _ := newStubSlskd(t, stub) + + got, err := s.Search(context.Background(), Download{ + ReleaseMBID: "rel", Artist: "Radiohead", Album: "OK Computer (Deluxe Edition)", + }) + if err != nil { + t.Fatalf("Search: %v", err) + } + + if len(got) != 1 { + t.Fatalf("got %d candidates, want the one folder once", len(got)) + } + + if got[0].Files[0].LengthMillis != 284_000 { + t.Errorf("length = %d ms, want 284000 from slskd's seconds", got[0].Files[0].LengthMillis) + } + + stub.mu.Lock() + searches := append([]map[string]any(nil), stub.searches...) + gets := append([]string(nil), stub.searchGets...) + stub.mu.Unlock() + + if len(searches) != 2 { + t.Fatalf("ran %d searches, want 2", len(searches)) + } + + for _, body := range searches { + for _, key := range []string{ + "searchTimeout", "responseLimit", "fileLimit", + "minimumResponseFileCount", "maximumPeerQueueLength", + } { + if _, ok := body[key]; !ok { + t.Errorf("search %q does not state %s", body["searchText"], key) + } + } + } + + // The responses are fetched once at the end, not with every poll. + for _, uri := range gets { + if strings.Contains(uri, "includeResponses") { + t.Errorf("poll %s asked for every response", uri) + } + } +} + +// A daemon without the responses endpoint still returns results. +func TestSlskdSearchFallsBackForAnOlderDaemon(t *testing.T) { + t.Parallel() + + stub := newSlskdStub(t) + stub.noResponsesEndpoint = true + stub.responses = []slskdResponse{{ + Username: "peer", + Files: []slskdFile{ + {Filename: `\m\Album\01 A.flac`, Size: 1}, + {Filename: `\m\Album\02 B.flac`, Size: 1}, + }, + }} + + s, _ := newStubSlskd(t, stub) + + got, err := s.Search(context.Background(), Download{Query: "album"}) + if err != nil { + t.Fatalf("Search: %v", err) + } + + if len(got) != 1 { + t.Errorf("got %d candidates, want 1 through the fallback", len(got)) + } +} + +func TestDurationAgreement(t *testing.T) { + t.Parallel() + + cases := []struct { + got, want int64 + score float64 + }{ + {300_000, 300_000, 1}, + {301_500, 300_000, 1}, // a second of silence + {300_000, 316_500, 0.5}, + {300_000, 345_000, 0}, // a different edit + } + + for _, tc := range cases { + if got := durationAgreement(tc.got, tc.want); got < tc.score-0.01 || got > tc.score+0.01 { + t.Errorf("durationAgreement(%d, %d) = %f, want %f", tc.got, tc.want, got, tc.score) + } + } +} + +// Two folders with the same track names are told apart by their +// lengths: one is the album, the other a live record of the same songs. +func TestDurationsSeparateTheRightRecording(t *testing.T) { + t.Parallel() + + dl := okComputer() + + timed := func(id string, lengths ...int64) Candidate { + c := candidateFor(id, allTitles(), ".flac", 30_000_000) + for i := range c.Files { + c.Files[i].LengthMillis = lengths[i] + } + + return c + } + + studio := timed("studio", trackMillis, trackMillis+1_000, trackMillis, trackMillis-500) + live := timed( + "live", + trackMillis+60_000, + trackMillis+75_000, + trackMillis+50_000, + trackMillis+90_000, + ) + + ranked := Rank(dl, []Candidate{live, studio}, nil, AutoDownloadPrefs{}) + + if ranked[0].ID != "studio" { + t.Fatalf("winner = %s, want the recording whose lengths match", ranked[0].ID) + } + + if !ranked[0].Match.DurationKnown || ranked[0].Match.DurationFit < 0.99 { + t.Errorf( + "studio duration fit = %f known=%v", + ranked[0].Match.DurationFit, + ranked[0].Match.DurationKnown, + ) + } + + if ranked[1].Match.DurationFit != 0 { + t.Errorf("live duration fit = %f, want 0", ranked[1].Match.DurationFit) + } +} + +// Without lengths the score is exactly what it was before durations +// were read, so a provider that reports none is not penalised. +func TestUnknownDurationsLeaveTheScoreAlone(t *testing.T) { + t.Parallel() + + dl := okComputer() + c := Score(dl, candidateFor("c", allTitles(), ".flac", 30_000_000), 50, AutoDownloadPrefs{}) + + if c.Match.DurationKnown { + t.Fatal("no file states a length, yet durations are known") + } + + want := weightTitleFit*c.Match.TitleFit + + weightCompleteness*c.Match.Completeness + + weightAlbumFit*c.Match.AlbumFit + + weightArtistFit*c.Match.ArtistFit + + if c.Match.Overall != want { + t.Errorf("match = %f, want the untimed formula's %f", c.Match.Overall, want) + } +} diff --git a/backend/download/types.go b/backend/download/types.go index d76c2ec..b753267 100644 --- a/backend/download/types.go +++ b/backend/download/types.go @@ -232,12 +232,17 @@ type Candidate struct { // results give paths and sizes but no tags, so Format and duration are // inferred from the path and size where possible. type CandidateFile struct { - Path string `json:"path"` - Size int64 `json:"size"` - Format Format `json:"format"` - Bitrate int `json:"bitrate,omitempty"` // kbps, 0 when unknown - IsAudio bool `json:"isAudio"` - MatchedTo int `json:"matchedTo,omitempty"` // expected track position + Path string `json:"path"` + Size int64 `json:"size"` + Format Format `json:"format"` + Bitrate int `json:"bitrate,omitempty"` // kbps, 0 when unknown + IsAudio bool `json:"isAudio"` + + // LengthMillis is the file's duration as the source reports it, or + // 0 when it does not. Soulseek reports it for most audio files. + LengthMillis int64 `json:"lengthMillis,omitempty"` + + MatchedTo int `json:"matchedTo,omitempty"` // expected track position } // Format is a normalized audio container/codec name. @@ -286,7 +291,14 @@ type MatchScore struct { TitleFit float64 `json:"titleFit"` // filenames vs expected titles ArtistFit float64 `json:"artistFit"` // path/origin vs expected artist AlbumFit float64 `json:"albumFit"` // folder name vs album title - Completeness float64 `json:"completeness"` // audio files vs expected count + Completeness float64 `json:"completeness"` // aligned tracks vs expected count + + // DurationFit is how well the aligned files' lengths agree with the + // expected tracks', and DurationKnown whether enough of them stated + // a length for that to count. When it does not, the score is the + // four text signals alone, exactly as before durations were read. + DurationFit float64 `json:"durationFit"` + DurationKnown bool `json:"durationKnown"` // Anchored records whether an MBID drove this score. Unanchored // matches are capped, because there is nothing to be right about. diff --git a/backend/library/coverart_storage_test.go b/backend/library/coverart_storage_test.go index 2c2b60d..1f0e6b7 100644 --- a/backend/library/coverart_storage_test.go +++ b/backend/library/coverart_storage_test.go @@ -8,6 +8,7 @@ import ( "yellowjacket/backend/coverart" "yellowjacket/backend/database/sql/sqlcgen" + "yellowjacket/internal/testfixtures" ) // TestScan_StoresOnlyCoverTiers pins the size decision: a scan writes @@ -26,10 +27,9 @@ func TestScan_StoresOnlyCoverTiers(t *testing.T) { lib, db := setupTestLibrary(t) - root, err := filepath.Abs("../../test_data/music_library_test") - if err != nil { - t.Fatalf("resolve fixture path: %v", err) - } + // Load skips when the fixture library has not been generated, as + // every other fixture test does. + root := testfixtures.Load(t).Root() library, err := db.Queries.CreateLibrary(lib.ctx, sqlcgen.CreateLibraryParams{ Name: "Fixtures", diff --git a/frontend/bindings/yellowjacket/backend/download/models.ts b/frontend/bindings/yellowjacket/backend/download/models.ts index fbe539a..0db4f82 100644 --- a/frontend/bindings/yellowjacket/backend/download/models.ts +++ b/frontend/bindings/yellowjacket/backend/download/models.ts @@ -121,6 +121,12 @@ export interface CandidateFile { "bitrate"?: number; "isAudio": boolean; + /** + * LengthMillis is the file's duration as the source reports it, or + * 0 when it does not. Soulseek reports it for most audio files. + */ + "lengthMillis"?: number; + /** * expected track position */ @@ -453,10 +459,19 @@ export interface MatchScore { "albumFit": number; /** - * audio files vs expected count + * aligned tracks vs expected count */ "completeness": number; + /** + * DurationFit is how well the aligned files' lengths agree with the + * expected tracks', and DurationKnown whether enough of them stated + * a length for that to count. When it does not, the score is the + * four text signals alone, exactly as before durations were read. + */ + "durationFit": number; + "durationKnown": boolean; + /** * Anchored records whether an MBID drove this score. Unanchored * matches are capped, because there is nothing to be right about.