navidrome/plugins/manager_watcher.go
Deluan Quintão b293b96256
refactor(persistence): stateless repositories with per-call context (#6149)
* refactor(persistence): adopt generic deluan/rest repository API

Pin deluan/rest to the refactor branch. REST-facing repository methods
take a context and return typed values. Drop DataStore.Resource and
ResourceRepository; the native API names typed repositories directly
through a per-request adapter that later commits remove.

* refactor(persistence): base repository helpers take a context

* refactor(persistence): LibraryRepository takes a context per call

* refactor(persistence): PropertyRepository takes a context per call

* refactor(persistence): UserPropsRepository takes a context per call

* refactor(persistence): TranscodingRepository takes a context per call

* refactor(persistence): ShareRepository takes a context per call

* refactor(persistence): PlayerRepository takes a context per call

* refactor(persistence): RadioRepository takes a context per call

* refactor(persistence): PlayQueueRepository takes a context per call

* refactor(persistence): Tag and Genre repositories take a context per call

* refactor(persistence): PluginRepository takes a context per call

* refactor(persistence): Scrobble repositories take a context per call

* refactor(persistence): FolderRepository takes a context per call

* refactor(persistence): Artwork repositories take a context per call

* refactor(persistence): UserRepository takes a context per call

* refactor(persistence): ArtistRepository takes a context per call

ReadAll no longer rewrites the shared sort mappings for the role filter;
it works on a per-call copy.

* test(persistence): assert artist role sort sanitization in ReadAll

* refactor(persistence): AlbumRepository takes a context per call

* test(persistence): pass the test context to album repository helpers

* refactor(persistence): MediaFileRepository takes a context per call

* refactor(persistence): Playlist repositories take a context per call

* refactor(persistence): build all repositories once per store

* refactor(core): REST repository wrappers are built once

* refactor(persistence): repositories are stateless

Remove the context field from the base repository and the per-request
REST adapter. Enable the containedctx linter so no repository can hold a
request context again.

* chore(lint): skip containedctx in test files

* refactor: share simplifications from the stateless repositories sweep

Add deleteOwnedAll on sqlRepository and use it in player/share Delete
to remove the duplicated bulk-delete loop; have Share.Repository()
return model.ShareRepository so subsonic sharing.go drops its repeated
type assertions.

* chore(core): assert REST wrappers implement Persistable

* chore: reformat imports

* perf(persistence): build repositories on first use

Each transaction store used to construct all 21 repositories up front,
paying for filter and sort mapping setup the block never touched. Fields
are now sync.OnceValue thunks, so a store only builds what it uses.

* fix(persistence): clean plugin references per deleted user

A bulk user delete that fails on a later id had already removed the
earlier rows but skipped their plugin cleanup. Cleanup now runs right
after each successful delete.

* fix(core): unload disabled plugins even when a user delete fails

A bulk delete can fail on a later id after earlier users were removed
and their plugins auto-disabled. The wrapper returned before unloading,
leaving those plugins running until the next successful delete or a
restart.

* chore(deps): pin deluan/rest to v1.0.1

Replaces the pseudo-version of the refactor branch with the tagged
release. REST error messages now name the bare type (Artist, not
model.Artist).

* test: use the spec context instead of context.Background()

Replace the context.Background()/context.TODO() calls this branch added
to tests with the spec's ctx, GinkgoT().Context(), or t/b.Context(), so
repository calls are bound to the running spec's lifetime.

* test: declare the spec context once per Describe

Set ctx from GinkgoT().Context() first in each top-level BeforeEach and reuse it, building user contexts on top of it instead of repeating inline calls.
2026-09-25 18:06:10 -04:00

221 lines
6.4 KiB
Go

package plugins
import (
"os"
"path/filepath"
"strings"
"sync"
"time"
"github.com/navidrome/navidrome/conf"
"github.com/navidrome/navidrome/log"
"github.com/rjeczalik/notify"
)
// debounceDuration is the time to wait before acting on file events
// to handle multiple rapid events for the same file.
const debounceDuration = 2 * time.Second
// startWatcher starts the file watcher for the plugins folder.
// It watches for CREATE, WRITE, and REMOVE events on .wasm files.
func (m *Manager) startWatcher() error {
folder := conf.Server.Plugins.Folder.String()
if folder == "" {
return nil
}
m.watcherEvents = make(chan notify.EventInfo, 10)
m.watcherDone = make(chan struct{})
m.debounceTimers = make(map[string]*time.Timer)
m.debounceMu = sync.Mutex{}
// Watch the plugins folder (not recursive)
// We filter for .wasm files in the event handler
if err := notify.Watch(folder, m.watcherEvents, notify.Create, notify.Write, notify.Remove, notify.Rename); err != nil {
close(m.watcherEvents)
return err
}
log.Info(m.ctx, "Started plugin file watcher", "folder", folder)
go m.watcherLoop()
return nil
}
// stopWatcher stops the file watcher
func (m *Manager) stopWatcher() {
if m.watcherEvents == nil {
return
}
notify.Stop(m.watcherEvents)
close(m.watcherDone)
// Cancel any pending debounce timers
m.debounceMu.Lock()
for _, timer := range m.debounceTimers {
timer.Stop()
}
m.debounceTimers = nil
m.debounceMu.Unlock()
log.Debug(m.ctx, "Stopped plugin file watcher")
}
// watcherLoop processes file watcher events
func (m *Manager) watcherLoop() {
for {
select {
case event, ok := <-m.watcherEvents:
if !ok {
return
}
m.handleWatcherEvent(event)
case <-m.ctx.Done():
return
case <-m.watcherDone:
return
}
}
}
// handleWatcherEvent processes a single file watcher event with debouncing
func (m *Manager) handleWatcherEvent(event notify.EventInfo) {
path := event.Path()
// Only process .ndp package files
if !strings.HasSuffix(path, PackageExtension) {
return
}
pluginName, ok := pluginIDFromPath(path)
if !ok {
log.Warn(m.ctx, "Ignoring plugin file with unusable name", "path", path)
return
}
log.Trace(m.ctx, "Plugin file event", "plugin", pluginName, "event", event.Event(), "path", path)
// Debounce: cancel any pending timer for this plugin and start a new one
m.debounceMu.Lock()
if timer, exists := m.debounceTimers[pluginName]; exists {
timer.Stop()
}
// Note: We don't capture the event type here. Instead, processPluginEvent
// checks if the file exists when the timer fires. This handles sequences like
// Remove+Create+Rename correctly by checking actual file state after debounce.
m.debounceTimers[pluginName] = time.AfterFunc(debounceDuration, func() {
m.processPluginEvent(pluginName)
})
m.debounceMu.Unlock()
}
// pluginAction represents the action to take on a plugin based on file state
type pluginAction int
const (
actionNone pluginAction = iota // No action needed
actionUpdate // File exists: add new or update existing plugin in DB
actionRemove // File gone: remove plugin from DB (unload if enabled)
)
// determinePluginAction decides what action to take based on file existence.
// We check file existence rather than relying on event type because:
// 1. Events can be coalesced on some systems (macOS FSEvents)
// 2. Rename events can mean either "renamed away" (remove) or "renamed to" (add)
// 3. Build tools often do atomic writes (write temp file, rename to target)
// By checking existence, we handle all these cases correctly.
func determinePluginAction(path string) pluginAction {
if _, err := os.Stat(path); err == nil {
// File exists - treat as add/update
return actionUpdate
}
// File doesn't exist - it was removed
return actionRemove
}
// processPluginEvent handles the actual plugin load/unload/reload after debouncing.
// - If file exists: extract manifest, add or update plugin in DB
// - If file gone: unload if enabled, delete from DB
func (m *Manager) processPluginEvent(pluginName string) {
// Don't process if manager is stopping/stopped (atomic check to avoid race with Stop())
if m.stopped.Load() {
return
}
// Clean up debounce timer entry
m.debounceMu.Lock()
delete(m.debounceTimers, pluginName)
m.debounceMu.Unlock()
folder := conf.Server.Plugins.Folder.String()
ndpPath := filepath.Join(folder, pluginName+PackageExtension)
action := determinePluginAction(ndpPath)
log.Debug(m.ctx, "Plugin event action", "plugin", pluginName, "action", action, "path", ndpPath)
ctx := adminContext(m.ctx)
repo := m.ds.Plugin()
switch action {
case actionUpdate:
// File changed - check SHA256 first, then extract manifest if needed
sha256Hash, err := ComputeFileSHA256(ndpPath)
if err != nil {
log.Error(m.ctx, "Failed to compute SHA256 for changed plugin", "plugin", pluginName, err)
return
}
dbPlugin, err := repo.Get(ctx, pluginName)
if err != nil {
// Plugin not in DB yet, need full manifest extraction to add it
metadata, extractErr := m.extractManifest(ndpPath)
if extractErr != nil {
log.Error(m.ctx, "Failed to extract manifest from new plugin", "plugin", pluginName, extractErr)
return
}
if addErr := m.addPluginToDB(ctx, repo, pluginName, ndpPath, metadata); addErr != nil {
log.Error(m.ctx, "Failed to add plugin to DB", "plugin", pluginName, addErr)
}
return
}
// Check if actually changed using lightweight SHA256 comparison
if dbPlugin.SHA256 == sha256Hash {
return // No actual change
}
// Plugin changed - now extract full manifest
metadata, err := m.extractManifest(ndpPath)
if err != nil {
log.Error(m.ctx, "Failed to extract manifest from changed plugin", "plugin", pluginName, err)
// Update error in DB
dbPlugin.LastError = err.Error()
dbPlugin.UpdatedAt = time.Now()
if dbPlugin.Enabled {
_ = m.unloadPlugin(pluginName)
dbPlugin.Enabled = false
}
_ = repo.Put(ctx, dbPlugin)
return
}
if err := m.updatePluginInDB(ctx, repo, dbPlugin, ndpPath, metadata); err != nil {
log.Error(m.ctx, "Failed to update plugin in DB", "plugin", pluginName, err)
}
case actionRemove:
// File removed - unload if enabled, delete from DB
dbPlugin, err := repo.Get(ctx, pluginName)
if err != nil {
log.Debug(m.ctx, "Removed plugin not in DB", "plugin", pluginName)
return
}
if err := m.removePluginFromDB(ctx, repo, dbPlugin); err != nil {
log.Error(m.ctx, "Failed to delete plugin from DB", "plugin", pluginName, err)
}
}
}