mirror of
https://github.com/navidrome/navidrome.git
synced 2026-08-01 07:21:17 +00:00
feat(plugins UI): implement plugin synchronization with database for add, update, and remove actions
Signed-off-by: Deluan <deluan@navidrome.org>
This commit is contained in:
parent
9c626183d0
commit
690785120a
@ -4,6 +4,8 @@ package plugins
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
@ -37,6 +39,10 @@ var _ = Describe("ArtworkService", Ordered, func() {
|
||||
err = os.WriteFile(destPath, data, 0600)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
// Compute SHA256 for the plugin
|
||||
hash := sha256.Sum256(data)
|
||||
hashHex := hex.EncodeToString(hash[:])
|
||||
|
||||
// Setup config
|
||||
DeferCleanup(configtest.SetupConfig())
|
||||
conf.Server.Plugins.Enabled = true
|
||||
@ -48,10 +54,25 @@ var _ = Describe("ArtworkService", Ordered, func() {
|
||||
ds := &tests.MockDataStore{MockedProperty: &tests.MockedPropertyRepo{}}
|
||||
auth.Init(ds)
|
||||
|
||||
// Setup mock DataStore with pre-enabled plugin
|
||||
mockPluginRepo := tests.CreateMockPluginRepo()
|
||||
mockPluginRepo.Permitted = true
|
||||
mockPluginRepo.SetData(model.Plugins{{
|
||||
ID: "test-artwork",
|
||||
Path: destPath,
|
||||
SHA256: hashHex,
|
||||
Enabled: true,
|
||||
}})
|
||||
dataStore := &tests.MockDataStore{
|
||||
MockedProperty: &tests.MockedPropertyRepo{},
|
||||
MockedPlugin: mockPluginRepo,
|
||||
}
|
||||
|
||||
// Create and start manager
|
||||
manager = &Manager{
|
||||
plugins: make(map[string]*plugin),
|
||||
}
|
||||
manager.SetDataStore(dataStore)
|
||||
err = manager.Start(GinkgoT().Context())
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
|
||||
@ -4,6 +4,8 @@ package plugins
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"os"
|
||||
@ -12,6 +14,8 @@ import (
|
||||
|
||||
"github.com/navidrome/navidrome/conf"
|
||||
"github.com/navidrome/navidrome/conf/configtest"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
"github.com/navidrome/navidrome/tests"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
)
|
||||
@ -331,6 +335,10 @@ var _ = Describe("CacheService Integration", Ordered, func() {
|
||||
err = os.WriteFile(destPath, data, 0600)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
// Compute SHA256 for the plugin
|
||||
hash := sha256.Sum256(data)
|
||||
hashHex := hex.EncodeToString(hash[:])
|
||||
|
||||
// Setup config
|
||||
DeferCleanup(configtest.SetupConfig())
|
||||
conf.Server.Plugins.Enabled = true
|
||||
@ -338,10 +346,22 @@ var _ = Describe("CacheService Integration", Ordered, func() {
|
||||
conf.Server.Plugins.AutoReload = false
|
||||
conf.Server.CacheFolder = filepath.Join(tmpDir, "cache")
|
||||
|
||||
// Setup mock DataStore with pre-enabled plugin
|
||||
mockPluginRepo := tests.CreateMockPluginRepo()
|
||||
mockPluginRepo.Permitted = true
|
||||
mockPluginRepo.SetData(model.Plugins{{
|
||||
ID: "test-cache-plugin",
|
||||
Path: destPath,
|
||||
SHA256: hashHex,
|
||||
Enabled: true,
|
||||
}})
|
||||
dataStore := &tests.MockDataStore{MockedPlugin: mockPluginRepo}
|
||||
|
||||
// Create and start manager
|
||||
manager = &Manager{
|
||||
plugins: make(map[string]*plugin),
|
||||
}
|
||||
manager.SetDataStore(dataStore)
|
||||
err = manager.Start(GinkgoT().Context())
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
|
||||
@ -4,6 +4,8 @@ package plugins
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
@ -11,7 +13,9 @@ import (
|
||||
|
||||
"github.com/navidrome/navidrome/conf"
|
||||
"github.com/navidrome/navidrome/conf/configtest"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
"github.com/navidrome/navidrome/scheduler"
|
||||
"github.com/navidrome/navidrome/tests"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
)
|
||||
@ -39,6 +43,10 @@ var _ = Describe("SchedulerService", Ordered, func() {
|
||||
err = os.WriteFile(destPath, data, 0600)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
// Compute SHA256 for the plugin
|
||||
hash := sha256.Sum256(data)
|
||||
hashHex := hex.EncodeToString(hash[:])
|
||||
|
||||
// Setup config
|
||||
DeferCleanup(configtest.SetupConfig())
|
||||
conf.Server.Plugins.Enabled = true
|
||||
@ -54,10 +62,22 @@ var _ = Describe("SchedulerService", Ordered, func() {
|
||||
origAfterFn = timeAfterFunc
|
||||
timeAfterFunc = mockTimers.AfterFunc
|
||||
|
||||
// Setup mock DataStore with pre-enabled plugin
|
||||
mockPluginRepo := tests.CreateMockPluginRepo()
|
||||
mockPluginRepo.Permitted = true
|
||||
mockPluginRepo.SetData(model.Plugins{{
|
||||
ID: "test-scheduler",
|
||||
Path: destPath,
|
||||
SHA256: hashHex,
|
||||
Enabled: true,
|
||||
}})
|
||||
dataStore := &tests.MockDataStore{MockedPlugin: mockPluginRepo}
|
||||
|
||||
// Create and start manager
|
||||
manager = &Manager{
|
||||
plugins: make(map[string]*plugin),
|
||||
}
|
||||
manager.SetDataStore(dataStore)
|
||||
err = manager.Start(GinkgoT().Context())
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
|
||||
@ -3,6 +3,8 @@
|
||||
package plugins
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"os"
|
||||
@ -69,6 +71,24 @@ var _ = Describe("SubsonicAPI Host Function", Ordered, func() {
|
||||
manager.SetSubsonicRouter(router)
|
||||
manager.SetDataStore(dataStore)
|
||||
|
||||
// Pre-enable the plugin in the mock repo so it loads on startup
|
||||
// Compute SHA256 of the plugin file to match what SyncPlugins will compute
|
||||
pluginPath := filepath.Join(tmpDir, "test-subsonicapi-plugin.wasm")
|
||||
wasmData, err := os.ReadFile(pluginPath)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
hash := sha256.Sum256(wasmData)
|
||||
hashHex := hex.EncodeToString(hash[:])
|
||||
|
||||
mockPluginRepo := dataStore.Plugin(GinkgoT().Context()).(*tests.MockPluginRepo)
|
||||
mockPluginRepo.Permitted = true
|
||||
enabledPlugin := model.Plugin{
|
||||
ID: "test-subsonicapi-plugin",
|
||||
Path: pluginPath,
|
||||
SHA256: hashHex,
|
||||
Enabled: true,
|
||||
}
|
||||
mockPluginRepo.SetData(model.Plugins{enabledPlugin})
|
||||
|
||||
// Start the manager
|
||||
err = manager.Start(GinkgoT().Context())
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
@ -4,7 +4,9 @@ package plugins
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/base64"
|
||||
"encoding/hex"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
@ -16,6 +18,8 @@ import (
|
||||
"github.com/gorilla/websocket"
|
||||
"github.com/navidrome/navidrome/conf"
|
||||
"github.com/navidrome/navidrome/conf/configtest"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
"github.com/navidrome/navidrome/tests"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
)
|
||||
@ -40,6 +44,10 @@ var _ = Describe("WebSocketService", Ordered, func() {
|
||||
err = os.WriteFile(destPath, data, 0600)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
// Compute SHA256 for the plugin
|
||||
hash := sha256.Sum256(data)
|
||||
hashHex := hex.EncodeToString(hash[:])
|
||||
|
||||
// Setup config
|
||||
DeferCleanup(configtest.SetupConfig())
|
||||
conf.Server.Plugins.Enabled = true
|
||||
@ -47,10 +55,22 @@ var _ = Describe("WebSocketService", Ordered, func() {
|
||||
conf.Server.Plugins.AutoReload = false
|
||||
conf.Server.CacheFolder = filepath.Join(tmpDir, "cache")
|
||||
|
||||
// Setup mock DataStore with pre-enabled plugin
|
||||
mockPluginRepo := tests.CreateMockPluginRepo()
|
||||
mockPluginRepo.Permitted = true
|
||||
mockPluginRepo.SetData(model.Plugins{{
|
||||
ID: "test-websocket",
|
||||
Path: destPath,
|
||||
SHA256: hashHex,
|
||||
Enabled: true,
|
||||
}})
|
||||
dataStore := &tests.MockDataStore{MockedPlugin: mockPluginRepo}
|
||||
|
||||
// Create and start manager
|
||||
manager = &Manager{
|
||||
plugins: make(map[string]*plugin),
|
||||
}
|
||||
manager.SetDataStore(dataStore)
|
||||
err = manager.Start(GinkgoT().Context())
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
|
||||
@ -4,6 +4,8 @@ import (
|
||||
"cmp"
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
@ -25,6 +27,7 @@ import (
|
||||
"github.com/navidrome/navidrome/core/scrobbler"
|
||||
"github.com/navidrome/navidrome/log"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
"github.com/navidrome/navidrome/model/request"
|
||||
"github.com/navidrome/navidrome/plugins/host"
|
||||
"github.com/navidrome/navidrome/scheduler"
|
||||
"github.com/navidrome/navidrome/utils/singleton"
|
||||
@ -113,6 +116,73 @@ func GetManager() *Manager {
|
||||
})
|
||||
}
|
||||
|
||||
// adminContext returns a context with admin privileges for DB operations.
|
||||
func adminContext(ctx context.Context) context.Context {
|
||||
return request.WithUser(ctx, model.User{IsAdmin: true})
|
||||
}
|
||||
|
||||
// marshalManifest marshals a manifest to JSON string, returning empty string on error.
|
||||
func marshalManifest(m *Manifest) string {
|
||||
b, _ := json.Marshal(m)
|
||||
return string(b)
|
||||
}
|
||||
|
||||
// addPluginToDB adds a new plugin to the database as disabled.
|
||||
func (m *Manager) addPluginToDB(ctx context.Context, repo model.PluginRepository, name, path string, metadata *PluginMetadata) error {
|
||||
now := time.Now()
|
||||
newPlugin := &model.Plugin{
|
||||
ID: name,
|
||||
Path: path,
|
||||
Manifest: marshalManifest(metadata.Manifest),
|
||||
SHA256: metadata.SHA256,
|
||||
Enabled: false,
|
||||
CreatedAt: now,
|
||||
UpdatedAt: now,
|
||||
}
|
||||
if err := repo.Put(newPlugin); err != nil {
|
||||
return fmt.Errorf("adding plugin to DB: %w", err)
|
||||
}
|
||||
log.Info(ctx, "Discovered new plugin", "plugin", name)
|
||||
return nil
|
||||
}
|
||||
|
||||
// updatePluginInDB updates an existing plugin in the database after a file change.
|
||||
// If the plugin was enabled, it will be unloaded and disabled.
|
||||
func (m *Manager) updatePluginInDB(ctx context.Context, repo model.PluginRepository, dbPlugin *model.Plugin, path string, metadata *PluginMetadata) error {
|
||||
wasEnabled := dbPlugin.Enabled
|
||||
if wasEnabled {
|
||||
if err := m.UnloadPlugin(dbPlugin.ID); err != nil {
|
||||
log.Debug(ctx, "Plugin not loaded during change", "plugin", dbPlugin.ID)
|
||||
}
|
||||
}
|
||||
dbPlugin.Path = path
|
||||
dbPlugin.Manifest = marshalManifest(metadata.Manifest)
|
||||
dbPlugin.SHA256 = metadata.SHA256
|
||||
dbPlugin.Enabled = false
|
||||
dbPlugin.LastError = ""
|
||||
dbPlugin.UpdatedAt = time.Now()
|
||||
if err := repo.Put(dbPlugin); err != nil {
|
||||
return fmt.Errorf("updating plugin in DB: %w", err)
|
||||
}
|
||||
log.Info(ctx, "Plugin file changed", "plugin", dbPlugin.ID, "wasEnabled", wasEnabled)
|
||||
return nil
|
||||
}
|
||||
|
||||
// removePluginFromDB removes a plugin from the database.
|
||||
// If the plugin was enabled, it will be unloaded first.
|
||||
func (m *Manager) removePluginFromDB(ctx context.Context, repo model.PluginRepository, dbPlugin *model.Plugin) error {
|
||||
if dbPlugin.Enabled {
|
||||
if err := m.UnloadPlugin(dbPlugin.ID); err != nil {
|
||||
log.Debug(ctx, "Plugin not loaded during removal", "plugin", dbPlugin.ID)
|
||||
}
|
||||
}
|
||||
if err := repo.Delete(dbPlugin.ID); err != nil {
|
||||
return fmt.Errorf("deleting plugin from DB: %w", err)
|
||||
}
|
||||
log.Info(ctx, "Plugin removed", "plugin", dbPlugin.ID)
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetSubsonicRouter sets the Subsonic router for SubsonicAPI host functions.
|
||||
// This should be called after the subsonic router is created but before plugins
|
||||
// that require SubsonicAPI access are loaded.
|
||||
@ -128,6 +198,9 @@ func (m *Manager) SetDataStore(ds model.DataStore) {
|
||||
|
||||
// Start initializes the plugin manager and loads plugins from the configured folder.
|
||||
// It should be called once during application startup when plugins are enabled.
|
||||
// The startup flow is:
|
||||
// 1. Sync plugins folder with DB (discover new, update changed, remove deleted)
|
||||
// 2. Load only enabled plugins from DB
|
||||
func (m *Manager) Start(ctx context.Context) error {
|
||||
if !conf.Server.Plugins.Enabled {
|
||||
log.Debug(ctx, "Plugin system is disabled")
|
||||
@ -168,10 +241,16 @@ func (m *Manager) Start(ctx context.Context) error {
|
||||
|
||||
log.Info(ctx, "Starting plugin manager", "folder", folder)
|
||||
|
||||
// Discover and load plugins
|
||||
if err := m.discoverPlugins(folder); err != nil {
|
||||
log.Error(ctx, "Error discovering plugins", err)
|
||||
return fmt.Errorf("discovering plugins: %w", err)
|
||||
// Sync plugins folder with DB
|
||||
if err := m.SyncPlugins(ctx, folder); err != nil {
|
||||
log.Error(ctx, "Error syncing plugins with DB", err)
|
||||
// Continue - we can still try to load plugins
|
||||
}
|
||||
|
||||
// Load enabled plugins from DB
|
||||
if err := m.loadEnabledPlugins(ctx); err != nil {
|
||||
log.Error(ctx, "Error loading enabled plugins", err)
|
||||
return fmt.Errorf("loading enabled plugins: %w", err)
|
||||
}
|
||||
|
||||
// Start file watcher if auto-reload is enabled
|
||||
@ -304,66 +383,281 @@ func (m *Manager) GetPluginInfo() map[string]PluginInfo {
|
||||
return info
|
||||
}
|
||||
|
||||
// discoverPlugins scans the plugins folder and loads all .wasm files in parallel
|
||||
func (m *Manager) discoverPlugins(folder string) error {
|
||||
// PluginMetadata holds the extracted information from a plugin file
|
||||
// without fully initializing the plugin.
|
||||
type PluginMetadata struct {
|
||||
Manifest *Manifest
|
||||
SHA256 string
|
||||
}
|
||||
|
||||
// compiledPluginInfo holds the intermediate compilation result used by both
|
||||
// ExtractManifest and loadPluginWithConfig.
|
||||
type compiledPluginInfo struct {
|
||||
wasmBytes []byte
|
||||
sha256 string
|
||||
manifest *Manifest
|
||||
compiled *extism.CompiledPlugin
|
||||
}
|
||||
|
||||
// stubHostFunctions returns the list of stub host functions needed for initial plugin compilation.
|
||||
func stubHostFunctions() []extism.HostFunction {
|
||||
stubs := host.RegisterSubsonicAPIHostFunctions(nil)
|
||||
stubs = append(stubs, host.RegisterSchedulerHostFunctions(nil)...)
|
||||
stubs = append(stubs, host.RegisterWebSocketHostFunctions(nil)...)
|
||||
stubs = append(stubs, host.RegisterArtworkHostFunctions(nil)...)
|
||||
stubs = append(stubs, host.RegisterCacheHostFunctions(nil)...)
|
||||
return stubs
|
||||
}
|
||||
|
||||
// compileAndExtractManifest reads a wasm file, compiles it with cache, and extracts the manifest.
|
||||
// The caller is responsible for closing the returned compiled plugin when done.
|
||||
func (m *Manager) compileAndExtractManifest(ctx context.Context, wasmPath string, config map[string]string) (*compiledPluginInfo, error) {
|
||||
wasmBytes, err := os.ReadFile(wasmPath)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("reading wasm file: %w", err)
|
||||
}
|
||||
|
||||
// Compute SHA-256 hash
|
||||
hash := sha256.Sum256(wasmBytes)
|
||||
hashHex := hex.EncodeToString(hash[:])
|
||||
|
||||
// Extract plugin name from path for logging
|
||||
pluginName := strings.TrimSuffix(filepath.Base(wasmPath), ".wasm")
|
||||
|
||||
pluginManifest := extism.Manifest{
|
||||
Wasm: []extism.Wasm{
|
||||
extism.WasmData{Data: wasmBytes, Name: "main"},
|
||||
},
|
||||
Config: config,
|
||||
Timeout: uint64(defaultTimeout.Milliseconds()),
|
||||
}
|
||||
extismConfig := extism.PluginConfig{
|
||||
EnableWasi: true,
|
||||
RuntimeConfig: wazero.NewRuntimeConfig().WithCompilationCache(m.cache),
|
||||
}
|
||||
|
||||
compiled, err := extism.NewCompiledPlugin(ctx, pluginManifest, extismConfig, stubHostFunctions())
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("compiling plugin: %w", err)
|
||||
}
|
||||
|
||||
instance, err := compiled.Instance(ctx, extism.PluginInstanceConfig{})
|
||||
if err != nil {
|
||||
compiled.Close(ctx)
|
||||
return nil, fmt.Errorf("creating instance: %w", err)
|
||||
}
|
||||
defer instance.Close(ctx)
|
||||
instance.SetLogger(extismLogger(pluginName))
|
||||
|
||||
exit, manifestBytes, err := instance.Call(manifestFunction, nil)
|
||||
if err != nil {
|
||||
compiled.Close(ctx)
|
||||
return nil, fmt.Errorf("calling manifest function: %w", err)
|
||||
}
|
||||
if exit != 0 {
|
||||
compiled.Close(ctx)
|
||||
return nil, fmt.Errorf("manifest function exited with code %d", exit)
|
||||
}
|
||||
|
||||
var manifest Manifest
|
||||
if err := json.Unmarshal(manifestBytes, &manifest); err != nil {
|
||||
compiled.Close(ctx)
|
||||
return nil, fmt.Errorf("parsing manifest: %w", err)
|
||||
}
|
||||
|
||||
return &compiledPluginInfo{
|
||||
wasmBytes: wasmBytes,
|
||||
sha256: hashHex,
|
||||
manifest: &manifest,
|
||||
compiled: compiled,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// ExtractManifest loads a wasm file, computes its SHA-256 hash, extracts the manifest,
|
||||
// and immediately closes without full plugin initialization.
|
||||
// This is a lightweight operation used for plugin discovery and change detection.
|
||||
// The compilation is cached to speed up subsequent EnablePlugin calls.
|
||||
func (m *Manager) ExtractManifest(wasmPath string) (*PluginMetadata, error) {
|
||||
if m.stopped.Load() {
|
||||
return nil, fmt.Errorf("manager is stopped")
|
||||
}
|
||||
|
||||
info, err := m.compileAndExtractManifest(context.Background(), wasmPath, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer info.compiled.Close(context.Background())
|
||||
|
||||
return &PluginMetadata{
|
||||
Manifest: info.manifest,
|
||||
SHA256: info.sha256,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// SyncPlugins scans the plugins folder and synchronizes with the database.
|
||||
// It handles new, changed, and removed plugins by comparing SHA-256 hashes.
|
||||
// - New plugins are added to DB as disabled
|
||||
// - Changed plugins are updated in DB and disabled if they were enabled
|
||||
// - Removed plugins are deleted from DB (after unloading if enabled)
|
||||
func (m *Manager) SyncPlugins(ctx context.Context, folder string) error {
|
||||
if m.ds == nil {
|
||||
return fmt.Errorf("datastore not configured")
|
||||
}
|
||||
|
||||
adminCtx := adminContext(ctx)
|
||||
|
||||
// Read current plugins from folder
|
||||
entries, err := os.ReadDir(folder)
|
||||
if err != nil {
|
||||
if os.IsNotExist(err) {
|
||||
log.Debug("Plugins folder does not exist", "folder", folder)
|
||||
log.Debug(ctx, "Plugins folder does not exist", "folder", folder)
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
return fmt.Errorf("reading plugins folder: %w", err)
|
||||
}
|
||||
|
||||
// Collect all plugin files to load
|
||||
type pluginFile struct {
|
||||
name string
|
||||
path string
|
||||
}
|
||||
var pluginFiles []pluginFile
|
||||
// Build map of files in folder
|
||||
filesOnDisk := make(map[string]string) // name -> path
|
||||
for _, entry := range entries {
|
||||
if entry.IsDir() || !strings.HasSuffix(entry.Name(), ".wasm") {
|
||||
continue
|
||||
}
|
||||
pluginFiles = append(pluginFiles, pluginFile{
|
||||
name: strings.TrimSuffix(entry.Name(), ".wasm"),
|
||||
path: filepath.Join(folder, entry.Name()),
|
||||
})
|
||||
name := strings.TrimSuffix(entry.Name(), ".wasm")
|
||||
filesOnDisk[name] = filepath.Join(folder, entry.Name())
|
||||
}
|
||||
|
||||
if len(pluginFiles) == 0 {
|
||||
log.Trace(m.ctx, "No plugins found", "folder", folder)
|
||||
return nil
|
||||
// Get all plugins from DB
|
||||
repo := m.ds.Plugin(adminCtx)
|
||||
dbPlugins, err := repo.GetAll()
|
||||
if err != nil {
|
||||
return fmt.Errorf("reading plugins from DB: %w", err)
|
||||
}
|
||||
pluginsInDB := make(map[string]*model.Plugin)
|
||||
for i := range dbPlugins {
|
||||
pluginsInDB[dbPlugins[i].ID] = &dbPlugins[i]
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
|
||||
// Process files on disk
|
||||
for name, path := range filesOnDisk {
|
||||
metadata, err := m.ExtractManifest(path)
|
||||
if err != nil {
|
||||
log.Error(ctx, "Failed to extract manifest from plugin", "plugin", name, "path", path, err)
|
||||
// Store error in DB if plugin exists
|
||||
if dbPlugin, exists := pluginsInDB[name]; exists {
|
||||
dbPlugin.LastError = err.Error()
|
||||
dbPlugin.UpdatedAt = now
|
||||
if dbPlugin.Enabled {
|
||||
// Unload broken plugin
|
||||
if unloadErr := m.UnloadPlugin(name); unloadErr != nil {
|
||||
log.Debug(ctx, "Plugin not loaded", "plugin", name)
|
||||
}
|
||||
dbPlugin.Enabled = false
|
||||
}
|
||||
if putErr := repo.Put(dbPlugin); putErr != nil {
|
||||
log.Error(ctx, "Failed to update plugin in DB", "plugin", name, err)
|
||||
}
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
dbPlugin, exists := pluginsInDB[name]
|
||||
if !exists {
|
||||
// New plugin - add to DB as disabled
|
||||
if err := m.addPluginToDB(ctx, repo, name, path, metadata); err != nil {
|
||||
log.Error(ctx, "Failed to add plugin to DB", "plugin", name, err)
|
||||
}
|
||||
} else if dbPlugin.SHA256 != metadata.SHA256 {
|
||||
// Plugin changed - update DB
|
||||
if err := m.updatePluginInDB(ctx, repo, dbPlugin, path, metadata); err != nil {
|
||||
log.Error(ctx, "Failed to update plugin in DB", "plugin", name, err)
|
||||
}
|
||||
} else {
|
||||
// Plugin unchanged - update path in case folder moved
|
||||
if dbPlugin.Path != path {
|
||||
dbPlugin.Path = path
|
||||
dbPlugin.UpdatedAt = now
|
||||
if err := repo.Put(dbPlugin); err != nil {
|
||||
log.Error(ctx, "Failed to update plugin path in DB", "plugin", name, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
// Mark as processed
|
||||
delete(pluginsInDB, name)
|
||||
}
|
||||
|
||||
// Remove plugins no longer on disk
|
||||
for _, dbPlugin := range pluginsInDB {
|
||||
if err := m.removePluginFromDB(ctx, repo, dbPlugin); err != nil {
|
||||
log.Error(ctx, "Failed to delete plugin from DB", "plugin", dbPlugin.ID, err)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// loadEnabledPlugins loads all enabled plugins from the database.
|
||||
func (m *Manager) loadEnabledPlugins(ctx context.Context) error {
|
||||
if m.ds == nil {
|
||||
return fmt.Errorf("datastore not configured")
|
||||
}
|
||||
|
||||
adminCtx := adminContext(ctx)
|
||||
repo := m.ds.Plugin(adminCtx)
|
||||
|
||||
plugins, err := repo.GetAll()
|
||||
if err != nil {
|
||||
return fmt.Errorf("reading plugins from DB: %w", err)
|
||||
}
|
||||
|
||||
g := errgroup.Group{}
|
||||
g.SetLimit(maxPluginLoadConcurrency)
|
||||
|
||||
for _, pf := range pluginFiles {
|
||||
for _, p := range plugins {
|
||||
if !p.Enabled {
|
||||
continue
|
||||
}
|
||||
|
||||
plugin := p // Capture for goroutine
|
||||
g.Go(func() error {
|
||||
start := time.Now()
|
||||
log.Debug(m.ctx, "Loading plugin", "plugin", pf.name, "path", pf.path)
|
||||
defer func() {
|
||||
log.Debug(m.ctx, "Finished loading plugin", "plugin", pf.name, "duration", time.Since(start))
|
||||
}()
|
||||
log.Debug(ctx, "Loading enabled plugin", "plugin", plugin.ID, "path", plugin.Path)
|
||||
|
||||
// Panic recovery to prevent one plugin from crashing the loading process
|
||||
// Panic recovery
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
log.Error(m.ctx, "Panic while loading plugin", "plugin", pf.name, "panic", r)
|
||||
log.Error(ctx, "Panic while loading plugin", "plugin", plugin.ID, "panic", r)
|
||||
}
|
||||
}()
|
||||
|
||||
if err := m.loadPlugin(pf.name, pf.path); err != nil {
|
||||
log.Error(m.ctx, "Failed to load plugin", "plugin", pf.name, "path", pf.path, err)
|
||||
if err := m.loadPluginWithConfig(plugin.ID, plugin.Path, plugin.Config); err != nil {
|
||||
// Store error in DB
|
||||
plugin.LastError = err.Error()
|
||||
plugin.Enabled = false
|
||||
plugin.UpdatedAt = time.Now()
|
||||
if putErr := repo.Put(&plugin); putErr != nil {
|
||||
log.Error(ctx, "Failed to update plugin error in DB", "plugin", plugin.ID, putErr)
|
||||
}
|
||||
log.Error(ctx, "Failed to load plugin", "plugin", plugin.ID, err)
|
||||
return nil
|
||||
}
|
||||
|
||||
// Clear any previous error
|
||||
if plugin.LastError != "" {
|
||||
plugin.LastError = ""
|
||||
plugin.UpdatedAt = time.Now()
|
||||
if putErr := repo.Put(&plugin); putErr != nil {
|
||||
log.Error(ctx, "Failed to clear plugin error in DB", "plugin", plugin.ID, putErr)
|
||||
}
|
||||
}
|
||||
|
||||
m.mu.RLock()
|
||||
p := m.plugins[pf.name]
|
||||
loadedPlugin := m.plugins[plugin.ID]
|
||||
m.mu.RUnlock()
|
||||
if p != nil {
|
||||
log.Info(m.ctx, "Loaded plugin", "plugin", pf.name, "manifest", p.manifest.Name, "capabilities", p.capabilities)
|
||||
if loadedPlugin != nil {
|
||||
log.Info(ctx, "Loaded plugin", "plugin", plugin.ID, "manifest", loadedPlugin.manifest.Name,
|
||||
"capabilities", loadedPlugin.capabilities, "duration", time.Since(start))
|
||||
}
|
||||
return nil
|
||||
})
|
||||
@ -372,96 +666,64 @@ func (m *Manager) discoverPlugins(folder string) error {
|
||||
return g.Wait()
|
||||
}
|
||||
|
||||
// loadPlugin loads a single plugin from a wasm file
|
||||
func (m *Manager) loadPlugin(name, wasmPath string) error {
|
||||
// loadPluginWithConfig loads a plugin with configuration from DB.
|
||||
func (m *Manager) loadPluginWithConfig(name, wasmPath, configJSON string) error {
|
||||
if m.stopped.Load() {
|
||||
return fmt.Errorf("manager is stopped")
|
||||
}
|
||||
|
||||
// Track this operation so Stop() can wait for it to complete
|
||||
// Track this operation
|
||||
m.loadWg.Add(1)
|
||||
defer m.loadWg.Done()
|
||||
|
||||
// Double-check after adding to WaitGroup (Stop may have been called between check and Add)
|
||||
if m.stopped.Load() {
|
||||
return fmt.Errorf("manager is stopped")
|
||||
}
|
||||
|
||||
wasmBytes, err := os.ReadFile(wasmPath)
|
||||
// Parse config from JSON
|
||||
var pluginConfig map[string]string
|
||||
if configJSON != "" {
|
||||
if err := json.Unmarshal([]byte(configJSON), &pluginConfig); err != nil {
|
||||
return fmt.Errorf("parsing plugin config: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Compile and extract manifest using shared helper
|
||||
info, err := m.compileAndExtractManifest(m.ctx, wasmPath, pluginConfig)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
pluginConfig := m.getPluginConfig(name)
|
||||
// Create instance to detect capabilities
|
||||
instance, err := info.compiled.Instance(m.ctx, extism.PluginInstanceConfig{})
|
||||
if err != nil {
|
||||
info.compiled.Close(m.ctx)
|
||||
return fmt.Errorf("creating instance: %w", err)
|
||||
}
|
||||
instance.SetLogger(extismLogger(name))
|
||||
capabilities := detectCapabilities(instance)
|
||||
instance.Close(m.ctx)
|
||||
|
||||
// Build host functions based on permissions
|
||||
var hostFunctions []extism.HostFunction
|
||||
var closers []io.Closer
|
||||
|
||||
// Build extism manifest for potential recompilation
|
||||
pluginManifest := extism.Manifest{
|
||||
Wasm: []extism.Wasm{
|
||||
extism.WasmData{Data: wasmBytes, Name: "main"},
|
||||
extism.WasmData{Data: info.wasmBytes, Name: "main"},
|
||||
},
|
||||
Config: pluginConfig,
|
||||
Timeout: uint64(defaultTimeout.Milliseconds()),
|
||||
}
|
||||
extismConfig := extism.PluginConfig{
|
||||
EnableWasi: true,
|
||||
RuntimeConfig: wazero.NewRuntimeConfig().WithCompilationCache(m.cache),
|
||||
}
|
||||
|
||||
// Register stub host functions for initial compilation.
|
||||
// This is necessary because plugins that import host functions will fail to compile if those
|
||||
// functions aren't available at compile time.
|
||||
// The real service will be registered during recompilation.
|
||||
stubHostFunctions := host.RegisterSubsonicAPIHostFunctions(nil)
|
||||
stubHostFunctions = append(stubHostFunctions, host.RegisterSchedulerHostFunctions(nil)...)
|
||||
stubHostFunctions = append(stubHostFunctions, host.RegisterWebSocketHostFunctions(nil)...)
|
||||
stubHostFunctions = append(stubHostFunctions, host.RegisterArtworkHostFunctions(nil)...)
|
||||
stubHostFunctions = append(stubHostFunctions, host.RegisterCacheHostFunctions(nil)...)
|
||||
|
||||
// Create initial compiled plugin with stub host functions
|
||||
compiled, err := extism.NewCompiledPlugin(m.ctx, pluginManifest, extismConfig, stubHostFunctions)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Create instance to read manifest and detect capabilities
|
||||
instance, err := compiled.Instance(m.ctx, extism.PluginInstanceConfig{})
|
||||
if err != nil {
|
||||
compiled.Close(m.ctx)
|
||||
return err
|
||||
}
|
||||
instance.SetLogger(extismLogger(name))
|
||||
|
||||
exit, manifestBytes, err := instance.Call(manifestFunction, nil)
|
||||
if err != nil {
|
||||
instance.Close(m.ctx)
|
||||
compiled.Close(m.ctx)
|
||||
return err
|
||||
}
|
||||
if exit != 0 {
|
||||
instance.Close(m.ctx)
|
||||
compiled.Close(m.ctx)
|
||||
return fmt.Errorf("calling %s: %d", manifestFunction, exit)
|
||||
}
|
||||
|
||||
var manifest Manifest
|
||||
if err := json.Unmarshal(manifestBytes, &manifest); err != nil {
|
||||
instance.Close(m.ctx)
|
||||
compiled.Close(m.ctx)
|
||||
return fmt.Errorf("invalid plugin manifest: %w", err)
|
||||
}
|
||||
|
||||
// Detect capabilities using the instance before closing it
|
||||
capabilities := detectCapabilities(instance)
|
||||
instance.Close(m.ctx)
|
||||
|
||||
var hostFunctions []extism.HostFunction
|
||||
var closers []io.Closer
|
||||
|
||||
if hosts := manifest.AllowedHosts(); len(hosts) > 0 {
|
||||
if hosts := info.manifest.AllowedHosts(); len(hosts) > 0 {
|
||||
pluginManifest.AllowedHosts = hosts
|
||||
}
|
||||
|
||||
// Register SubsonicAPI host functions if permission is granted
|
||||
if manifest.Permissions != nil && manifest.Permissions.Subsonicapi != nil {
|
||||
perm := manifest.Permissions.Subsonicapi
|
||||
if info.manifest.Permissions != nil && info.manifest.Permissions.Subsonicapi != nil {
|
||||
perm := info.manifest.Permissions.Subsonicapi
|
||||
if m.subsonicRouter != nil && m.ds != nil {
|
||||
service := newSubsonicAPIService(name, m.subsonicRouter, m.ds, perm)
|
||||
hostFunctions = append(hostFunctions, host.RegisterSubsonicAPIHostFunctions(service)...)
|
||||
@ -471,41 +733,44 @@ func (m *Manager) loadPlugin(name, wasmPath string) error {
|
||||
}
|
||||
|
||||
// Register Scheduler host functions if permission is granted
|
||||
if manifest.Permissions != nil && manifest.Permissions.Scheduler != nil {
|
||||
if info.manifest.Permissions != nil && info.manifest.Permissions.Scheduler != nil {
|
||||
service := newSchedulerService(name, m, scheduler.GetInstance())
|
||||
closers = append(closers, service)
|
||||
hostFunctions = append(hostFunctions, host.RegisterSchedulerHostFunctions(service)...)
|
||||
}
|
||||
|
||||
// Register WebSocket host functions if permission is granted
|
||||
if manifest.Permissions != nil && manifest.Permissions.Websocket != nil {
|
||||
perm := manifest.Permissions.Websocket
|
||||
if info.manifest.Permissions != nil && info.manifest.Permissions.Websocket != nil {
|
||||
perm := info.manifest.Permissions.Websocket
|
||||
service := newWebSocketService(name, m, perm.AllowedHosts)
|
||||
closers = append(closers, service)
|
||||
hostFunctions = append(hostFunctions, host.RegisterWebSocketHostFunctions(service)...)
|
||||
}
|
||||
|
||||
// Register Artwork host functions if permission is granted
|
||||
if manifest.Permissions != nil && manifest.Permissions.Artwork != nil {
|
||||
if info.manifest.Permissions != nil && info.manifest.Permissions.Artwork != nil {
|
||||
service := newArtworkService()
|
||||
hostFunctions = append(hostFunctions, host.RegisterArtworkHostFunctions(service)...)
|
||||
}
|
||||
|
||||
// Register Cache host functions if permission is granted
|
||||
if manifest.Permissions != nil && manifest.Permissions.Cache != nil {
|
||||
if info.manifest.Permissions != nil && info.manifest.Permissions.Cache != nil {
|
||||
service := newCacheService(name)
|
||||
closers = append(closers, service)
|
||||
hostFunctions = append(hostFunctions, host.RegisterCacheHostFunctions(service)...)
|
||||
}
|
||||
|
||||
// Check if the plugin needs to be recompiled with real host functions
|
||||
compiled := info.compiled
|
||||
needsRecompile := len(pluginManifest.AllowedHosts) > 0 || len(hostFunctions) > 0
|
||||
|
||||
// Recompile if needed. It is actually not a "recompile" since the first compilation
|
||||
// should be cached by wazero. We just need to do it this way to provide the real host functions.
|
||||
if needsRecompile {
|
||||
log.Trace(m.ctx, "Recompiling plugin", "plugin", name)
|
||||
compiled.Close(m.ctx)
|
||||
log.Trace(m.ctx, "Recompiling plugin with host functions", "plugin", name)
|
||||
info.compiled.Close(m.ctx)
|
||||
extismConfig := extism.PluginConfig{
|
||||
EnableWasi: true,
|
||||
RuntimeConfig: wazero.NewRuntimeConfig().WithCompilationCache(m.cache),
|
||||
}
|
||||
compiled, err = extism.NewCompiledPlugin(m.ctx, pluginManifest, extismConfig, hostFunctions)
|
||||
if err != nil {
|
||||
return err
|
||||
@ -516,25 +781,135 @@ func (m *Manager) loadPlugin(name, wasmPath string) error {
|
||||
m.plugins[name] = &plugin{
|
||||
name: name,
|
||||
path: wasmPath,
|
||||
manifest: &manifest,
|
||||
manifest: info.manifest,
|
||||
compiled: compiled,
|
||||
capabilities: capabilities,
|
||||
closers: closers,
|
||||
}
|
||||
m.mu.Unlock()
|
||||
|
||||
// Call plugin init function if the plugin has the Lifecycle capability
|
||||
// Call plugin init function
|
||||
callPluginInit(m.ctx, m.plugins[name])
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// getPluginConfig returns the configuration for a specific plugin
|
||||
func (m *Manager) getPluginConfig(name string) map[string]string {
|
||||
if conf.Server.PluginConfig == nil {
|
||||
return nil
|
||||
// EnablePlugin enables a plugin by loading it and updating the DB.
|
||||
// Returns an error if the plugin is not found in DB or fails to load.
|
||||
func (m *Manager) EnablePlugin(ctx context.Context, id string) error {
|
||||
if m.ds == nil {
|
||||
return fmt.Errorf("datastore not configured")
|
||||
}
|
||||
return conf.Server.PluginConfig[name]
|
||||
|
||||
adminCtx := adminContext(ctx)
|
||||
repo := m.ds.Plugin(adminCtx)
|
||||
|
||||
plugin, err := repo.Get(id)
|
||||
if err != nil {
|
||||
return fmt.Errorf("getting plugin from DB: %w", err)
|
||||
}
|
||||
|
||||
if plugin.Enabled {
|
||||
return nil // Already enabled
|
||||
}
|
||||
|
||||
// Try to load the plugin
|
||||
if err := m.loadPluginWithConfig(plugin.ID, plugin.Path, plugin.Config); err != nil {
|
||||
// Store error and return
|
||||
plugin.LastError = err.Error()
|
||||
plugin.UpdatedAt = time.Now()
|
||||
_ = repo.Put(plugin)
|
||||
return fmt.Errorf("loading plugin: %w", err)
|
||||
}
|
||||
|
||||
// Update DB
|
||||
plugin.Enabled = true
|
||||
plugin.LastError = ""
|
||||
plugin.UpdatedAt = time.Now()
|
||||
if err := repo.Put(plugin); err != nil {
|
||||
// Unload since we couldn't update DB
|
||||
_ = m.UnloadPlugin(id)
|
||||
return fmt.Errorf("updating plugin in DB: %w", err)
|
||||
}
|
||||
|
||||
log.Info(ctx, "Enabled plugin", "plugin", id)
|
||||
return nil
|
||||
}
|
||||
|
||||
// DisablePlugin disables a plugin by unloading it and updating the DB.
|
||||
// Returns an error if the plugin is not found in DB.
|
||||
func (m *Manager) DisablePlugin(ctx context.Context, id string) error {
|
||||
if m.ds == nil {
|
||||
return fmt.Errorf("datastore not configured")
|
||||
}
|
||||
|
||||
adminCtx := adminContext(ctx)
|
||||
repo := m.ds.Plugin(adminCtx)
|
||||
|
||||
plugin, err := repo.Get(id)
|
||||
if err != nil {
|
||||
return fmt.Errorf("getting plugin from DB: %w", err)
|
||||
}
|
||||
|
||||
if !plugin.Enabled {
|
||||
return nil // Already disabled
|
||||
}
|
||||
|
||||
// Unload the plugin
|
||||
if err := m.UnloadPlugin(id); err != nil {
|
||||
log.Debug(ctx, "Plugin was not loaded", "plugin", id)
|
||||
}
|
||||
|
||||
// Update DB
|
||||
plugin.Enabled = false
|
||||
plugin.UpdatedAt = time.Now()
|
||||
if err := repo.Put(plugin); err != nil {
|
||||
return fmt.Errorf("updating plugin in DB: %w", err)
|
||||
}
|
||||
|
||||
log.Info(ctx, "Disabled plugin", "plugin", id)
|
||||
return nil
|
||||
}
|
||||
|
||||
// UpdatePluginConfig updates the configuration for a plugin.
|
||||
// If the plugin is enabled, it will be reloaded with the new config.
|
||||
func (m *Manager) UpdatePluginConfig(ctx context.Context, id, configJSON string) error {
|
||||
if m.ds == nil {
|
||||
return fmt.Errorf("datastore not configured")
|
||||
}
|
||||
|
||||
adminCtx := adminContext(ctx)
|
||||
repo := m.ds.Plugin(adminCtx)
|
||||
|
||||
plugin, err := repo.Get(id)
|
||||
if err != nil {
|
||||
return fmt.Errorf("getting plugin from DB: %w", err)
|
||||
}
|
||||
|
||||
wasEnabled := plugin.Enabled
|
||||
|
||||
// Update config in DB
|
||||
plugin.Config = configJSON
|
||||
plugin.UpdatedAt = time.Now()
|
||||
if err := repo.Put(plugin); err != nil {
|
||||
return fmt.Errorf("updating plugin config in DB: %w", err)
|
||||
}
|
||||
|
||||
// Reload if enabled
|
||||
if wasEnabled {
|
||||
if err := m.UnloadPlugin(id); err != nil {
|
||||
log.Debug(ctx, "Plugin was not loaded", "plugin", id)
|
||||
}
|
||||
if err := m.loadPluginWithConfig(plugin.ID, plugin.Path, configJSON); err != nil {
|
||||
plugin.LastError = err.Error()
|
||||
plugin.Enabled = false
|
||||
_ = repo.Put(plugin)
|
||||
return fmt.Errorf("reloading plugin with new config: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
log.Info(ctx, "Updated plugin config", "plugin", id)
|
||||
return nil
|
||||
}
|
||||
|
||||
// UnloadPlugin removes a plugin from the manager and closes its resources.
|
||||
@ -570,54 +945,6 @@ func (m *Manager) UnloadPlugin(name string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// LoadPlugin loads a new plugin by name from the plugins folder.
|
||||
// The plugin file must exist at <plugins_folder>/<name>.wasm.
|
||||
// Returns an error if the plugin is already loaded or fails to load.
|
||||
func (m *Manager) LoadPlugin(name string) error {
|
||||
m.mu.RLock()
|
||||
_, exists := m.plugins[name]
|
||||
m.mu.RUnlock()
|
||||
|
||||
if exists {
|
||||
return fmt.Errorf("plugin %q is already loaded", name)
|
||||
}
|
||||
|
||||
folder := conf.Server.Plugins.Folder
|
||||
if folder == "" {
|
||||
return fmt.Errorf("no plugins folder configured")
|
||||
}
|
||||
|
||||
wasmPath := filepath.Join(folder, name+".wasm")
|
||||
if _, err := os.Stat(wasmPath); err != nil {
|
||||
if os.IsNotExist(err) {
|
||||
return fmt.Errorf("plugin file not found: %s", wasmPath)
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
if err := m.loadPlugin(name, wasmPath); err != nil {
|
||||
return fmt.Errorf("failed to load plugin %q: %w", name, err)
|
||||
}
|
||||
|
||||
log.Info(m.ctx, "Loaded plugin", "plugin", name)
|
||||
return nil
|
||||
}
|
||||
|
||||
// ReloadPlugin unloads and reloads a plugin by name.
|
||||
// If the plugin was loaded and unload succeeds but reload fails,
|
||||
// the plugin remains unloaded and the error is returned.
|
||||
func (m *Manager) ReloadPlugin(name string) error {
|
||||
if err := m.UnloadPlugin(name); err != nil {
|
||||
return fmt.Errorf("failed to unload plugin %q: %w", name, err)
|
||||
}
|
||||
|
||||
if err := m.LoadPlugin(name); err != nil {
|
||||
log.Error(m.ctx, "Failed to reload plugin, plugin remains unloaded", "plugin", name, err)
|
||||
return fmt.Errorf("failed to reload plugin %q: %w", name, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
var errFunctionNotFound = errors.New("function not found")
|
||||
|
||||
// callPluginFunction is a helper to call a plugin function with input and output types.
|
||||
|
||||
@ -9,7 +9,6 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/dustin/go-humanize"
|
||||
"github.com/navidrome/navidrome/conf"
|
||||
"github.com/navidrome/navidrome/core/agents"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
@ -18,56 +17,16 @@ import (
|
||||
var _ = Describe("Manager", Ordered, func() {
|
||||
var ctx context.Context
|
||||
|
||||
// Ensure plugin is loaded at the start (might have been unloaded by previous tests)
|
||||
BeforeAll(func() {
|
||||
ctx = GinkgoT().Context()
|
||||
if _, ok := testManager.plugins["test-metadata-agent"]; !ok {
|
||||
err := testManager.LoadPlugin("test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
}
|
||||
})
|
||||
|
||||
// Ensure plugin is restored after all tests in this block
|
||||
AfterAll(func() {
|
||||
if _, ok := testManager.plugins["test-metadata-agent"]; !ok {
|
||||
_ = testManager.LoadPlugin("test-metadata-agent")
|
||||
}
|
||||
})
|
||||
|
||||
Describe("LoadPlugin", func() {
|
||||
It("auto-loads plugins from folder on Start", func() {
|
||||
// Plugin is already loaded by testManager.Start() via discoverPlugins
|
||||
Describe("Plugin Loading", func() {
|
||||
It("loads enabled plugins from DB on Start", func() {
|
||||
// Plugin is already loaded by testManager.Start() via loadEnabledPlugins
|
||||
names := testManager.PluginNames(string(CapabilityMetadataAgent))
|
||||
Expect(names).To(ContainElement("test-metadata-agent"))
|
||||
})
|
||||
|
||||
It("returns error when plugin file does not exist", func() {
|
||||
err := testManager.LoadPlugin("nonexistent")
|
||||
Expect(err).To(HaveOccurred())
|
||||
Expect(err.Error()).To(ContainSubstring("plugin file not found"))
|
||||
})
|
||||
|
||||
It("returns error when plugin is already loaded", func() {
|
||||
// Plugin was loaded on Start, try to load again
|
||||
err := testManager.LoadPlugin("test-metadata-agent")
|
||||
Expect(err).To(HaveOccurred())
|
||||
Expect(err.Error()).To(ContainSubstring("already loaded"))
|
||||
})
|
||||
|
||||
It("returns error when plugins folder is not configured", func() {
|
||||
originalFolder := conf.Server.Plugins.Folder
|
||||
originalDataFolder := conf.Server.DataFolder
|
||||
conf.Server.Plugins.Folder = ""
|
||||
conf.Server.DataFolder = ""
|
||||
defer func() {
|
||||
conf.Server.Plugins.Folder = originalFolder
|
||||
conf.Server.DataFolder = originalDataFolder
|
||||
}()
|
||||
|
||||
err := testManager.LoadPlugin("test")
|
||||
Expect(err).To(HaveOccurred())
|
||||
Expect(err.Error()).To(ContainSubstring("no plugins folder configured"))
|
||||
})
|
||||
})
|
||||
|
||||
Describe("UnloadPlugin", func() {
|
||||
@ -80,15 +39,6 @@ var _ = Describe("Manager", Ordered, func() {
|
||||
Expect(names).ToNot(ContainElement("test-metadata-agent"))
|
||||
})
|
||||
|
||||
It("can reload after unload", func() {
|
||||
// Reload the plugin we just unloaded
|
||||
err := testManager.LoadPlugin("test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
names := testManager.PluginNames(string(CapabilityMetadataAgent))
|
||||
Expect(names).To(ContainElement("test-metadata-agent"))
|
||||
})
|
||||
|
||||
It("returns error when plugin not found", func() {
|
||||
err := testManager.UnloadPlugin("nonexistent")
|
||||
Expect(err).To(HaveOccurred())
|
||||
@ -96,23 +46,41 @@ var _ = Describe("Manager", Ordered, func() {
|
||||
})
|
||||
})
|
||||
|
||||
Describe("ReloadPlugin", func() {
|
||||
It("unloads and reloads a plugin", func() {
|
||||
err := testManager.ReloadPlugin("test-metadata-agent")
|
||||
Describe("EnablePlugin", func() {
|
||||
It("enables and loads a disabled plugin", func() {
|
||||
// First disable the plugin (which also unloads it)
|
||||
err := testManager.DisablePlugin(ctx, "test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(testManager.PluginNames(string(CapabilityMetadataAgent))).ToNot(ContainElement("test-metadata-agent"))
|
||||
|
||||
// Enable it
|
||||
err = testManager.EnablePlugin(ctx, "test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
names := testManager.PluginNames(string(CapabilityMetadataAgent))
|
||||
Expect(names).To(ContainElement("test-metadata-agent"))
|
||||
})
|
||||
})
|
||||
|
||||
It("returns error when plugin not found", func() {
|
||||
err := testManager.ReloadPlugin("nonexistent")
|
||||
Expect(err).To(HaveOccurred())
|
||||
Expect(err.Error()).To(ContainSubstring("failed to unload"))
|
||||
Describe("DisablePlugin", func() {
|
||||
It("disables and unloads an enabled plugin", func() {
|
||||
// Ensure the plugin is loaded first
|
||||
_ = testManager.EnablePlugin(ctx, "test-metadata-agent")
|
||||
|
||||
err := testManager.DisablePlugin(ctx, "test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
names := testManager.PluginNames(string(CapabilityMetadataAgent))
|
||||
Expect(names).ToNot(ContainElement("test-metadata-agent"))
|
||||
})
|
||||
})
|
||||
|
||||
Describe("GetPluginInfo", func() {
|
||||
BeforeEach(func() {
|
||||
// Ensure plugin is loaded for this test
|
||||
_ = testManager.EnablePlugin(ctx, "test-metadata-agent")
|
||||
})
|
||||
|
||||
It("returns information about all loaded plugins", func() {
|
||||
info := testManager.GetPluginInfo()
|
||||
Expect(info).To(HaveKey("test-metadata-agent"))
|
||||
@ -122,7 +90,8 @@ var _ = Describe("Manager", Ordered, func() {
|
||||
})
|
||||
|
||||
It("can call the plugin concurrently", func() {
|
||||
// Plugin is already loaded
|
||||
// Ensure plugin is loaded
|
||||
_ = testManager.EnablePlugin(ctx, "test-metadata-agent")
|
||||
|
||||
const concurrency = 30
|
||||
errs := make(chan error, concurrency)
|
||||
|
||||
@ -3,6 +3,9 @@
|
||||
package plugins
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
@ -12,6 +15,7 @@ import (
|
||||
"github.com/navidrome/navidrome/conf"
|
||||
"github.com/navidrome/navidrome/conf/configtest"
|
||||
"github.com/navidrome/navidrome/log"
|
||||
"github.com/navidrome/navidrome/model"
|
||||
"github.com/navidrome/navidrome/tests"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
@ -60,7 +64,8 @@ func createTestManagerWithPlugins(pluginConfig map[string]map[string]string, plu
|
||||
tmpDir, err := os.MkdirTemp("", "plugins-test-*")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
// Copy test plugins to temp dir
|
||||
// Copy test plugins to temp dir and build plugin list with SHA256
|
||||
var enabledPlugins model.Plugins
|
||||
for _, plugin := range plugins {
|
||||
srcPath := filepath.Join(testdataDir, plugin)
|
||||
destPath := filepath.Join(tmpDir, plugin)
|
||||
@ -68,6 +73,28 @@ func createTestManagerWithPlugins(pluginConfig map[string]map[string]string, plu
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
err = os.WriteFile(destPath, data, 0600)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
// Compute SHA256 for the plugin
|
||||
hash := sha256.Sum256(data)
|
||||
hashHex := hex.EncodeToString(hash[:])
|
||||
pluginName := plugin[:len(plugin)-5] // Remove .wasm extension
|
||||
|
||||
// Build config JSON if provided
|
||||
configJSON := ""
|
||||
if pluginConfig != nil && pluginConfig[pluginName] != nil {
|
||||
// Encode config to JSON
|
||||
configBytes, err := json.Marshal(pluginConfig[pluginName])
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
configJSON = string(configBytes)
|
||||
}
|
||||
|
||||
enabledPlugins = append(enabledPlugins, model.Plugin{
|
||||
ID: pluginName,
|
||||
Path: destPath,
|
||||
SHA256: hashHex,
|
||||
Enabled: true,
|
||||
Config: configJSON,
|
||||
})
|
||||
}
|
||||
|
||||
// Setup config
|
||||
@ -78,10 +105,17 @@ func createTestManagerWithPlugins(pluginConfig map[string]map[string]string, plu
|
||||
conf.Server.CacheFolder = filepath.Join(tmpDir, "cache")
|
||||
conf.Server.PluginConfig = pluginConfig
|
||||
|
||||
// Setup mock DataStore with pre-enabled plugins
|
||||
mockPluginRepo := tests.CreateMockPluginRepo()
|
||||
mockPluginRepo.Permitted = true
|
||||
mockPluginRepo.SetData(enabledPlugins)
|
||||
dataStore := &tests.MockDataStore{MockedPlugin: mockPluginRepo}
|
||||
|
||||
// Create and start manager
|
||||
manager := &Manager{
|
||||
plugins: make(map[string]*plugin),
|
||||
}
|
||||
manager.SetDataStore(dataStore)
|
||||
err = manager.Start(GinkgoT().Context())
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
|
||||
@ -110,34 +110,28 @@ type pluginAction int
|
||||
|
||||
const (
|
||||
actionNone pluginAction = iota // No action needed
|
||||
actionLoad // Load the plugin
|
||||
actionUnload // Unload the plugin
|
||||
actionReload // Reload the plugin
|
||||
actionAdd // Add new plugin to DB (disabled)
|
||||
actionUpdate // Update existing plugin in DB (disable if enabled)
|
||||
actionRemove // Remove plugin from DB (unload if enabled)
|
||||
)
|
||||
|
||||
// determinePluginAction decides what action to take based on the file event type
|
||||
// and whether the plugin is currently loaded. This is a pure function with no side effects.
|
||||
func determinePluginAction(eventType notify.Event, isLoaded bool) pluginAction {
|
||||
// determinePluginAction decides what action to take based on the file event type.
|
||||
func determinePluginAction(eventType notify.Event) pluginAction {
|
||||
switch {
|
||||
case eventType¬ify.Remove != 0 || eventType¬ify.Rename != 0:
|
||||
// File removed or renamed away - unload if loaded
|
||||
return actionUnload
|
||||
|
||||
return actionRemove
|
||||
case eventType¬ify.Create != 0:
|
||||
// New file - load it
|
||||
return actionLoad
|
||||
|
||||
return actionAdd
|
||||
case eventType¬ify.Write != 0:
|
||||
// File modified - reload if loaded, otherwise load
|
||||
if isLoaded {
|
||||
return actionReload
|
||||
}
|
||||
return actionLoad
|
||||
return actionUpdate
|
||||
}
|
||||
return actionNone
|
||||
}
|
||||
|
||||
// processPluginEvent handles the actual plugin load/unload/reload after debouncing
|
||||
// processPluginEvent handles the actual plugin load/unload/reload after debouncing.
|
||||
// - On file add: extract manifest, create DB record as disabled
|
||||
// - On file change: extract manifest, update DB, disable if was enabled
|
||||
// - On file remove: unload if enabled, delete DB record
|
||||
func (m *Manager) processPluginEvent(pluginName string, eventType notify.Event) {
|
||||
// Don't process if manager is stopping/stopped (atomic check to avoid race with Stop())
|
||||
if m.stopped.Load() {
|
||||
@ -149,29 +143,72 @@ func (m *Manager) processPluginEvent(pluginName string, eventType notify.Event)
|
||||
delete(m.debounceTimers, pluginName)
|
||||
m.debounceMu.Unlock()
|
||||
|
||||
// Check if plugin is currently loaded
|
||||
m.mu.RLock()
|
||||
_, isLoaded := m.plugins[pluginName]
|
||||
m.mu.RUnlock()
|
||||
action := determinePluginAction(eventType)
|
||||
log.Debug(m.ctx, "Plugin event action (DB mode)", "plugin", pluginName, "action", action)
|
||||
|
||||
ctx := adminContext(m.ctx)
|
||||
repo := m.ds.Plugin(ctx)
|
||||
folder := conf.Server.Plugins.Folder
|
||||
wasmPath := filepath.Join(folder, pluginName+".wasm")
|
||||
|
||||
// Determine and execute the appropriate action
|
||||
action := determinePluginAction(eventType, isLoaded)
|
||||
log.Debug("Plugin event action", "plugin", pluginName, "action", action)
|
||||
switch action {
|
||||
case actionLoad:
|
||||
log.Debug("Loading new Plugin", "plugin", pluginName)
|
||||
if err := m.LoadPlugin(pluginName); err != nil {
|
||||
log.Error(m.ctx, "Failed to load plugin", "plugin", pluginName, err)
|
||||
case actionAdd:
|
||||
// New file - extract manifest and add to DB as disabled
|
||||
metadata, err := m.ExtractManifest(wasmPath)
|
||||
if err != nil {
|
||||
log.Error(m.ctx, "Failed to extract manifest from new plugin", "plugin", pluginName, err)
|
||||
return
|
||||
}
|
||||
case actionUnload:
|
||||
log.Debug("Unloading removed Plugin", "plugin", pluginName)
|
||||
if err := m.UnloadPlugin(pluginName); err != nil {
|
||||
log.Debug(m.ctx, "Plugin not loaded, skipping unload", "plugin", pluginName, err)
|
||||
if err := m.addPluginToDB(m.ctx, repo, pluginName, wasmPath, metadata); err != nil {
|
||||
log.Error(m.ctx, "Failed to add plugin to DB", "plugin", pluginName, err)
|
||||
}
|
||||
case actionReload:
|
||||
log.Debug("Reloading modified Plugin", "plugin", pluginName)
|
||||
if err := m.ReloadPlugin(pluginName); err != nil {
|
||||
log.Error(m.ctx, "Failed to reload plugin", "plugin", pluginName, err)
|
||||
|
||||
case actionUpdate:
|
||||
// File changed - extract manifest, update DB, disable if enabled
|
||||
metadata, err := m.ExtractManifest(wasmPath)
|
||||
if err != nil {
|
||||
log.Error(m.ctx, "Failed to extract manifest from changed plugin", "plugin", pluginName, err)
|
||||
// Try to update error in DB if plugin exists
|
||||
if dbPlugin, getErr := repo.Get(pluginName); getErr == nil {
|
||||
dbPlugin.LastError = err.Error()
|
||||
dbPlugin.UpdatedAt = time.Now()
|
||||
if dbPlugin.Enabled {
|
||||
_ = m.UnloadPlugin(pluginName)
|
||||
dbPlugin.Enabled = false
|
||||
}
|
||||
_ = repo.Put(dbPlugin)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
dbPlugin, err := repo.Get(pluginName)
|
||||
if err != nil {
|
||||
// Plugin not in DB yet, add it
|
||||
if addErr := m.addPluginToDB(m.ctx, repo, pluginName, wasmPath, metadata); addErr != nil {
|
||||
log.Error(m.ctx, "Failed to add plugin to DB", "plugin", pluginName, addErr)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
// Check if actually changed
|
||||
if dbPlugin.SHA256 == metadata.SHA256 {
|
||||
return // No actual change
|
||||
}
|
||||
|
||||
if err := m.updatePluginInDB(m.ctx, repo, dbPlugin, wasmPath, 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(pluginName)
|
||||
if err != nil {
|
||||
log.Debug(m.ctx, "Removed plugin not in DB", "plugin", pluginName)
|
||||
return
|
||||
}
|
||||
|
||||
if err := m.removePluginFromDB(m.ctx, repo, dbPlugin); err != nil {
|
||||
log.Error(m.ctx, "Failed to delete plugin from DB", "plugin", pluginName, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@ -6,6 +6,7 @@ import (
|
||||
"path/filepath"
|
||||
|
||||
"github.com/navidrome/navidrome/conf"
|
||||
"github.com/navidrome/navidrome/tests"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
"github.com/rjeczalik/notify"
|
||||
@ -42,7 +43,7 @@ var _ = Describe("Plugin Watcher", func() {
|
||||
}
|
||||
|
||||
Describe("Plugin event processing (integration)", func() {
|
||||
// These tests verify the full flow with actual WASM plugin loading.
|
||||
// These tests verify the DB-driven flow with actual WASM plugin loading.
|
||||
|
||||
AfterEach(func() {
|
||||
// Clean up: unload plugin if loaded, remove copied file
|
||||
@ -50,28 +51,70 @@ var _ = Describe("Plugin Watcher", func() {
|
||||
_ = os.Remove(filepath.Join(tmpDir, "test-metadata-agent.wasm"))
|
||||
})
|
||||
|
||||
It("loads a plugin on CREATE event", func() {
|
||||
It("adds plugin to DB on CREATE event", func() {
|
||||
copyTestPlugin()
|
||||
manager.processPluginEvent("test-metadata-agent", notify.Create)
|
||||
Expect(manager.PluginNames(string(CapabilityMetadataAgent))).To(ContainElement("test-metadata-agent"))
|
||||
})
|
||||
|
||||
It("reloads a plugin on WRITE event", func() {
|
||||
copyTestPlugin()
|
||||
err := manager.LoadPlugin("test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
manager.processPluginEvent("test-metadata-agent", notify.Write)
|
||||
Expect(manager.PluginNames(string(CapabilityMetadataAgent))).To(ContainElement("test-metadata-agent"))
|
||||
})
|
||||
|
||||
It("unloads a plugin on REMOVE event", func() {
|
||||
copyTestPlugin()
|
||||
err := manager.LoadPlugin("test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
manager.processPluginEvent("test-metadata-agent", notify.Remove)
|
||||
// Plugin should be in DB but not loaded (starts disabled)
|
||||
Expect(manager.PluginNames(string(CapabilityMetadataAgent))).ToNot(ContainElement("test-metadata-agent"))
|
||||
|
||||
// Verify it was added to DB
|
||||
repo := manager.ds.Plugin(ctx)
|
||||
plugin, err := repo.Get("test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(plugin.ID).To(Equal("test-metadata-agent"))
|
||||
Expect(plugin.Enabled).To(BeFalse())
|
||||
})
|
||||
|
||||
It("updates DB and disables plugin on WRITE event when file changes", func() {
|
||||
copyTestPlugin()
|
||||
|
||||
// First add and enable the plugin
|
||||
manager.processPluginEvent("test-metadata-agent", notify.Create)
|
||||
err := manager.EnablePlugin(ctx, "test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(manager.PluginNames(string(CapabilityMetadataAgent))).To(ContainElement("test-metadata-agent"))
|
||||
|
||||
// Modify the stored SHA256 in DB to simulate a file change
|
||||
// (In reality, the file would have different content)
|
||||
repo := manager.ds.Plugin(ctx)
|
||||
plugin, err := repo.Get("test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
plugin.SHA256 = "different-hash-to-simulate-change"
|
||||
err = repo.Put(plugin)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
// Simulate modification - the plugin should be disabled and unloaded
|
||||
manager.processPluginEvent("test-metadata-agent", notify.Write)
|
||||
|
||||
// Should be unloaded
|
||||
Expect(manager.PluginNames(string(CapabilityMetadataAgent))).ToNot(ContainElement("test-metadata-agent"))
|
||||
|
||||
// But still in DB (just disabled)
|
||||
plugin, err = repo.Get("test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
Expect(plugin.Enabled).To(BeFalse())
|
||||
})
|
||||
|
||||
It("removes plugin from DB on REMOVE event", func() {
|
||||
copyTestPlugin()
|
||||
|
||||
// First add and enable the plugin
|
||||
manager.processPluginEvent("test-metadata-agent", notify.Create)
|
||||
err := manager.EnablePlugin(ctx, "test-metadata-agent")
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
|
||||
// Simulate removal - plugin should be unloaded and removed from DB
|
||||
_ = os.Remove(filepath.Join(tmpDir, "test-metadata-agent.wasm"))
|
||||
manager.processPluginEvent("test-metadata-agent", notify.Remove)
|
||||
|
||||
// Should be unloaded
|
||||
Expect(manager.PluginNames(string(CapabilityMetadataAgent))).ToNot(ContainElement("test-metadata-agent"))
|
||||
|
||||
// And removed from DB
|
||||
repo := manager.ds.Plugin(ctx)
|
||||
_, err = repo.Get("test-metadata-agent")
|
||||
Expect(err).To(HaveOccurred())
|
||||
})
|
||||
})
|
||||
|
||||
@ -85,9 +128,16 @@ var _ = Describe("Plugin Watcher", func() {
|
||||
_ = manager.Stop()
|
||||
|
||||
conf.Server.Plugins.AutoReload = true
|
||||
|
||||
// Set up a mock DataStore for the auto-reload manager
|
||||
mockPluginRepo := tests.CreateMockPluginRepo()
|
||||
mockPluginRepo.Permitted = true
|
||||
dataStore := &tests.MockDataStore{MockedPlugin: mockPluginRepo}
|
||||
|
||||
autoReloadManager := &Manager{
|
||||
plugins: make(map[string]*plugin),
|
||||
}
|
||||
autoReloadManager.SetDataStore(dataStore)
|
||||
err := autoReloadManager.Start(ctx)
|
||||
Expect(err).ToNot(HaveOccurred())
|
||||
DeferCleanup(autoReloadManager.Stop)
|
||||
@ -102,31 +152,26 @@ var _ = Describe("Plugin Watcher", func() {
|
||||
// These are fast unit tests for the pure routing logic.
|
||||
// No WASM compilation, no file I/O - runs in microseconds.
|
||||
|
||||
DescribeTable("returns correct action for event type and loaded state",
|
||||
func(eventType notify.Event, isLoaded bool, expected pluginAction) {
|
||||
Expect(determinePluginAction(eventType, isLoaded)).To(Equal(expected))
|
||||
DescribeTable("returns correct action for event type",
|
||||
func(eventType notify.Event, expected pluginAction) {
|
||||
Expect(determinePluginAction(eventType)).To(Equal(expected))
|
||||
},
|
||||
// CREATE events - always load
|
||||
Entry("CREATE when not loaded", notify.Create, false, actionLoad),
|
||||
Entry("CREATE when loaded", notify.Create, true, actionLoad),
|
||||
// CREATE events - add to DB
|
||||
Entry("CREATE", notify.Create, actionAdd),
|
||||
|
||||
// WRITE events - reload if loaded, load if not
|
||||
Entry("WRITE when not loaded", notify.Write, false, actionLoad),
|
||||
Entry("WRITE when loaded", notify.Write, true, actionReload),
|
||||
// WRITE events - update in DB
|
||||
Entry("WRITE", notify.Write, actionUpdate),
|
||||
|
||||
// REMOVE events - always unload
|
||||
Entry("REMOVE when not loaded", notify.Remove, false, actionUnload),
|
||||
Entry("REMOVE when loaded", notify.Remove, true, actionUnload),
|
||||
// REMOVE events - remove from DB
|
||||
Entry("REMOVE", notify.Remove, actionRemove),
|
||||
|
||||
// RENAME events - treated same as REMOVE
|
||||
Entry("RENAME when not loaded", notify.Rename, false, actionUnload),
|
||||
Entry("RENAME when loaded", notify.Rename, true, actionUnload),
|
||||
Entry("RENAME", notify.Rename, actionRemove),
|
||||
)
|
||||
|
||||
It("returns actionNone for unknown event types", func() {
|
||||
// Event type 0 or other unknown values
|
||||
Expect(determinePluginAction(0, false)).To(Equal(actionNone))
|
||||
Expect(determinePluginAction(0, true)).To(Equal(actionNone))
|
||||
Expect(determinePluginAction(0)).To(Equal(actionNone))
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user