mirror of
https://github.com/navidrome/navidrome.git
synced 2026-10-11 11:57:12 +02:00
fix(plugins): cancel abandoned lyrics lookups
This commit is contained in:
parent
f1b5660532
commit
6d4ac1506f
3 changed files with 120 additions and 19 deletions
|
|
@ -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()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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"},
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue