From fc68908ee5412c2f11999909236ad7f17af6025d Mon Sep 17 00:00:00 2001 From: ranokay Date: Sun, 23 Aug 2026 06:42:39 +0300 Subject: [PATCH 1/4] perf(plugins): coalesce lyrics requests --- plugins/lyrics_adapter.go | 54 +++++++++++++++++- plugins/lyrics_adapter_test.go | 100 +++++++++++++++++++++++++++++++++ 2 files changed, 152 insertions(+), 2 deletions(-) diff --git a/plugins/lyrics_adapter.go b/plugins/lyrics_adapter.go index 9e02115e7..2b41b7603 100644 --- a/plugins/lyrics_adapter.go +++ b/plugins/lyrics_adapter.go @@ -2,10 +2,13 @@ package plugins import ( "context" + "fmt" + "time" "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/plugins/capabilities" + "golang.org/x/sync/singleflight" ) const CapabilityLyrics Capability = "Lyrics" @@ -18,6 +21,10 @@ 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 + func init() { registerCapability( CapabilityLyrics, @@ -33,11 +40,54 @@ func newLyricsPlugin(p *plugin) *LyricsPlugin { type LyricsPlugin struct { name string plugin *plugin + calls singleflight.Group } -// 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 is +// detached from any one request so one disconnected client does not cancel it +// for the remaining callers. func (l *LyricsPlugin) GetLyrics(ctx context.Context, mf *model.MediaFile) (model.LyricList, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + + result := l.calls.DoChan(lyricsPluginCallKey(mf), func() (any, error) { + callCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), lyricsPluginCallTimeout) + defer cancel() + return l.getLyrics(callCtx, mf) + }) + + select { + case call := <-result: + if call.Err != nil { + return nil, call.Err + } + lyricsList, ok := call.Val.(model.LyricList) + if !ok { + return nil, fmt.Errorf("unexpected lyrics plugin result type %T", call.Val) + } + return lyricsList, nil + case <-ctx.Done(): + return nil, ctx.Err() + } +} + +func lyricsPluginCallKey(mf *model.MediaFile) string { + return fmt.Sprintf( + "%s\x00%s\x00%d\x00%s\x00%s\x00%s\x00%.3f", + mf.ID, + mf.Path, + mf.UpdatedAt.UnixNano(), + mf.Title, + mf.Artist, + mf.Album, + mf.Duration, + ) +} + +// 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) (model.LyricList, error) { select { case l.plugin.lyricsSem <- struct{}{}: defer func() { <-l.plugin.lyricsSem }() diff --git a/plugins/lyrics_adapter_test.go b/plugins/lyrics_adapter_test.go index d110665f5..5aba65adb 100644 --- a/plugins/lyrics_adapter_test.go +++ b/plugins/lyrics_adapter_test.go @@ -58,6 +58,106 @@ 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, + ) + + 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 + } + }) + + 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 range 2 { + go func() { + <-start + lyrics, err := coalescingProvider.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("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("defaults language to 'xxx' when plugin does not provide one", func() { manager, _ := createTestManagerWithPlugins(map[string]map[string]string{ "test-lyrics": {"no_lang": "true"}, From 944ba23390cf90e9c4f6469ec6b5afd5b9a387af Mon Sep 17 00:00:00 2001 From: ranokay Date: Sun, 23 Aug 2026 06:57:10 +0300 Subject: [PATCH 2/4] fix(plugins): key coalescing by request metadata --- plugins/lyrics_adapter.go | 35 +++++++++++++------------ plugins/lyrics_adapter_test.go | 48 ++++++++++++++++++++++++++++++++++ 2 files changed, 66 insertions(+), 17 deletions(-) diff --git a/plugins/lyrics_adapter.go b/plugins/lyrics_adapter.go index 2b41b7603..5ad70be02 100644 --- a/plugins/lyrics_adapter.go +++ b/plugins/lyrics_adapter.go @@ -2,6 +2,7 @@ package plugins import ( "context" + "encoding/json" "fmt" "time" @@ -51,10 +52,18 @@ func (l *LyricsPlugin) GetLyrics(ctx context.Context, mf *model.MediaFile) (mode return nil, err } - result := l.calls.DoChan(lyricsPluginCallKey(mf), func() (any, error) { + req := capabilities.GetLyricsRequest{ + Track: mediaFileToTrackInfo(l.plugin, mf), + } + key, err := lyricsPluginCallKey(req) + if err != nil { + return nil, err + } + + result := l.calls.DoChan(key, func() (any, error) { callCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), lyricsPluginCallTimeout) defer cancel() - return l.getLyrics(callCtx, mf) + return l.getLyrics(callCtx, mf, req) }) select { @@ -72,31 +81,23 @@ func (l *LyricsPlugin) GetLyrics(ctx context.Context, mf *model.MediaFile) (mode } } -func lyricsPluginCallKey(mf *model.MediaFile) string { - return fmt.Sprintf( - "%s\x00%s\x00%d\x00%s\x00%s\x00%s\x00%.3f", - mf.ID, - mf.Path, - mf.UpdatedAt.UnixNano(), - mf.Title, - mf.Artist, - mf.Album, - mf.Duration, - ) +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) (model.LyricList, error) { +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 5aba65adb..b366a387e 100644 --- a/plugins/lyrics_adapter_test.go +++ b/plugins/lyrics_adapter_test.go @@ -111,6 +111,54 @@ var _ = Describe("LyricsPlugin", Ordered, func() { 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( From f1b56605322d48dd9644db0ba56be4dc03e3fd89 Mon Sep 17 00:00:00 2001 From: ranokay Date: Sun, 23 Aug 2026 07:03:28 +0300 Subject: [PATCH 3/4] fix(plugins): share lyrics coalescing state --- plugins/lyrics_adapter.go | 4 +--- plugins/lyrics_adapter_test.go | 14 +++++++++----- plugins/manager_plugin.go | 6 ++++-- 3 files changed, 14 insertions(+), 10 deletions(-) diff --git a/plugins/lyrics_adapter.go b/plugins/lyrics_adapter.go index 5ad70be02..20479c333 100644 --- a/plugins/lyrics_adapter.go +++ b/plugins/lyrics_adapter.go @@ -9,7 +9,6 @@ import ( "github.com/navidrome/navidrome/log" "github.com/navidrome/navidrome/model" "github.com/navidrome/navidrome/plugins/capabilities" - "golang.org/x/sync/singleflight" ) const CapabilityLyrics Capability = "Lyrics" @@ -41,7 +40,6 @@ func newLyricsPlugin(p *plugin) *LyricsPlugin { type LyricsPlugin struct { name string plugin *plugin - calls singleflight.Group } // GetLyrics coalesces concurrent lookups for the same track. The shared call is @@ -60,7 +58,7 @@ func (l *LyricsPlugin) GetLyrics(ctx context.Context, mf *model.MediaFile) (mode return nil, err } - result := l.calls.DoChan(key, func() (any, error) { + result := l.plugin.lyricsCalls.DoChan(key, func() (any, error) { callCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), lyricsPluginCallTimeout) defer cancel() return l.getLyrics(callCtx, mf, req) diff --git a/plugins/lyrics_adapter_test.go b/plugins/lyrics_adapter_test.go index b366a387e..dc273a2a1 100644 --- a/plugins/lyrics_adapter_test.go +++ b/plugins/lyrics_adapter_test.go @@ -66,11 +66,15 @@ var _ = Describe("LyricsPlugin", Ordered, func() { "test-lyrics"+PackageExtension, ) - p, ok := manager.LoadLyricsProvider("test-lyrics") + first, ok := manager.LoadLyricsProvider("test-lyrics") Expect(ok).To(BeTrue()) - coalescingProvider := p.(*LyricsPlugin) + second, ok := manager.LoadLyricsProvider("test-lyrics") + Expect(ok).To(BeTrue()) + firstProvider := first.(*LyricsPlugin) + secondProvider := second.(*LyricsPlugin) + Expect(firstProvider).ToNot(BeIdenticalTo(secondProvider)) - sem := coalescingProvider.plugin.lyricsSem + sem := firstProvider.plugin.lyricsSem for range cap(sem) { sem <- struct{}{} } @@ -87,10 +91,10 @@ var _ = Describe("LyricsPlugin", Ordered, func() { start := make(chan struct{}) results := make(chan callResult, 2) track := &model.MediaFile{ID: "shared-track", Title: "Test Song", Artist: "Test Artist"} - for range 2 { + for _, provider := range []*LyricsPlugin{firstProvider, secondProvider} { go func() { <-start - lyrics, err := coalescingProvider.GetLyrics(GinkgoT().Context(), track) + lyrics, err := provider.GetLyrics(GinkgoT().Context(), track) results <- callResult{lyrics: lyrics, err: err} }() } diff --git a/plugins/manager_plugin.go b/plugins/manager_plugin.go index 13375a70f..394b4c881 100644 --- a/plugins/manager_plugin.go +++ b/plugins/manager_plugin.go @@ -10,6 +10,7 @@ import ( extism "github.com/extism/go-sdk" "github.com/navidrome/navidrome/model" "github.com/tetratelabs/wazero" + "golang.org/x/sync/singleflight" ) // plugin represents a loaded plugin @@ -24,8 +25,9 @@ type plugin struct { allowedUserIDs []string // User IDs this plugin can access (from DB configuration) allUsers bool // If true, plugin can access all users libraries libraryAccess - lyricsSem chan struct{} // Caps concurrent lyrics calls (see LyricsPlugin.GetLyrics) - fsConfig wazero.FSConfig // Sandboxed library mounts, nil if no filesystem permission + lyricsSem chan struct{} // Caps concurrent lyrics calls (see LyricsPlugin.GetLyrics) + lyricsCalls singleflight.Group // Shared by the transient LyricsPlugin adapters + fsConfig wazero.FSConfig // Sandboxed library mounts, nil if no filesystem permission } // instanceConfig is used by every call site, so all instances get the sandboxed mounts. From 6d4ac1506ff7f435418f79de1fc4cd7b097d0e60 Mon Sep 17 00:00:00 2001 From: ranokay Date: Tue, 25 Aug 2026 20:53:01 +0300 Subject: [PATCH 4/4] fix(plugins): cancel abandoned lyrics lookups --- plugins/lyrics_adapter.go | 102 ++++++++++++++++++++++++++++----- plugins/lyrics_adapter_test.go | 30 ++++++++++ plugins/manager_plugin.go | 7 +-- 3 files changed, 120 insertions(+), 19 deletions(-) diff --git a/plugins/lyrics_adapter.go b/plugins/lyrics_adapter.go index 20479c333..76a0fb885 100644 --- a/plugins/lyrics_adapter.go +++ b/plugins/lyrics_adapter.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "fmt" + "sync" "time" "github.com/navidrome/navidrome/log" @@ -25,6 +26,86 @@ const maxConcurrentLyricsCalls = 2 // 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, @@ -42,9 +123,8 @@ type LyricsPlugin struct { plugin *plugin } -// GetLyrics coalesces concurrent lookups for the same track. The shared call is -// detached from any one request so one disconnected client does not cancel it -// for the remaining callers. +// 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 @@ -58,22 +138,14 @@ func (l *LyricsPlugin) GetLyrics(ctx context.Context, mf *model.MediaFile) (mode return nil, err } - result := l.plugin.lyricsCalls.DoChan(key, func() (any, error) { - callCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), lyricsPluginCallTimeout) - defer cancel() + 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 := <-result: - if call.Err != nil { - return nil, call.Err - } - lyricsList, ok := call.Val.(model.LyricList) - if !ok { - return nil, fmt.Errorf("unexpected lyrics plugin result type %T", call.Val) - } - return lyricsList, nil + case <-call.done: + return call.lyrics, call.err case <-ctx.Done(): return nil, ctx.Err() } diff --git a/plugins/lyrics_adapter_test.go b/plugins/lyrics_adapter_test.go index dc273a2a1..ac8a5b158 100644 --- a/plugins/lyrics_adapter_test.go +++ b/plugins/lyrics_adapter_test.go @@ -210,6 +210,36 @@ var _ = Describe("LyricsPlugin", Ordered, func() { 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 394b4c881..cda2cd1f5 100644 --- a/plugins/manager_plugin.go +++ b/plugins/manager_plugin.go @@ -10,7 +10,6 @@ import ( extism "github.com/extism/go-sdk" "github.com/navidrome/navidrome/model" "github.com/tetratelabs/wazero" - "golang.org/x/sync/singleflight" ) // plugin represents a loaded plugin @@ -25,9 +24,9 @@ type plugin struct { allowedUserIDs []string // User IDs this plugin can access (from DB configuration) allUsers bool // If true, plugin can access all users libraries libraryAccess - lyricsSem chan struct{} // Caps concurrent lyrics calls (see LyricsPlugin.GetLyrics) - lyricsCalls singleflight.Group // Shared by the transient LyricsPlugin adapters - fsConfig wazero.FSConfig // Sandboxed library mounts, nil if no filesystem permission + 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 } // instanceConfig is used by every call site, so all instances get the sandboxed mounts.