Plugins: Plugin Store API returns DTO model (#41340)

* toying around

* fix refs

* remove unused fields

* go further

* add context

* ensure streaming handler is set
This commit is contained in:
Will Browne
2021-11-17 12:04:22 +01:00
committed by GitHub
parent dbb8246b6b
commit 2e3e7a7e55
24 changed files with 494 additions and 353 deletions
+56 -151
View File
@@ -11,7 +11,6 @@ import (
"net/url"
"os"
"path/filepath"
"strings"
"sync"
"time"
@@ -45,7 +44,7 @@ type PluginManager struct {
cfg *setting.Cfg
requestValidator models.PluginRequestValidator
sqlStore *sqlstore.SQLStore
plugins map[string]*plugins.Plugin
store map[string]*plugins.Plugin
pluginInstaller plugins.Installer
pluginLoader plugins.Loader
pluginsMu sync.RWMutex
@@ -68,7 +67,7 @@ func newManager(cfg *setting.Cfg, pluginRequestValidator models.PluginRequestVal
requestValidator: pluginRequestValidator,
sqlStore: sqlStore,
pluginLoader: pluginLoader,
plugins: map[string]*plugins.Plugin{},
store: map[string]*plugins.Plugin{},
log: log.New("plugin.manager"),
pluginInstaller: installer.New(false, cfg.BuildVersion, newInstallerLogger("plugin.installer", true)),
}
@@ -138,10 +137,36 @@ func (m *PluginManager) Run(ctx context.Context) error {
}
<-ctx.Done()
m.stop(ctx)
m.shutdown(ctx)
return ctx.Err()
}
func (m *PluginManager) plugin(pluginID string) (*plugins.Plugin, bool) {
m.pluginsMu.RLock()
defer m.pluginsMu.RUnlock()
p, exists := m.store[pluginID]
if !exists || (p.IsDecommissioned()) {
return nil, false
}
return p, true
}
func (m *PluginManager) plugins() []*plugins.Plugin {
m.pluginsMu.RLock()
defer m.pluginsMu.RUnlock()
res := make([]*plugins.Plugin, 0)
for _, p := range m.store {
if !p.IsDecommissioned() {
res = append(res, p)
}
}
return res
}
func (m *PluginManager) loadPlugins(paths ...string) error {
if len(paths) == 0 {
return nil
@@ -171,52 +196,15 @@ func (m *PluginManager) loadPlugins(paths ...string) error {
func (m *PluginManager) registeredPlugins() map[string]struct{} {
pluginsByID := make(map[string]struct{})
m.pluginsMu.RLock()
defer m.pluginsMu.RUnlock()
for _, p := range m.plugins {
for _, p := range m.plugins() {
pluginsByID[p.ID] = struct{}{}
}
return pluginsByID
}
func (m *PluginManager) Plugin(pluginID string) *plugins.Plugin {
m.pluginsMu.RLock()
p, ok := m.plugins[pluginID]
m.pluginsMu.RUnlock()
if ok && (p.IsDecommissioned()) {
return nil
}
return p
}
func (m *PluginManager) Plugins(pluginTypes ...plugins.Type) []*plugins.Plugin {
// if no types passed, assume all
if len(pluginTypes) == 0 {
pluginTypes = plugins.PluginTypes
}
var requestedTypes = make(map[plugins.Type]struct{})
for _, pt := range pluginTypes {
requestedTypes[pt] = struct{}{}
}
m.pluginsMu.RLock()
var pluginsList []*plugins.Plugin
for _, p := range m.plugins {
if _, exists := requestedTypes[p.Type]; exists {
pluginsList = append(pluginsList, p)
}
}
m.pluginsMu.RUnlock()
return pluginsList
}
func (m *PluginManager) Renderer() *plugins.Plugin {
for _, p := range m.plugins {
for _, p := range m.plugins() {
if p.IsRenderer() {
return p
}
@@ -226,8 +214,8 @@ func (m *PluginManager) Renderer() *plugins.Plugin {
}
func (m *PluginManager) QueryData(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
plugin := m.Plugin(req.PluginContext.PluginID)
if plugin == nil {
plugin, exists := m.plugin(req.PluginContext.PluginID)
if !exists {
return nil, backendplugin.ErrPluginNotRegistered
}
@@ -291,8 +279,8 @@ func (m *PluginManager) CallResource(pCtx backend.PluginContext, reqCtx *models.
}
func (m *PluginManager) callResourceInternal(w http.ResponseWriter, req *http.Request, pCtx backend.PluginContext) error {
p := m.Plugin(pCtx.PluginID)
if p == nil {
p, exists := m.plugin(pCtx.PluginID)
if !exists {
return backendplugin.ErrPluginNotRegistered
}
@@ -419,8 +407,8 @@ func flushStream(plugin backendplugin.Plugin, stream callResourceClientResponseS
}
func (m *PluginManager) CollectMetrics(ctx context.Context, pluginID string) (*backend.CollectMetricsResult, error) {
p := m.Plugin(pluginID)
if p == nil {
p, exists := m.plugin(pluginID)
if !exists {
return nil, backendplugin.ErrPluginNotRegistered
}
@@ -450,8 +438,8 @@ func (m *PluginManager) CheckHealth(ctx context.Context, req *backend.CheckHealt
}, nil
}
p := m.Plugin(req.PluginContext.PluginID)
if p == nil {
p, exists := m.plugin(req.PluginContext.PluginID)
if !exists {
return nil, backendplugin.ErrPluginNotRegistered
}
@@ -504,96 +492,14 @@ func (m *PluginManager) RunStream(ctx context.Context, req *backend.RunStreamReq
}
func (m *PluginManager) isRegistered(pluginID string) bool {
p := m.Plugin(pluginID)
if p == nil {
p, exists := m.plugin(pluginID)
if !exists {
return false
}
return !p.IsDecommissioned()
}
func (m *PluginManager) Add(ctx context.Context, pluginID, version string, opts plugins.AddOpts) error {
var pluginZipURL string
if opts.PluginRepoURL == "" {
opts.PluginRepoURL = grafanaComURL
}
plugin := m.Plugin(pluginID)
if plugin != nil {
if !plugin.IsExternalPlugin() {
return plugins.ErrInstallCorePlugin
}
if plugin.Info.Version == version {
return plugins.DuplicateError{
PluginID: plugin.ID,
ExistingPluginDir: plugin.PluginDir,
}
}
// get plugin update information to confirm if upgrading is possible
updateInfo, err := m.pluginInstaller.GetUpdateInfo(ctx, pluginID, version, opts.PluginRepoURL)
if err != nil {
return err
}
pluginZipURL = updateInfo.PluginZipURL
// remove existing installation of plugin
err = m.Remove(ctx, plugin.ID)
if err != nil {
return err
}
}
if opts.PluginInstallDir == "" {
opts.PluginInstallDir = m.cfg.PluginsPath
}
if opts.PluginZipURL == "" {
opts.PluginZipURL = pluginZipURL
}
err := m.pluginInstaller.Install(ctx, pluginID, version, opts.PluginInstallDir, opts.PluginZipURL, opts.PluginRepoURL)
if err != nil {
return err
}
err = m.loadPlugins(opts.PluginInstallDir)
if err != nil {
return err
}
return nil
}
func (m *PluginManager) Remove(ctx context.Context, pluginID string) error {
plugin := m.Plugin(pluginID)
if plugin == nil {
return plugins.ErrPluginNotInstalled
}
if !plugin.IsExternalPlugin() {
return plugins.ErrUninstallCorePlugin
}
// extra security check to ensure we only remove plugins that are located in the configured plugins directory
path, err := filepath.Rel(m.cfg.PluginsPath, plugin.PluginDir)
if err != nil || strings.HasPrefix(path, ".."+string(filepath.Separator)) {
return plugins.ErrUninstallOutsideOfPluginDir
}
if m.isRegistered(pluginID) {
err := m.unregisterAndStop(ctx, plugin)
if err != nil {
return err
}
}
return m.pluginInstaller.Uninstall(ctx, plugin.PluginDir)
}
func (m *PluginManager) LoadAndRegister(pluginID string, factory backendplugin.PluginFactoryFunc) error {
if m.isRegistered(pluginID) {
return fmt.Errorf("backend plugin %s already registered", pluginID)
@@ -620,9 +526,9 @@ func (m *PluginManager) LoadAndRegister(pluginID string, factory backendplugin.P
}
func (m *PluginManager) Routes() []*plugins.StaticRoute {
staticRoutes := []*plugins.StaticRoute{}
staticRoutes := make([]*plugins.StaticRoute, 0)
for _, p := range m.Plugins() {
for _, p := range m.plugins() {
if p.StaticRoute() != nil {
staticRoutes = append(staticRoutes, p.StaticRoute())
}
@@ -644,18 +550,16 @@ func (m *PluginManager) registerAndStart(ctx context.Context, plugin *plugins.Pl
}
func (m *PluginManager) register(p *plugins.Plugin) error {
m.pluginsMu.Lock()
defer m.pluginsMu.Unlock()
pluginID := p.ID
if _, exists := m.plugins[pluginID]; exists {
return fmt.Errorf("plugin %s already registered", pluginID)
if m.isRegistered(p.ID) {
return fmt.Errorf("plugin %s is already registered", p.ID)
}
m.plugins[pluginID] = p
m.pluginsMu.Lock()
m.store[p.ID] = p
m.pluginsMu.Unlock()
if !p.IsCorePlugin() {
m.log.Info("Plugin registered", "pluginId", pluginID)
m.log.Info("Plugin registered", "pluginId", p.ID)
}
return nil
@@ -663,6 +567,9 @@ func (m *PluginManager) register(p *plugins.Plugin) error {
func (m *PluginManager) unregisterAndStop(ctx context.Context, p *plugins.Plugin) error {
m.log.Debug("Stopping plugin process", "pluginId", p.ID)
m.pluginsMu.Lock()
defer m.pluginsMu.Unlock()
if err := p.Decommission(); err != nil {
return err
}
@@ -671,7 +578,7 @@ func (m *PluginManager) unregisterAndStop(ctx context.Context, p *plugins.Plugin
return err
}
delete(m.plugins, p.ID)
delete(m.store, p.ID)
m.log.Debug("Plugin unregistered", "pluginId", p.ID)
return nil
@@ -742,12 +649,10 @@ func restartKilledProcess(ctx context.Context, p *plugins.Plugin) error {
}
}
// stop stops a backend plugin process
func (m *PluginManager) stop(ctx context.Context) {
m.pluginsMu.RLock()
defer m.pluginsMu.RUnlock()
// shutdown stops all backend plugin processes
func (m *PluginManager) shutdown(ctx context.Context) {
var wg sync.WaitGroup
for _, p := range m.plugins {
for _, p := range m.plugins() {
wg.Add(1)
go func(p backendplugin.Plugin, ctx context.Context) {
defer wg.Done()