mirror of
https://github.com/navidrome/navidrome.git
synced 2026-10-11 03:47:18 +02:00
Eliminate redundant work and minor issues found during code review: - Replace manual PlaylistTrack construction in syncPlaylist with the existing Playlist.AddMediaFiles helper, removing duplicated logic - Pre-sanitize track fields once per artist batch in the matcher's fuzzy matching loop, avoiding redundant sanitization in both findBestMatch and computeSpecificityLevel on every iteration - Cache resolved usernames in discoverAndSync to avoid N+1 DB lookups when multiple playlists share the same owner - Use the local loadedPlugin variable instead of reading m.plugins[p.ID] after releasing the lock in loadPluginWithConfig - Fix misleading uint32 comparison (<=0 to ==0) in durationProximity - Update stale comment on checkTracksEditable to mention plugin playlists
274 lines
8.7 KiB
Go
274 lines
8.7 KiB
Go
package plugins
|
|
|
|
import (
|
|
"context"
|
|
"slices"
|
|
"strings"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"github.com/navidrome/navidrome/core/matcher"
|
|
"github.com/navidrome/navidrome/log"
|
|
"github.com/navidrome/navidrome/model"
|
|
"github.com/navidrome/navidrome/model/id"
|
|
"github.com/navidrome/navidrome/plugins/capabilities"
|
|
)
|
|
|
|
const (
|
|
CapabilityPlaylistProvider Capability = "PlaylistProvider"
|
|
|
|
FuncPlaylistProviderGetAvailablePlaylists = "nd_playlist_provider_get_available_playlists"
|
|
FuncPlaylistProviderGetPlaylist = "nd_playlist_provider_get_playlist"
|
|
|
|
// workChCapacity is the buffer size for the work channel.
|
|
workChCapacity = 64
|
|
|
|
// discoveryRetryDelay is how long to wait before retrying a failed GetAvailablePlaylists call.
|
|
discoveryRetryDelay = 5 * time.Minute
|
|
)
|
|
|
|
func init() {
|
|
registerCapability(
|
|
CapabilityPlaylistProvider,
|
|
FuncPlaylistProviderGetAvailablePlaylists,
|
|
FuncPlaylistProviderGetPlaylist,
|
|
)
|
|
}
|
|
|
|
type workType int
|
|
|
|
const (
|
|
workDiscover workType = iota // run discoverAndSync
|
|
workSync // run syncPlaylist for a single playlist
|
|
)
|
|
|
|
type workItem struct {
|
|
typ workType
|
|
info capabilities.PlaylistInfo // only for workSync
|
|
dbID string // only for workSync
|
|
ownerID string // only for workSync
|
|
}
|
|
|
|
// playlistSyncer manages playlist synchronization for a single plugin.
|
|
// All mutable state (refreshTimers, discoveryTimer) is owned exclusively by the
|
|
// worker goroutine — no synchronization needed. The retryInterval and
|
|
// refreshTimerCount fields use atomics so tests can observe them race-free.
|
|
type playlistSyncer struct {
|
|
pluginName string
|
|
plugin *plugin
|
|
ds model.DataStore
|
|
matcher *matcher.Matcher
|
|
ctx context.Context
|
|
cancel context.CancelFunc
|
|
workCh chan workItem // serialized work queue
|
|
refreshTimers map[string]*time.Timer // keyed by playlist DB ID — worker-only
|
|
discoveryTimer *time.Timer // worker-only
|
|
retryInterval atomic.Int64 // nanoseconds; from last GetAvailablePlaylists response
|
|
refreshTimerCount atomic.Int32 // number of active refresh timers
|
|
done chan struct{} // closed when worker exits
|
|
}
|
|
|
|
func newPlaylistSyncer(parentCtx context.Context, pluginName string, p *plugin, ds model.DataStore, m *matcher.Matcher) *playlistSyncer {
|
|
ctx, cancel := context.WithCancel(parentCtx)
|
|
return &playlistSyncer{
|
|
pluginName: pluginName,
|
|
plugin: p,
|
|
ds: ds,
|
|
matcher: m,
|
|
ctx: ctx,
|
|
cancel: cancel,
|
|
workCh: make(chan workItem, workChCapacity),
|
|
refreshTimers: make(map[string]*time.Timer),
|
|
done: make(chan struct{}),
|
|
}
|
|
}
|
|
|
|
// run is the single worker goroutine that processes all work items sequentially.
|
|
// It performs an initial discovery before entering the main loop.
|
|
func (o *playlistSyncer) run() {
|
|
defer close(o.done)
|
|
|
|
// Run initial discovery before entering the loop
|
|
o.discoverAndSync()
|
|
|
|
for {
|
|
select {
|
|
case <-o.ctx.Done():
|
|
o.stopAllTimers()
|
|
return
|
|
case item := <-o.workCh:
|
|
switch item.typ {
|
|
case workDiscover:
|
|
o.discoverAndSync()
|
|
case workSync:
|
|
o.syncPlaylist(item.info, item.dbID, item.ownerID)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// discoverAndSync calls GetAvailablePlaylists, then GetPlaylist for each, matches tracks, and upserts.
|
|
func (o *playlistSyncer) discoverAndSync() {
|
|
ctx := o.ctx
|
|
resp, err := callPluginFunction[capabilities.GetAvailablePlaylistsRequest, capabilities.GetAvailablePlaylistsResponse](
|
|
ctx, o.plugin, FuncPlaylistProviderGetAvailablePlaylists, capabilities.GetAvailablePlaylistsRequest{},
|
|
)
|
|
if err != nil {
|
|
log.Error(ctx, "Failed to call GetAvailablePlaylists, retrying later", "plugin", o.pluginName, err)
|
|
o.scheduleDiscovery(discoveryRetryDelay)
|
|
return
|
|
}
|
|
|
|
// Store retry interval from response
|
|
if resp.RetryInterval > 0 {
|
|
o.retryInterval.Store(int64(time.Duration(resp.RetryInterval) * time.Second))
|
|
}
|
|
|
|
resolvedUsers := map[string]string{} // username -> userID cache
|
|
for _, info := range resp.Playlists {
|
|
// Resolve username to user ID (cached)
|
|
ownerID, ok := resolvedUsers[info.OwnerUsername]
|
|
if !ok {
|
|
user, err := o.ds.User(adminContext(ctx)).FindByUsername(info.OwnerUsername)
|
|
if err != nil {
|
|
log.Error(ctx, "Failed to resolve playlist owner", "plugin", o.pluginName,
|
|
"playlistID", info.ID, "username", info.OwnerUsername, err)
|
|
continue
|
|
}
|
|
ownerID = user.ID
|
|
resolvedUsers[info.OwnerUsername] = ownerID
|
|
}
|
|
|
|
// Validate that the plugin is permitted to create playlists for this user
|
|
if !o.plugin.allUsers && !slices.Contains(o.plugin.allowedUserIDs, ownerID) {
|
|
log.Error(ctx, "Plugin not permitted to create playlists for user", "plugin", o.pluginName,
|
|
"playlistID", info.ID, "username", info.OwnerUsername)
|
|
continue
|
|
}
|
|
|
|
dbID := id.NewHash(o.pluginName, info.ID, ownerID)
|
|
o.syncPlaylist(info, dbID, ownerID)
|
|
}
|
|
|
|
// Schedule re-discovery if RefreshInterval > 0
|
|
if resp.RefreshInterval > 0 {
|
|
o.scheduleDiscovery(time.Duration(resp.RefreshInterval) * time.Second)
|
|
}
|
|
}
|
|
|
|
// syncPlaylist calls GetPlaylist, matches tracks, and upserts the playlist in the DB.
|
|
func (o *playlistSyncer) syncPlaylist(info capabilities.PlaylistInfo, dbID string, ownerID string) {
|
|
ctx := o.ctx
|
|
resp, err := callPluginFunction[capabilities.GetPlaylistRequest, capabilities.GetPlaylistResponse](
|
|
ctx, o.plugin, FuncPlaylistProviderGetPlaylist, capabilities.GetPlaylistRequest{ID: info.ID},
|
|
)
|
|
if err != nil {
|
|
if isPlaylistNotFoundError(err) {
|
|
log.Info(ctx, "Playlist not found, skipping", "plugin", o.pluginName, "playlistID", info.ID)
|
|
// Stop any existing refresh timer for this playlist
|
|
if timer, ok := o.refreshTimers[dbID]; ok {
|
|
timer.Stop()
|
|
delete(o.refreshTimers, dbID)
|
|
o.refreshTimerCount.Store(int32(len(o.refreshTimers)))
|
|
}
|
|
return
|
|
}
|
|
log.Warn(ctx, "Failed to call GetPlaylist", "plugin", o.pluginName, "playlistID", info.ID, err)
|
|
// Schedule retry for transient errors if retryInterval is configured
|
|
if ri := time.Duration(o.retryInterval.Load()); ri > 0 {
|
|
o.schedulePlaylistRefresh(info, dbID, ownerID, ri)
|
|
}
|
|
return
|
|
}
|
|
|
|
// Convert SongRef → agents.Song and match against library
|
|
songs := songRefsToAgentSongs(resp.Tracks)
|
|
matched, err := o.matcher.MatchSongsToLibrary(ctx, songs, len(songs))
|
|
if err != nil {
|
|
log.Error(ctx, "Failed to match songs to library", "plugin", o.pluginName, "playlistID", info.ID, err)
|
|
return
|
|
}
|
|
|
|
// Build playlist model
|
|
pls := &model.Playlist{
|
|
ID: dbID,
|
|
Name: resp.Name,
|
|
Comment: resp.Description,
|
|
OwnerID: ownerID,
|
|
Public: false,
|
|
ExternalImageURL: resp.CoverArtURL,
|
|
PluginID: o.pluginName,
|
|
PluginPlaylistID: info.ID,
|
|
}
|
|
|
|
// Set tracks from matched media files
|
|
pls.AddMediaFiles(matched)
|
|
|
|
// Upsert via repository
|
|
plsRepo := o.ds.Playlist(ctx)
|
|
if err := plsRepo.Put(pls); err != nil {
|
|
log.Error(ctx, "Failed to upsert plugin playlist", "plugin", o.pluginName, "playlistID", info.ID, err)
|
|
return
|
|
}
|
|
|
|
log.Info(ctx, "Synced plugin playlist", "plugin", o.pluginName, "playlistID", info.ID,
|
|
"name", resp.Name, "tracks", len(matched), "owner", ownerID)
|
|
|
|
// Schedule refresh if ValidUntil > 0
|
|
if resp.ValidUntil > 0 {
|
|
validUntil := time.Unix(resp.ValidUntil, 0)
|
|
delay := time.Until(validUntil)
|
|
if delay <= 0 {
|
|
delay = 1 * time.Second // Already expired, refresh soon
|
|
}
|
|
o.schedulePlaylistRefresh(info, dbID, ownerID, delay)
|
|
}
|
|
}
|
|
|
|
func (o *playlistSyncer) schedulePlaylistRefresh(info capabilities.PlaylistInfo, dbID string, ownerID string, delay time.Duration) {
|
|
// Cancel existing timer if any
|
|
if timer, ok := o.refreshTimers[dbID]; ok {
|
|
timer.Stop()
|
|
}
|
|
o.refreshTimers[dbID] = time.AfterFunc(delay, func() {
|
|
select {
|
|
case o.workCh <- workItem{typ: workSync, info: info, dbID: dbID, ownerID: ownerID}:
|
|
case <-o.ctx.Done():
|
|
}
|
|
})
|
|
o.refreshTimerCount.Store(int32(len(o.refreshTimers)))
|
|
}
|
|
|
|
func (o *playlistSyncer) scheduleDiscovery(delay time.Duration) {
|
|
if o.discoveryTimer != nil {
|
|
o.discoveryTimer.Stop()
|
|
}
|
|
o.discoveryTimer = time.AfterFunc(delay, func() {
|
|
select {
|
|
case o.workCh <- workItem{typ: workDiscover}:
|
|
case <-o.ctx.Done():
|
|
}
|
|
})
|
|
}
|
|
|
|
// isPlaylistNotFoundError checks if the error contains a NotFound sentinel from the plugin.
|
|
func isPlaylistNotFoundError(err error) bool {
|
|
return err != nil && strings.Contains(err.Error(), capabilities.PlaylistProviderErrorNotFound.Error())
|
|
}
|
|
|
|
// stopAllTimers stops the discovery timer and all refresh timers.
|
|
func (o *playlistSyncer) stopAllTimers() {
|
|
if o.discoveryTimer != nil {
|
|
o.discoveryTimer.Stop()
|
|
}
|
|
for _, timer := range o.refreshTimers {
|
|
timer.Stop()
|
|
}
|
|
}
|
|
|
|
// Close cancels the context and waits for the worker goroutine to finish.
|
|
func (o *playlistSyncer) Close() error {
|
|
o.cancel()
|
|
<-o.done
|
|
return nil
|
|
}
|