diff --git a/plugins/lyrics_adapter.go b/plugins/lyrics_adapter.go index 9e02115e7..76a0fb885 100644 --- a/plugins/lyrics_adapter.go +++ b/plugins/lyrics_adapter.go @@ -2,6 +2,10 @@ package plugins import ( "context" + "encoding/json" + "fmt" + "sync" + "time" "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" @@ -18,6 +22,90 @@ const ( // lyrics for whole queues, and the resulting burst can rate-limit upstream providers. const maxConcurrentLyricsCalls = 2 +// lyricsPluginCallTimeout bounds work detached from one caller's request so a +// disconnected client cannot leave a shared plugin lookup running forever. +const lyricsPluginCallTimeout = time.Minute + +// lyricsCallGroup coalesces lookups while keeping their lifetime tied to the +// callers that are still waiting. One caller may leave without interrupting the +// others, but the shared work is cancelled once the last waiter is gone. +type lyricsCallGroup struct { + mu sync.Mutex + calls map[string]*lyricsCall +} + +type lyricsCall struct { + group *lyricsCallGroup + key string + ctx context.Context + cancel context.CancelFunc + done chan struct{} + lyrics model.LyricList + err error + waiters int + finished bool +} + +func (g *lyricsCallGroup) join( + parent context.Context, + key string, + lookup func(context.Context) (model.LyricList, error), +) *lyricsCall { + g.mu.Lock() + if call := g.calls[key]; call != nil { + call.waiters++ + g.mu.Unlock() + return call + } + + ctx, cancel := context.WithTimeout(context.WithoutCancel(parent), lyricsPluginCallTimeout) + call := &lyricsCall{ + group: g, + key: key, + ctx: ctx, + cancel: cancel, + done: make(chan struct{}), + waiters: 1, + } + if g.calls == nil { + g.calls = make(map[string]*lyricsCall) + } + g.calls[key] = call + g.mu.Unlock() + + go call.run(lookup) + return call +} + +func (c *lyricsCall) run(lookup func(context.Context) (model.LyricList, error)) { + lyrics, err := lookup(c.ctx) + + c.group.mu.Lock() + c.lyrics = lyrics + c.err = err + c.finished = true + if c.group.calls[c.key] == c { + delete(c.group.calls, c.key) + } + close(c.done) + c.group.mu.Unlock() + c.cancel() +} + +func (c *lyricsCall) release() { + c.group.mu.Lock() + c.waiters-- + shouldCancel := c.waiters == 0 && !c.finished + if shouldCancel && c.group.calls[c.key] == c { + delete(c.group.calls, c.key) + } + c.group.mu.Unlock() + + if shouldCancel { + c.cancel() + } +} + func init() { registerCapability( CapabilityLyrics, @@ -35,18 +123,51 @@ type LyricsPlugin struct { plugin *plugin } -// GetLyrics calls the plugin to fetch lyrics, then content-sniffs each response -// via model.ParseLyrics (TTML/SRT/YAML/LRC/plain). +// GetLyrics coalesces concurrent lookups for the same track. The shared call +// survives individual disconnections while another caller is still waiting. func (l *LyricsPlugin) GetLyrics(ctx context.Context, mf *model.MediaFile) (model.LyricList, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + + req := capabilities.GetLyricsRequest{ + Track: mediaFileToTrackInfo(l.plugin, mf), + } + key, err := lyricsPluginCallKey(req) + if err != nil { + return nil, err + } + + call := l.plugin.lyricsCalls.join(ctx, key, func(callCtx context.Context) (model.LyricList, error) { + return l.getLyrics(callCtx, mf, req) + }) + defer call.release() + + select { + case <-call.done: + return call.lyrics, call.err + case <-ctx.Done(): + return nil, ctx.Err() + } +} + +func lyricsPluginCallKey(req capabilities.GetLyricsRequest) (string, error) { + value, err := json.Marshal(req) + if err != nil { + return "", fmt.Errorf("encode lyrics plugin request key: %w", err) + } + return string(value), nil +} + +// getLyrics calls the plugin, then content-sniffs each response via +// model.ParseLyrics (TTML/SRT/YAML/LRC/plain). +func (l *LyricsPlugin) getLyrics(ctx context.Context, mf *model.MediaFile, req capabilities.GetLyricsRequest) (model.LyricList, error) { select { case l.plugin.lyricsSem <- struct{}{}: defer func() { <-l.plugin.lyricsSem }() case <-ctx.Done(): return nil, ctx.Err() } - req := capabilities.GetLyricsRequest{ - Track: mediaFileToTrackInfo(l.plugin, mf), - } resp, err := callPluginFunction[capabilities.GetLyricsRequest, capabilities.GetLyricsResponse]( ctx, l.plugin, FuncLyricsGetLyrics, req, ) diff --git a/plugins/lyrics_adapter_test.go b/plugins/lyrics_adapter_test.go index d110665f5..ac8a5b158 100644 --- a/plugins/lyrics_adapter_test.go +++ b/plugins/lyrics_adapter_test.go @@ -58,6 +58,188 @@ var _ = Describe("LyricsPlugin", Ordered, func() { Expect(result[0].Line[0].Value).To(ContainSubstring("Test Song")) }) + It("coalesces concurrent requests for the same track", func() { + metrics := &mockMetricsRecorder{} + manager, _ := createTestManagerWithPluginsAndMetrics( + nil, + metrics, + "test-lyrics"+PackageExtension, + ) + + first, ok := manager.LoadLyricsProvider("test-lyrics") + Expect(ok).To(BeTrue()) + second, ok := manager.LoadLyricsProvider("test-lyrics") + Expect(ok).To(BeTrue()) + firstProvider := first.(*LyricsPlugin) + secondProvider := second.(*LyricsPlugin) + Expect(firstProvider).ToNot(BeIdenticalTo(secondProvider)) + + sem := firstProvider.plugin.lyricsSem + for range cap(sem) { + sem <- struct{}{} + } + DeferCleanup(func() { + for len(sem) > 0 { + <-sem + } + }) + + type callResult struct { + lyrics model.LyricList + err error + } + start := make(chan struct{}) + results := make(chan callResult, 2) + track := &model.MediaFile{ID: "shared-track", Title: "Test Song", Artist: "Test Artist"} + for _, provider := range []*LyricsPlugin{firstProvider, secondProvider} { + go func() { + <-start + lyrics, err := provider.GetLyrics(GinkgoT().Context(), track) + results <- callResult{lyrics: lyrics, err: err} + }() + } + close(start) + + Consistently(results, "500ms").ShouldNot(Receive()) + <-sem + + for range 2 { + var result callResult + Eventually(results).Should(Receive(&result)) + Expect(result.err).ToNot(HaveOccurred()) + Expect(result.lyrics).To(HaveLen(1)) + } + + calls := metrics.getCalls() + Expect(calls).To(HaveLen(1)) + Expect(calls[0].method).To(Equal(FuncLyricsGetLyrics)) + }) + + It("does not coalesce requests with different plugin metadata", func() { + metrics := &mockMetricsRecorder{} + manager, _ := createTestManagerWithPluginsAndMetrics( + nil, + metrics, + "test-lyrics"+PackageExtension, + ) + + p, ok := manager.LoadLyricsProvider("test-lyrics") + Expect(ok).To(BeTrue()) + coalescingProvider := p.(*LyricsPlugin) + + sem := coalescingProvider.plugin.lyricsSem + for range cap(sem) { + sem <- struct{}{} + } + DeferCleanup(func() { + for len(sem) > 0 { + <-sem + } + }) + + start := make(chan struct{}) + results := make(chan error, 2) + tracks := []*model.MediaFile{ + {ID: "same-id", Title: "Test Song", Artist: "Test Artist", TrackNumber: 1}, + {ID: "same-id", Title: "Test Song", Artist: "Test Artist", TrackNumber: 2}, + } + for _, track := range tracks { + go func() { + <-start + _, err := coalescingProvider.GetLyrics(GinkgoT().Context(), track) + results <- err + }() + } + close(start) + + Consistently(results, "500ms").ShouldNot(Receive()) + for range cap(sem) { + <-sem + } + + for range tracks { + Eventually(results).Should(Receive(BeNil())) + } + Expect(metrics.getCalls()).To(HaveLen(2)) + }) + + It("keeps a shared request alive when one caller cancels", func() { + metrics := &mockMetricsRecorder{} + manager, _ := createTestManagerWithPluginsAndMetrics( + nil, + metrics, + "test-lyrics"+PackageExtension, + ) + + p, ok := manager.LoadLyricsProvider("test-lyrics") + Expect(ok).To(BeTrue()) + coalescingProvider := p.(*LyricsPlugin) + + sem := coalescingProvider.plugin.lyricsSem + for range cap(sem) { + sem <- struct{}{} + } + DeferCleanup(func() { + for len(sem) > 0 { + <-sem + } + }) + + track := &model.MediaFile{ID: "shared-track", Title: "Test Song", Artist: "Test Artist"} + firstCtx, cancelFirst := context.WithCancel(GinkgoT().Context()) + firstDone := make(chan error, 1) + go func() { + _, err := coalescingProvider.GetLyrics(firstCtx, track) + firstDone <- err + }() + Consistently(firstDone, "100ms").ShouldNot(Receive()) + + secondDone := make(chan error, 1) + go func() { + _, err := coalescingProvider.GetLyrics(GinkgoT().Context(), track) + secondDone <- err + }() + Consistently(secondDone, "100ms").ShouldNot(Receive()) + + cancelFirst() + Eventually(firstDone).Should(Receive(MatchError(context.Canceled))) + Consistently(secondDone, "100ms").ShouldNot(Receive()) + + <-sem + Eventually(secondDone).Should(Receive(BeNil())) + Expect(metrics.getCalls()).To(HaveLen(1)) + }) + + It("cancels a shared request after its last caller leaves", func() { + var group lyricsCallGroup + firstCtx, cancelFirst := context.WithCancel(GinkgoT().Context()) + secondCtx, cancelSecond := context.WithCancel(GinkgoT().Context()) + started := make(chan struct{}) + stopped := make(chan error, 1) + + firstCall := group.join(firstCtx, "shared", func(ctx context.Context) (model.LyricList, error) { + close(started) + <-ctx.Done() + stopped <- ctx.Err() + return nil, ctx.Err() + }) + Eventually(started).Should(BeClosed()) + + secondCall := group.join(secondCtx, "shared", func(context.Context) (model.LyricList, error) { + Fail("started a second lookup for the same key") + return nil, nil + }) + Expect(secondCall).To(BeIdenticalTo(firstCall)) + + cancelFirst() + firstCall.release() + Consistently(stopped, "100ms").ShouldNot(Receive()) + + cancelSecond() + secondCall.release() + Eventually(stopped).Should(Receive(MatchError(context.Canceled))) + }) + It("defaults language to 'xxx' when plugin does not provide one", func() { manager, _ := createTestManagerWithPlugins(map[string]map[string]string{ "test-lyrics": {"no_lang": "true"}, diff --git a/plugins/manager_plugin.go b/plugins/manager_plugin.go index 13375a70f..cda2cd1f5 100644 --- a/plugins/manager_plugin.go +++ b/plugins/manager_plugin.go @@ -25,6 +25,7 @@ type plugin struct { allUsers bool // If true, plugin can access all users libraries libraryAccess lyricsSem chan struct{} // Caps concurrent lyrics calls (see LyricsPlugin.GetLyrics) + lyricsCalls lyricsCallGroup // Shared by the transient LyricsPlugin adapters fsConfig wazero.FSConfig // Sandboxed library mounts, nil if no filesystem permission }