mirror of
https://github.com/navidrome/navidrome.git
synced 2026-08-31 07:30:32 +00:00
Merge 6d4ac1506ff7f435418f79de1fc4cd7b097d0e60 into dbd26ba2e71d0a5b79dba873a2beeff59f1cd8dd
This commit is contained in:
commit
eb2bdaf52e
@ -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,
|
||||
)
|
||||
|
||||
@ -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"},
|
||||
|
||||
@ -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
|
||||
}
|
||||
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user