This commit is contained in:
Yuuta 2026-10-06 20:09:06 +03:00 • committed by GitHub
commit f364644ee3
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 309 additions and 5 deletions

View file

@ -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,
)

View file

@ -56,6 +56,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"},

View file

@ -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
}